001/* 002 * Licensed to the Apache Software Foundation (ASF) under one or more 003 * contributor license agreements. See the NOTICE file distributed with 004 * this work for additional information regarding copyright ownership. 005 * The ASF licenses this file to You under the Apache License, Version 2.0 006 * (the "License"); you may not use this file except in compliance with 007 * the License. You may obtain a copy of the License at 008 * 009 * http://www.apache.org/licenses/LICENSE-2.0 010 * 011 * Unless required by applicable law or agreed to in writing, software 012 * distributed under the License is distributed on an "AS IS" BASIS, 013 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 014 * See the License for the specific language governing permissions and 015 * limitations under the License. 016 */ 017package org.apache.activemq.transport.amqp.protocol; 018 019import static org.apache.activemq.transport.amqp.AmqpSupport.ANONYMOUS_RELAY; 020import static org.apache.activemq.transport.amqp.AmqpSupport.CONNECTION_OPEN_FAILED; 021import static org.apache.activemq.transport.amqp.AmqpSupport.CONTAINER_ID; 022import static org.apache.activemq.transport.amqp.AmqpSupport.DELAYED_DELIVERY; 023import static org.apache.activemq.transport.amqp.AmqpSupport.INVALID_FIELD; 024import static org.apache.activemq.transport.amqp.AmqpSupport.PLATFORM; 025import static org.apache.activemq.transport.amqp.AmqpSupport.PRODUCT; 026import static org.apache.activemq.transport.amqp.AmqpSupport.QUEUE_PREFIX; 027import static org.apache.activemq.transport.amqp.AmqpSupport.TEMP_QUEUE_CAPABILITY; 028import static org.apache.activemq.transport.amqp.AmqpSupport.TEMP_TOPIC_CAPABILITY; 029import static org.apache.activemq.transport.amqp.AmqpSupport.TOPIC_PREFIX; 030import static org.apache.activemq.transport.amqp.AmqpSupport.VERSION; 031import static org.apache.activemq.transport.amqp.AmqpSupport.contains; 032 033import java.io.BufferedReader; 034import java.io.IOException; 035import java.io.InputStream; 036import java.io.InputStreamReader; 037import java.nio.ByteBuffer; 038import java.util.HashMap; 039import java.util.Map; 040import java.util.concurrent.ConcurrentHashMap; 041import java.util.concurrent.ConcurrentMap; 042import java.util.concurrent.TimeUnit; 043import java.util.concurrent.atomic.AtomicInteger; 044 045import javax.jms.InvalidClientIDException; 046 047import org.apache.activemq.broker.BrokerService; 048import org.apache.activemq.broker.region.AbstractRegion; 049import org.apache.activemq.broker.region.DurableTopicSubscription; 050import org.apache.activemq.broker.region.RegionBroker; 051import org.apache.activemq.broker.region.Subscription; 052import org.apache.activemq.broker.region.TopicRegion; 053import org.apache.activemq.command.ActiveMQDestination; 054import org.apache.activemq.command.ActiveMQTempDestination; 055import org.apache.activemq.command.ActiveMQTempQueue; 056import org.apache.activemq.command.ActiveMQTempTopic; 057import org.apache.activemq.command.Command; 058import org.apache.activemq.command.ConnectionError; 059import org.apache.activemq.command.ConnectionId; 060import org.apache.activemq.command.ConnectionInfo; 061import org.apache.activemq.command.ConsumerControl; 062import org.apache.activemq.command.ConsumerId; 063import org.apache.activemq.command.ConsumerInfo; 064import org.apache.activemq.command.DestinationInfo; 065import org.apache.activemq.command.ExceptionResponse; 066import org.apache.activemq.command.LocalTransactionId; 067import org.apache.activemq.command.MessageDispatch; 068import org.apache.activemq.command.RemoveInfo; 069import org.apache.activemq.command.Response; 070import org.apache.activemq.command.SessionId; 071import org.apache.activemq.command.ShutdownInfo; 072import org.apache.activemq.command.TransactionId; 073import org.apache.activemq.transport.InactivityIOException; 074import org.apache.activemq.transport.amqp.AmqpHeader; 075import org.apache.activemq.transport.amqp.AmqpInactivityMonitor; 076import org.apache.activemq.transport.amqp.AmqpProtocolConverter; 077import org.apache.activemq.transport.amqp.AmqpProtocolException; 078import org.apache.activemq.transport.amqp.AmqpTransport; 079import org.apache.activemq.transport.amqp.AmqpTransportFilter; 080import org.apache.activemq.transport.amqp.AmqpWireFormat; 081import org.apache.activemq.transport.amqp.sasl.AmqpAuthenticator; 082import org.apache.activemq.util.IOExceptionSupport; 083import org.apache.activemq.util.IdGenerator; 084import org.apache.qpid.proton.Proton; 085import org.apache.qpid.proton.amqp.Symbol; 086import org.apache.qpid.proton.amqp.transaction.Coordinator; 087import org.apache.qpid.proton.amqp.transport.AmqpError; 088import org.apache.qpid.proton.amqp.transport.ErrorCondition; 089import org.apache.qpid.proton.engine.Collector; 090import org.apache.qpid.proton.engine.Connection; 091import org.apache.qpid.proton.engine.Delivery; 092import org.apache.qpid.proton.engine.EndpointState; 093import org.apache.qpid.proton.engine.Event; 094import org.apache.qpid.proton.engine.Link; 095import org.apache.qpid.proton.engine.Receiver; 096import org.apache.qpid.proton.engine.Sender; 097import org.apache.qpid.proton.engine.Session; 098import org.apache.qpid.proton.engine.Transport; 099import org.apache.qpid.proton.engine.impl.CollectorImpl; 100import org.apache.qpid.proton.engine.impl.ProtocolTracer; 101import org.apache.qpid.proton.engine.impl.TransportImpl; 102import org.apache.qpid.proton.framing.TransportFrame; 103import org.fusesource.hawtbuf.Buffer; 104import org.slf4j.Logger; 105import org.slf4j.LoggerFactory; 106 107/** 108 * Implements the mechanics of managing a single remote peer connection. 109 */ 110public class AmqpConnection implements AmqpProtocolConverter { 111 112 private static final Logger TRACE_FRAMES = AmqpTransportFilter.TRACE_FRAMES; 113 private static final Logger LOG = LoggerFactory.getLogger(AmqpConnection.class); 114 private static final int CHANNEL_MAX = 32767; 115 private static final String BROKER_VERSION; 116 private static final String BROKER_PLATFORM; 117 118 static { 119 String javaVersion = System.getProperty("java.version"); 120 121 BROKER_PLATFORM = "Java/" + (javaVersion == null ? "unknown" : javaVersion); 122 123 InputStream in = null; 124 String version = "<unknown-5.x>"; 125 if ((in = AmqpConnection.class.getResourceAsStream("/org/apache/activemq/version.txt")) != null) { 126 BufferedReader reader = new BufferedReader(new InputStreamReader(in)); 127 try { 128 version = reader.readLine(); 129 } catch(Exception e) { 130 } 131 } 132 BROKER_VERSION = version; 133 } 134 135 private final Transport protonTransport = Proton.transport(); 136 private final Connection protonConnection = Proton.connection(); 137 private final Collector eventCollector = new CollectorImpl(); 138 139 private final AmqpTransport amqpTransport; 140 private final AmqpWireFormat amqpWireFormat; 141 private final BrokerService brokerService; 142 143 private static final IdGenerator CONNECTION_ID_GENERATOR = new IdGenerator(); 144 private final AtomicInteger lastCommandId = new AtomicInteger(); 145 private final ConnectionId connectionId = new ConnectionId(CONNECTION_ID_GENERATOR.generateId()); 146 private final ConnectionInfo connectionInfo = new ConnectionInfo(); 147 private long nextSessionId; 148 private long nextTempDestinationId; 149 private long nextTransactionId; 150 private boolean closing; 151 private boolean closedSocket; 152 private AmqpAuthenticator authenticator; 153 154 private final Map<TransactionId, AmqpTransactionCoordinator> transactions = new HashMap<>(); 155 private final ConcurrentMap<Integer, ResponseHandler> resposeHandlers = new ConcurrentHashMap<>(); 156 private final ConcurrentMap<ConsumerId, AmqpSender> subscriptionsByConsumerId = new ConcurrentHashMap<>(); 157 158 public AmqpConnection(AmqpTransport transport, BrokerService brokerService) { 159 this.amqpTransport = transport; 160 161 AmqpInactivityMonitor monitor = transport.getInactivityMonitor(); 162 if (monitor != null) { 163 monitor.setAmqpTransport(amqpTransport); 164 } 165 166 this.amqpWireFormat = transport.getWireFormat(); 167 this.brokerService = brokerService; 168 169 // the configured maxFrameSize on the URI. 170 int maxFrameSize = amqpWireFormat.getMaxAmqpFrameSize(); 171 if (maxFrameSize > AmqpWireFormat.NO_AMQP_MAX_FRAME_SIZE) { 172 this.protonTransport.setMaxFrameSize(maxFrameSize); 173 try { 174 this.protonTransport.setOutboundFrameSizeLimit(maxFrameSize); 175 } catch (Throwable e) { 176 // Ignore if older proton-j was injected. 177 } 178 } 179 180 this.protonTransport.bind(this.protonConnection); 181 this.protonTransport.setChannelMax(CHANNEL_MAX); 182 this.protonTransport.setEmitFlowEventOnSend(false); 183 184 this.protonConnection.collect(eventCollector); 185 186 updateTracer(); 187 } 188 189 /** 190 * Load and return a <code>[]Symbol</code> that contains the connection capabilities 191 * offered to new connections 192 * 193 * @return the capabilities that are offered to new clients on connect. 194 */ 195 protected Symbol[] getConnectionCapabilitiesOffered() { 196 return new Symbol[]{ ANONYMOUS_RELAY, DELAYED_DELIVERY }; 197 } 198 199 /** 200 * Load and return a <code>Map<Symbol, Object></code> that contains the properties 201 * that this connection supplies to incoming connections. 202 * 203 * @return the properties that are offered to the incoming connection. 204 */ 205 protected Map<Symbol, Object> getConnetionProperties() { 206 Map<Symbol, Object> properties = new HashMap<>(); 207 208 properties.put(QUEUE_PREFIX, "queue://"); 209 properties.put(TOPIC_PREFIX, "topic://"); 210 properties.put(PRODUCT, "ActiveMQ"); 211 properties.put(VERSION, BROKER_VERSION); 212 properties.put(PLATFORM, BROKER_PLATFORM); 213 214 return properties; 215 } 216 217 /** 218 * Load and return a <code>Map<Symbol, Object></code> that contains the properties 219 * that this connection supplies to incoming connections when the open has failed 220 * and the remote should expect a close to follow. 221 * 222 * @return the properties that are offered to the incoming connection. 223 */ 224 protected Map<Symbol, Object> getFailedConnetionProperties() { 225 Map<Symbol, Object> properties = new HashMap<>(); 226 227 properties.put(CONNECTION_OPEN_FAILED, true); 228 229 return properties; 230 } 231 232 @Override 233 public void updateTracer() { 234 if (amqpTransport.isTrace()) { 235 ((TransportImpl) protonTransport).setProtocolTracer(new ProtocolTracer() { 236 @Override 237 public void receivedFrame(TransportFrame transportFrame) { 238 TRACE_FRAMES.trace("{} | RECV: {}", AmqpConnection.this.amqpTransport.getRemoteAddress(), transportFrame.getBody()); 239 } 240 241 @Override 242 public void sentFrame(TransportFrame transportFrame) { 243 TRACE_FRAMES.trace("{} | SENT: {}", AmqpConnection.this.amqpTransport.getRemoteAddress(), transportFrame.getBody()); 244 } 245 }); 246 } 247 } 248 249 @Override 250 public long keepAlive() throws IOException { 251 long rescheduleAt = 0l; 252 253 LOG.trace("Performing connection:{} keep-alive processing", amqpTransport.getRemoteAddress()); 254 255 if (protonConnection.getLocalState() != EndpointState.CLOSED) { 256 // Using nano time since it is not related to the wall clock, which may change 257 long now = TimeUnit.NANOSECONDS.toMillis(System.nanoTime()); 258 long deadline = protonTransport.tick(now); 259 pumpProtonToSocket(); 260 if (protonTransport.isClosed()) { 261 LOG.debug("Transport closed after inactivity check."); 262 throw new InactivityIOException("Channel was inactive for too long"); 263 } else { 264 if(deadline != 0) { 265 // caller treats 0 as no-work, ensure value is at least 1 as there was a deadline 266 rescheduleAt = Math.max(deadline - now, 1); 267 } 268 } 269 } 270 271 LOG.trace("Connection:{} keep alive processing done, next update in {} milliseconds.", 272 amqpTransport.getRemoteAddress(), rescheduleAt); 273 274 return rescheduleAt; 275 } 276 277 //----- Connection Properties Accessors ----------------------------------// 278 279 /** 280 * @return the amount of credit assigned to AMQP receiver links created from 281 * sender links on the remote peer. 282 */ 283 public int getConfiguredReceiverCredit() { 284 return amqpWireFormat.getProducerCredit(); 285 } 286 287 /** 288 * @return the transformer type that was configured for this AMQP transport. 289 */ 290 public String getConfiguredTransformer() { 291 return amqpWireFormat.getTransformer(); 292 } 293 294 /** 295 * @return the ActiveMQ ConnectionId that identifies this AMQP Connection. 296 */ 297 public ConnectionId getConnectionId() { 298 return connectionId; 299 } 300 301 /** 302 * @return the Client ID used to create the connection with ActiveMQ 303 */ 304 public String getClientId() { 305 return connectionInfo.getClientId(); 306 } 307 308 /** 309 * @return the configured max frame size allowed for incoming messages. 310 */ 311 public long getMaxFrameSize() { 312 return amqpWireFormat.getMaxFrameSize(); 313 } 314 315 //----- Proton Event handling and IO support -----------------------------// 316 317 void pumpProtonToSocket() { 318 try { 319 boolean done = false; 320 while (!done) { 321 ByteBuffer toWrite = protonTransport.getOutputBuffer(); 322 if (toWrite != null && toWrite.hasRemaining()) { 323 LOG.trace("Server: Sending {} bytes out", toWrite.limit()); 324 amqpTransport.sendToAmqp(toWrite); 325 protonTransport.outputConsumed(); 326 } else { 327 done = true; 328 } 329 } 330 } catch (IOException e) { 331 amqpTransport.onException(e); 332 } 333 } 334 335 @SuppressWarnings("deprecation") 336 @Override 337 public void onAMQPData(Object command) throws Exception { 338 Buffer frame; 339 if (command.getClass() == AmqpHeader.class) { 340 AmqpHeader header = (AmqpHeader) command; 341 342 if (amqpWireFormat.isHeaderValid(header, authenticator != null)) { 343 LOG.trace("Connection from an AMQP v1.0 client initiated. {}", header); 344 } else { 345 LOG.warn("Connection attempt from non AMQP v1.0 client. {}", header); 346 AmqpHeader reply = amqpWireFormat.getMinimallySupportedHeader(); 347 amqpTransport.sendToAmqp(reply.getBuffer()); 348 handleException(new AmqpProtocolException( 349 "Connection from client using unsupported AMQP attempted", true)); 350 } 351 352 switch (header.getProtocolId()) { 353 case 0: 354 authenticator = null; 355 break; // nothing to do.. 356 case 3: // Client will be using SASL for auth.. 357 authenticator = new AmqpAuthenticator(amqpTransport, protonTransport.sasl(), brokerService); 358 break; 359 default: 360 } 361 frame = header.getBuffer(); 362 } else { 363 frame = (Buffer) command; 364 } 365 366 if (protonTransport.isClosed()) { 367 LOG.debug("Ignoring incoming AMQP data, transport is closed."); 368 return; 369 } 370 371 LOG.trace("Server: Received from client: {} bytes", frame.getLength()); 372 373 while (frame.length > 0) { 374 try { 375 int count = protonTransport.input(frame.data, frame.offset, frame.length); 376 frame.moveHead(count); 377 } catch (Throwable e) { 378 handleException(new AmqpProtocolException("Could not decode AMQP frame: " + frame, true, e)); 379 return; 380 } 381 382 if (authenticator != null) { 383 processSaslExchange(); 384 } else { 385 processProtonEvents(); 386 } 387 } 388 } 389 390 private void processSaslExchange() throws Exception { 391 authenticator.processSaslExchange(connectionInfo); 392 if (authenticator.isDone()) { 393 amqpTransport.getWireFormat().resetMagicRead(); 394 } 395 pumpProtonToSocket(); 396 } 397 398 private void processProtonEvents() throws Exception { 399 try { 400 Event event = null; 401 while ((event = eventCollector.peek()) != null) { 402 if (amqpTransport.isTrace()) { 403 LOG.trace("Server: Processing event: {}", event.getType()); 404 } 405 switch (event.getType()) { 406 case CONNECTION_REMOTE_OPEN: 407 processConnectionOpen(event.getConnection()); 408 break; 409 case CONNECTION_REMOTE_CLOSE: 410 processConnectionClose(event.getConnection()); 411 break; 412 case SESSION_REMOTE_OPEN: 413 processSessionOpen(event.getSession()); 414 break; 415 case SESSION_REMOTE_CLOSE: 416 processSessionClose(event.getSession()); 417 break; 418 case LINK_REMOTE_OPEN: 419 processLinkOpen(event.getLink()); 420 break; 421 case LINK_REMOTE_DETACH: 422 processLinkDetach(event.getLink()); 423 break; 424 case LINK_REMOTE_CLOSE: 425 processLinkClose(event.getLink()); 426 break; 427 case LINK_FLOW: 428 processLinkFlow(event.getLink()); 429 break; 430 case DELIVERY: 431 processDelivery(event.getDelivery()); 432 break; 433 default: 434 break; 435 } 436 437 eventCollector.pop(); 438 } 439 440 } catch (Throwable e) { 441 handleException(new AmqpProtocolException("Could not process AMQP commands", true, e)); 442 } 443 444 pumpProtonToSocket(); 445 } 446 447 protected void processConnectionOpen(Connection connection) throws Exception { 448 449 stopConnectionTimeoutChecker(); 450 451 connectionInfo.setResponseRequired(true); 452 connectionInfo.setConnectionId(connectionId); 453 454 String clientId = protonConnection.getRemoteContainer(); 455 if (clientId != null && !clientId.isEmpty()) { 456 connectionInfo.setClientId(clientId); 457 } 458 459 connectionInfo.setTransportContext(amqpTransport.getPeerCertificates()); 460 461 if (connection.getTransport().getRemoteIdleTimeout() > 0 && !amqpTransport.isUseInactivityMonitor()) { 462 // We cannot meet the requested Idle processing because the inactivity monitor is 463 // disabled so we won't send idle frames to match the request. 464 protonConnection.setProperties(getFailedConnetionProperties()); 465 protonConnection.open(); 466 protonConnection.setCondition(new ErrorCondition(AmqpError.PRECONDITION_FAILED, "Cannot send idle frames")); 467 protonConnection.close(); 468 pumpProtonToSocket(); 469 470 amqpTransport.onException(new IOException( 471 "Connection failed, remote requested idle processing but inactivity monitoring is disbaled.")); 472 return; 473 } 474 475 sendToActiveMQ(connectionInfo, new ResponseHandler() { 476 @Override 477 public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException { 478 Throwable exception = null; 479 try { 480 if (response.isException()) { 481 protonConnection.setProperties(getFailedConnetionProperties()); 482 protonConnection.open(); 483 484 exception = ((ExceptionResponse) response).getException(); 485 if (exception instanceof SecurityException) { 486 protonConnection.setCondition(new ErrorCondition(AmqpError.UNAUTHORIZED_ACCESS, exception.getMessage())); 487 } else if (exception instanceof InvalidClientIDException) { 488 ErrorCondition condition = new ErrorCondition(AmqpError.INVALID_FIELD, exception.getMessage()); 489 490 Map<Symbol, Object> infoMap = new HashMap<> (); 491 infoMap.put(INVALID_FIELD, CONTAINER_ID); 492 condition.setInfo(infoMap); 493 494 protonConnection.setCondition(condition); 495 } else { 496 protonConnection.setCondition(new ErrorCondition(AmqpError.ILLEGAL_STATE, exception.getMessage())); 497 } 498 499 protonConnection.close(); 500 } else { 501 if (amqpTransport.isUseInactivityMonitor() && amqpWireFormat.getIdleTimeout() > 0) { 502 LOG.trace("Connection requesting Idle timeout of: {} mills", amqpWireFormat.getIdleTimeout()); 503 protonTransport.setIdleTimeout(amqpWireFormat.getIdleTimeout()); 504 } 505 506 protonConnection.setOfferedCapabilities(getConnectionCapabilitiesOffered()); 507 protonConnection.setProperties(getConnetionProperties()); 508 protonConnection.setContainer(brokerService.getBrokerName()); 509 protonConnection.open(); 510 511 configureInactivityMonitor(); 512 } 513 } finally { 514 pumpProtonToSocket(); 515 516 if (response.isException()) { 517 amqpTransport.onException(IOExceptionSupport.create(exception)); 518 } 519 } 520 } 521 }); 522 } 523 524 protected void processConnectionClose(Connection connection) throws Exception { 525 if (!closing) { 526 closing = true; 527 sendToActiveMQ(new RemoveInfo(connectionId), new ResponseHandler() { 528 @Override 529 public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException { 530 protonConnection.close(); 531 protonConnection.free(); 532 533 if (!closedSocket) { 534 pumpProtonToSocket(); 535 } 536 } 537 }); 538 539 sendToActiveMQ(new ShutdownInfo()); 540 } 541 } 542 543 protected void processSessionOpen(Session protonSession) throws Exception { 544 new AmqpSession(this, getNextSessionId(), protonSession).open(); 545 } 546 547 protected void processSessionClose(Session protonSession) throws Exception { 548 if (protonSession.getContext() != null) { 549 ((AmqpResource) protonSession.getContext()).close(); 550 } else { 551 protonSession.close(); 552 protonSession.free(); 553 } 554 } 555 556 protected void processLinkOpen(Link link) throws Exception { 557 link.setSource(link.getRemoteSource()); 558 link.setTarget(link.getRemoteTarget()); 559 560 AmqpSession session = (AmqpSession) link.getSession().getContext(); 561 if (link instanceof Receiver) { 562 if (link.getRemoteTarget() instanceof Coordinator) { 563 session.createCoordinator((Receiver) link); 564 } else { 565 session.createReceiver((Receiver) link); 566 } 567 } else { 568 session.createSender((Sender) link); 569 } 570 } 571 572 protected void processLinkDetach(Link link) throws Exception { 573 Object context = link.getContext(); 574 575 if (context instanceof AmqpLink) { 576 ((AmqpLink) context).detach(); 577 } else { 578 link.detach(); 579 link.free(); 580 } 581 } 582 583 protected void processLinkClose(Link link) throws Exception { 584 Object context = link.getContext(); 585 586 if (context instanceof AmqpLink) { 587 ((AmqpLink) context).close();; 588 } else { 589 link.close(); 590 link.free(); 591 } 592 } 593 594 protected void processLinkFlow(Link link) throws Exception { 595 Object context = link.getContext(); 596 if (context instanceof AmqpLink) { 597 ((AmqpLink) context).flow(); 598 } 599 } 600 601 protected void processDelivery(Delivery delivery) throws Exception { 602 if (!delivery.isPartial()) { 603 Object context = delivery.getLink().getContext(); 604 if (context instanceof AmqpLink) { 605 AmqpLink amqpLink = (AmqpLink) context; 606 amqpLink.delivery(delivery); 607 } 608 } 609 } 610 611 //----- Event entry points for ActiveMQ commands and errors --------------// 612 613 @Override 614 public void onAMQPException(IOException error) { 615 closedSocket = true; 616 if (!closing) { 617 try { 618 closing = true; 619 // Attempt to inform the other end that we are going to close 620 // so that the client doesn't wait around forever. 621 protonConnection.setCondition(new ErrorCondition(AmqpError.DECODE_ERROR, error.getMessage())); 622 protonConnection.close(); 623 pumpProtonToSocket(); 624 } catch (Exception ignore) { 625 } 626 amqpTransport.sendToActiveMQ(error); 627 } else { 628 try { 629 amqpTransport.stop(); 630 } catch (Exception ignore) { 631 } 632 } 633 } 634 635 @Override 636 public void onActiveMQCommand(Command command) throws Exception { 637 if (command.isResponse()) { 638 Response response = (Response) command; 639 ResponseHandler rh = resposeHandlers.remove(Integer.valueOf(response.getCorrelationId())); 640 if (rh != null) { 641 rh.onResponse(this, response); 642 } else { 643 // Pass down any unexpected errors. Should this close the connection? 644 if (response.isException()) { 645 Throwable exception = ((ExceptionResponse) response).getException(); 646 handleException(exception); 647 } 648 } 649 } else if (command.isMessageDispatch()) { 650 MessageDispatch dispatch = (MessageDispatch) command; 651 AmqpSender sender = subscriptionsByConsumerId.get(dispatch.getConsumerId()); 652 if (sender != null) { 653 // End of Queue Browse will have no Message object. 654 if (dispatch.getMessage() != null) { 655 LOG.trace("Dispatching MessageId: {} to consumer", dispatch.getMessage().getMessageId()); 656 } else { 657 LOG.trace("Dispatching End of Browse Command to consumer {}", dispatch.getConsumerId()); 658 } 659 sender.onMessageDispatch(dispatch); 660 if (dispatch.getMessage() != null) { 661 LOG.trace("Finished Dispatch of MessageId: {} to consumer", dispatch.getMessage().getMessageId()); 662 } 663 } 664 } else if (command.getDataStructureType() == ConnectionError.DATA_STRUCTURE_TYPE) { 665 // Pass down any unexpected async errors. Should this close the connection? 666 Throwable exception = ((ConnectionError) command).getException(); 667 handleException(exception); 668 } else if (command.isConsumerControl()) { 669 ConsumerControl control = (ConsumerControl) command; 670 AmqpSender sender = subscriptionsByConsumerId.get(control.getConsumerId()); 671 if (sender != null) { 672 sender.onConsumerControl(control); 673 } 674 } else if (command.isBrokerInfo()) { 675 // ignore 676 } else { 677 LOG.debug("Do not know how to process ActiveMQ Command {}", command); 678 } 679 } 680 681 //----- Utility methods for connection resources to use ------------------// 682 683 void registerSender(ConsumerId consumerId, AmqpSender sender) { 684 subscriptionsByConsumerId.put(consumerId, sender); 685 } 686 687 void unregisterSender(ConsumerId consumerId) { 688 subscriptionsByConsumerId.remove(consumerId); 689 } 690 691 void registerTransaction(TransactionId txId, AmqpTransactionCoordinator coordinator) { 692 transactions.put(txId, coordinator); 693 } 694 695 void unregisterTransaction(TransactionId txId) { 696 transactions.remove(txId); 697 } 698 699 AmqpTransactionCoordinator getTxCoordinator(TransactionId txId) { 700 return transactions.get(txId); 701 } 702 703 LocalTransactionId getNextTransactionId() { 704 return new LocalTransactionId(getConnectionId(), ++nextTransactionId); 705 } 706 707 ConsumerInfo lookupSubscription(String subscriptionName) throws AmqpProtocolException { 708 ConsumerInfo result = null; 709 RegionBroker regionBroker; 710 711 try { 712 regionBroker = (RegionBroker) brokerService.getBroker().getAdaptor(RegionBroker.class); 713 } catch (Exception e) { 714 throw new AmqpProtocolException("Error finding subscription: " + subscriptionName + ": " + e.getMessage(), false, e); 715 } 716 717 final TopicRegion topicRegion = (TopicRegion) regionBroker.getTopicRegion(); 718 DurableTopicSubscription subscription = topicRegion.lookupSubscription(subscriptionName, connectionInfo.getClientId()); 719 if (subscription != null) { 720 result = subscription.getConsumerInfo(); 721 } 722 723 return result; 724 } 725 726 727 Subscription lookupPrefetchSubscription(ConsumerInfo consumerInfo) { 728 Subscription subscription = null; 729 try { 730 subscription = ((AbstractRegion)((RegionBroker) brokerService.getBroker().getAdaptor(RegionBroker.class)).getRegion(consumerInfo.getDestination())).getSubscriptions().get(consumerInfo.getConsumerId()); 731 } catch (Exception e) { 732 LOG.warn("Error finding subscription for: " + consumerInfo + ": " + e.getMessage(), false, e); 733 } 734 return subscription; 735 } 736 737 ActiveMQDestination createTemporaryDestination(final Link link, Symbol[] capabilities) { 738 ActiveMQDestination rc = null; 739 if (contains(capabilities, TEMP_TOPIC_CAPABILITY)) { 740 rc = new ActiveMQTempTopic(connectionId, nextTempDestinationId++); 741 } else if (contains(capabilities, TEMP_QUEUE_CAPABILITY)) { 742 rc = new ActiveMQTempQueue(connectionId, nextTempDestinationId++); 743 } else { 744 LOG.debug("Dynamic link request with no type capability, defaults to Temporary Queue"); 745 rc = new ActiveMQTempQueue(connectionId, nextTempDestinationId++); 746 } 747 748 DestinationInfo info = new DestinationInfo(); 749 info.setConnectionId(connectionId); 750 info.setOperationType(DestinationInfo.ADD_OPERATION_TYPE); 751 info.setDestination(rc); 752 753 sendToActiveMQ(info, new ResponseHandler() { 754 755 @Override 756 public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException { 757 if (response.isException()) { 758 link.setSource(null); 759 760 Throwable exception = ((ExceptionResponse) response).getException(); 761 if (exception instanceof SecurityException) { 762 link.setCondition(new ErrorCondition(AmqpError.UNAUTHORIZED_ACCESS, exception.getMessage())); 763 } else { 764 link.setCondition(new ErrorCondition(AmqpError.INTERNAL_ERROR, exception.getMessage())); 765 } 766 767 link.close(); 768 link.free(); 769 } 770 } 771 }); 772 773 return rc; 774 } 775 776 void deleteTemporaryDestination(ActiveMQTempDestination destination) { 777 DestinationInfo info = new DestinationInfo(); 778 info.setConnectionId(connectionId); 779 info.setOperationType(DestinationInfo.REMOVE_OPERATION_TYPE); 780 info.setDestination(destination); 781 782 sendToActiveMQ(info, new ResponseHandler() { 783 784 @Override 785 public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException { 786 if (response.isException()) { 787 Throwable exception = ((ExceptionResponse) response).getException(); 788 LOG.debug("Error during temp destination removeal: {}", exception.getMessage()); 789 } 790 } 791 }); 792 } 793 794 void sendToActiveMQ(Command command) { 795 sendToActiveMQ(command, null); 796 } 797 798 void sendToActiveMQ(Command command, ResponseHandler handler) { 799 command.setCommandId(lastCommandId.incrementAndGet()); 800 if (handler != null) { 801 command.setResponseRequired(true); 802 resposeHandlers.put(Integer.valueOf(command.getCommandId()), handler); 803 } 804 amqpTransport.sendToActiveMQ(command); 805 } 806 807 void handleException(Throwable exception) { 808 LOG.debug("Exception detail", exception); 809 if (exception instanceof AmqpProtocolException) { 810 onAMQPException((IOException) exception); 811 } else { 812 try { 813 // Must ensure that the broker removes Connection resources. 814 sendToActiveMQ(new ShutdownInfo()); 815 amqpTransport.stop(); 816 } catch (Throwable e) { 817 LOG.error("Failed to stop AMQP Transport ", e); 818 } 819 } 820 } 821 822 //----- Internal implementation ------------------------------------------// 823 824 private SessionId getNextSessionId() { 825 return new SessionId(connectionId, nextSessionId++); 826 } 827 828 private void stopConnectionTimeoutChecker() { 829 AmqpInactivityMonitor monitor = amqpTransport.getInactivityMonitor(); 830 if (monitor != null) { 831 monitor.stopConnectionTimeoutChecker(); 832 } 833 } 834 835 private void configureInactivityMonitor() { 836 AmqpInactivityMonitor monitor = amqpTransport.getInactivityMonitor(); 837 if (monitor == null) { 838 return; 839 } 840 841 // If either end has idle timeout requirements then the tick method 842 // will give us a deadline on the next time we need to tick() in order 843 // to meet those obligations. 844 // Using nano time since it is not related to the wall clock, which may change 845 long now = TimeUnit.NANOSECONDS.toMillis(System.nanoTime()); 846 long nextIdleCheck = protonTransport.tick(now); 847 if (nextIdleCheck != 0) { 848 // monitor treats <= 0 as no work, ensure value is at least 1 as there was a deadline 849 long delay = Math.max(nextIdleCheck - now, 1); 850 LOG.trace("Connection keep-alive processing starts in: {}", delay); 851 monitor.startKeepAliveTask(delay); 852 } else { 853 LOG.trace("Connection does not require keep-alive processing"); 854 } 855 } 856}