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}