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.COPY;
020import static org.apache.activemq.transport.amqp.AmqpSupport.JMS_SELECTOR_FILTER_IDS;
021import static org.apache.activemq.transport.amqp.AmqpSupport.JMS_SELECTOR_NAME;
022import static org.apache.activemq.transport.amqp.AmqpSupport.LIFETIME_POLICY;
023import static org.apache.activemq.transport.amqp.AmqpSupport.NO_LOCAL_FILTER_IDS;
024import static org.apache.activemq.transport.amqp.AmqpSupport.NO_LOCAL_NAME;
025import static org.apache.activemq.transport.amqp.AmqpSupport.createDestination;
026import static org.apache.activemq.transport.amqp.AmqpSupport.findFilter;
027
028import java.io.IOException;
029import java.util.Arrays;
030import java.util.HashMap;
031import java.util.List;
032import java.util.Map;
033
034import javax.jms.InvalidSelectorException;
035
036import org.apache.activemq.command.ActiveMQDestination;
037import org.apache.activemq.command.ActiveMQTempDestination;
038import org.apache.activemq.command.ConsumerId;
039import org.apache.activemq.command.ConsumerInfo;
040import org.apache.activemq.command.ExceptionResponse;
041import org.apache.activemq.command.LocalTransactionId;
042import org.apache.activemq.command.ProducerId;
043import org.apache.activemq.command.ProducerInfo;
044import org.apache.activemq.command.RemoveInfo;
045import org.apache.activemq.command.Response;
046import org.apache.activemq.command.SessionId;
047import org.apache.activemq.command.SessionInfo;
048import org.apache.activemq.command.TransactionId;
049import org.apache.activemq.selector.SelectorParser;
050import org.apache.activemq.transport.amqp.AmqpProtocolConverter;
051import org.apache.activemq.transport.amqp.AmqpProtocolException;
052import org.apache.activemq.transport.amqp.AmqpSupport;
053import org.apache.activemq.util.IntrospectionSupport;
054import org.apache.qpid.proton.amqp.DescribedType;
055import org.apache.qpid.proton.amqp.Symbol;
056import org.apache.qpid.proton.amqp.messaging.DeleteOnClose;
057import org.apache.qpid.proton.amqp.messaging.Target;
058import org.apache.qpid.proton.amqp.messaging.TerminusDurability;
059import org.apache.qpid.proton.amqp.messaging.TerminusExpiryPolicy;
060import org.apache.qpid.proton.amqp.transport.AmqpError;
061import org.apache.qpid.proton.amqp.transport.ErrorCondition;
062import org.apache.qpid.proton.engine.Receiver;
063import org.apache.qpid.proton.engine.Sender;
064import org.apache.qpid.proton.engine.Session;
065import org.slf4j.Logger;
066import org.slf4j.LoggerFactory;
067
068/**
069 * Wraps the AMQP Session and provides the services needed to manage the remote
070 * peer requests for link establishment.
071 */
072public class AmqpSession implements AmqpResource {
073
074    private static final Logger LOG = LoggerFactory.getLogger(AmqpSession.class);
075
076    private final Map<ConsumerId, AmqpSender> consumers = new HashMap<>();
077
078    private final AmqpConnection connection;
079    private final Session protonSession;
080    private final SessionId sessionId;
081
082    private boolean enlisted;
083    private long nextProducerId = 0;
084    private long nextConsumerId = 0;
085
086    /**
087     * Create new AmqpSession instance whose parent is the given AmqpConnection.
088     *
089     * @param connection
090     *        the parent connection for this session.
091     * @param sessionId
092     *        the ActiveMQ SessionId that is used to identify this session.
093     * @param session
094     *        the AMQP Session that this class manages.
095     */
096    public AmqpSession(AmqpConnection connection, SessionId sessionId, Session session) {
097        this.connection = connection;
098        this.sessionId = sessionId;
099        this.protonSession = session;
100    }
101
102    @Override
103    public void open() {
104        LOG.debug("Session {} opened", getSessionId());
105
106        getEndpoint().setContext(this);
107        getEndpoint().setIncomingCapacity(Integer.MAX_VALUE);
108        getEndpoint().open();
109
110        connection.sendToActiveMQ(new SessionInfo(getSessionId()));
111    }
112
113    @Override
114    public void close() {
115        LOG.debug("Session {} closed", getSessionId());
116
117        connection.sendToActiveMQ(new RemoveInfo(getSessionId()), new ResponseHandler() {
118
119            @Override
120            public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
121                getEndpoint().setContext(null);
122                getEndpoint().close();
123                getEndpoint().free();
124            }
125        });
126    }
127
128    /**
129     * Commits all pending work for all resources managed under this session.
130     *
131     * @param txId
132     *      The specific TransactionId that is being committed.
133     *
134     * @throws Exception if an error occurs while attempting to commit work.
135     */
136    public void commit(LocalTransactionId txId) throws Exception {
137        for (AmqpSender consumer : consumers.values()) {
138            consumer.commit(txId);
139        }
140
141        enlisted = false;
142    }
143
144    /**
145     * Rolls back any pending work being down under this session.
146     *
147     * @param txId
148     *      The specific TransactionId that is being rolled back.
149     *
150     * @throws Exception if an error occurs while attempting to roll back work.
151     */
152    public void rollback(LocalTransactionId txId) throws Exception {
153        for (AmqpSender consumer : consumers.values()) {
154            consumer.rollback(txId);
155        }
156
157        enlisted = false;
158    }
159
160    /**
161     * Used to direct all Session managed Senders to push any queued Messages
162     * out to the remote peer.
163     *
164     * @throws Exception if an error occurs while flushing the messages.
165     */
166    public void flushPendingMessages() throws Exception {
167        for (AmqpSender consumer : consumers.values()) {
168            consumer.pumpOutbound();
169        }
170    }
171
172    public void createCoordinator(final Receiver protonReceiver) throws Exception {
173        AmqpTransactionCoordinator txCoordinator = new AmqpTransactionCoordinator(this, protonReceiver);
174        txCoordinator.flow(connection.getConfiguredReceiverCredit());
175        txCoordinator.open();
176    }
177
178    public void createReceiver(final Receiver protonReceiver) throws Exception {
179        org.apache.qpid.proton.amqp.transport.Target remoteTarget = protonReceiver.getRemoteTarget();
180
181        ProducerInfo producerInfo = new ProducerInfo(getNextProducerId());
182        final AmqpReceiver receiver = new AmqpReceiver(this, protonReceiver, producerInfo);
183
184        LOG.debug("opening new receiver {} on link: {}", producerInfo.getProducerId(), protonReceiver.getName());
185
186        try {
187            Target target = (Target) remoteTarget;
188            ActiveMQDestination destination = null;
189            String targetNodeName = target.getAddress();
190
191            if (target.getDynamic()) {
192                destination = connection.createTemporaryDestination(protonReceiver, target.getCapabilities());
193
194                Map<Symbol, Object> dynamicNodeProperties = new HashMap<>();
195                dynamicNodeProperties.put(LIFETIME_POLICY, DeleteOnClose.getInstance());
196
197                // Currently we only support temporary destinations with delete on close lifetime policy.
198                Target actualTarget = new Target();
199                actualTarget.setAddress(destination.getQualifiedName());
200                actualTarget.setCapabilities(AmqpSupport.getDestinationTypeSymbol(destination));
201                actualTarget.setDynamic(true);
202                actualTarget.setDynamicNodeProperties(dynamicNodeProperties);
203
204                protonReceiver.setTarget(actualTarget);
205                receiver.addCloseAction(new Runnable() {
206
207                    @Override
208                    public void run() {
209                        connection.deleteTemporaryDestination((ActiveMQTempDestination) receiver.getDestination());
210                    }
211                });
212            } else if (targetNodeName != null && !targetNodeName.isEmpty()) {
213                destination = createDestination(remoteTarget);
214                if (destination.isTemporary()) {
215                    String connectionId = ((ActiveMQTempDestination) destination).getConnectionId();
216                    if (connectionId == null) {
217                        throw new AmqpProtocolException(AmqpError.PRECONDITION_FAILED.toString(), "Not a broker created temp destination");
218                    }
219                }
220            }
221
222            Symbol[] remoteDesiredCapabilities = protonReceiver.getRemoteDesiredCapabilities();
223            if (remoteDesiredCapabilities != null) {
224                List<Symbol> list = Arrays.asList(remoteDesiredCapabilities);
225                if (list.contains(AmqpSupport.DELAYED_DELIVERY)) {
226                    protonReceiver.setOfferedCapabilities(new Symbol[] { AmqpSupport.DELAYED_DELIVERY });
227                }
228            }
229
230            receiver.setDestination(destination);
231            connection.sendToActiveMQ(producerInfo, new ResponseHandler() {
232                @Override
233                public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
234                    if (response.isException()) {
235                        ErrorCondition error = null;
236                        Throwable exception = ((ExceptionResponse) response).getException();
237                        if (exception instanceof SecurityException) {
238                            error = new ErrorCondition(AmqpError.UNAUTHORIZED_ACCESS, exception.getMessage());
239                        } else {
240                            error = new ErrorCondition(AmqpError.INTERNAL_ERROR, exception.getMessage());
241                        }
242
243                        receiver.close(error);
244                    } else {
245                        receiver.flow(connection.getConfiguredReceiverCredit());
246                        receiver.open();
247                    }
248                    pumpProtonToSocket();
249                }
250            });
251
252        } catch (AmqpProtocolException exception) {
253            receiver.close(new ErrorCondition(Symbol.getSymbol(exception.getSymbolicName()), exception.getMessage()));
254        }
255    }
256
257    @SuppressWarnings("unchecked")
258    public void createSender(final Sender protonSender) throws Exception {
259        org.apache.qpid.proton.amqp.messaging.Source source = (org.apache.qpid.proton.amqp.messaging.Source) protonSender.getRemoteSource();
260
261        ConsumerInfo consumerInfo = new ConsumerInfo(getNextConsumerId());
262        final AmqpSender sender = new AmqpSender(this, protonSender, consumerInfo);
263
264        LOG.debug("opening new sender {} on link: {}", consumerInfo.getConsumerId(), protonSender.getName());
265
266        try {
267            final Map<Symbol, Object> supportedFilters = new HashMap<>();
268            protonSender.setContext(sender);
269
270            boolean noLocal = false;
271            String selector = null;
272
273            if (source != null) {
274                Map.Entry<Symbol, DescribedType> filter = findFilter(source.getFilter(), JMS_SELECTOR_FILTER_IDS);
275                if (filter != null) {
276                    selector = filter.getValue().getDescribed().toString();
277                    // Validate the Selector.
278                    try {
279                        SelectorParser.parse(selector);
280                    } catch (InvalidSelectorException e) {
281                        sender.close(new ErrorCondition(AmqpError.INVALID_FIELD, e.getMessage()));
282                        return;
283                    }
284
285                    supportedFilters.put(filter.getKey(), filter.getValue());
286                }
287
288                filter = findFilter(source.getFilter(), NO_LOCAL_FILTER_IDS);
289                if (filter != null) {
290                    noLocal = true;
291                    supportedFilters.put(filter.getKey(), filter.getValue());
292                }
293            }
294
295            ActiveMQDestination destination;
296            if (source == null) {
297                // Attempt to recover previous subscription
298                ConsumerInfo storedInfo = connection.lookupSubscription(protonSender.getName());
299
300                if (storedInfo != null) {
301                    destination = storedInfo.getDestination();
302
303                    source = new org.apache.qpid.proton.amqp.messaging.Source();
304                    source.setAddress(destination.getQualifiedName());
305                    source.setDurable(TerminusDurability.UNSETTLED_STATE);
306                    source.setExpiryPolicy(TerminusExpiryPolicy.NEVER);
307                    source.setDistributionMode(COPY);
308
309                    if (storedInfo.isNoLocal()) {
310                        supportedFilters.put(NO_LOCAL_NAME, AmqpNoLocalFilter.NO_LOCAL);
311                    }
312
313                    if (storedInfo.getSelector() != null && !storedInfo.getSelector().trim().equals("")) {
314                        supportedFilters.put(JMS_SELECTOR_NAME, new AmqpJmsSelectorFilter(storedInfo.getSelector()));
315                    }
316                } else {
317                    sender.close(new ErrorCondition(AmqpError.NOT_FOUND, "Unknown subscription link: " + protonSender.getName()));
318                    return;
319                }
320            } else if (source.getDynamic()) {
321                destination = connection.createTemporaryDestination(protonSender, source.getCapabilities());
322
323                Map<Symbol, Object> dynamicNodeProperties = new HashMap<>();
324                dynamicNodeProperties.put(LIFETIME_POLICY, DeleteOnClose.getInstance());
325
326                // Currently we only support temporary destinations with delete on close lifetime policy.
327                source = new org.apache.qpid.proton.amqp.messaging.Source();
328                source.setAddress(destination.getQualifiedName());
329                source.setCapabilities(AmqpSupport.getDestinationTypeSymbol(destination));
330                source.setDynamic(true);
331                source.setDynamicNodeProperties(dynamicNodeProperties);
332
333                sender.addCloseAction(new Runnable() {
334
335                    @Override
336                    public void run() {
337                        connection.deleteTemporaryDestination((ActiveMQTempDestination) sender.getDestination());
338                    }
339                });
340            } else {
341                destination = createDestination(source);
342                if (destination.isTemporary()) {
343                    String connectionId = ((ActiveMQTempDestination) destination).getConnectionId();
344                    if (connectionId == null) {
345                        throw new AmqpProtocolException(AmqpError.INVALID_FIELD.toString(), "Not a broker created temp destination");
346                    }
347                }
348            }
349
350            source.setFilter(supportedFilters.isEmpty() ? null : supportedFilters);
351            protonSender.setSource(source);
352
353            int senderCredit = protonSender.getRemoteCredit();
354
355            // Allows the options on the destination to configure the consumerInfo
356            if (destination.getOptions() != null) {
357                Map<String, Object> options = IntrospectionSupport.extractProperties(
358                    new HashMap<String, Object>(destination.getOptions()), "consumer.");
359                IntrospectionSupport.setProperties(consumerInfo, options);
360                if (options.size() > 0) {
361                    String msg = "There are " + options.size()
362                        + " consumer options that couldn't be set on the consumer."
363                        + " Check the options are spelled correctly."
364                        + " Unknown parameters=[" + options + "]."
365                        + " This consumer cannot be started.";
366                    LOG.warn(msg);
367                    throw new AmqpProtocolException(AmqpError.INVALID_FIELD.toString(), msg);
368                }
369            }
370
371            consumerInfo.setSelector(selector);
372            consumerInfo.setNoRangeAcks(true);
373            consumerInfo.setDestination(destination);
374            consumerInfo.setPrefetchSize(senderCredit >= 0 ? senderCredit : 0);
375            consumerInfo.setDispatchAsync(true);
376            consumerInfo.setNoLocal(noLocal);
377
378            if (source.getDistributionMode() == COPY && destination.isQueue()) {
379                consumerInfo.setBrowser(true);
380            }
381
382            if ((TerminusDurability.UNSETTLED_STATE.equals(source.getDurable()) ||
383                 TerminusDurability.CONFIGURATION.equals(source.getDurable())) && destination.isTopic()) {
384                consumerInfo.setSubscriptionName(protonSender.getName());
385            }
386
387            connection.sendToActiveMQ(consumerInfo, new ResponseHandler() {
388                @Override
389                public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
390                    if (response.isException()) {
391                        ErrorCondition error = null;
392                        Throwable exception = ((ExceptionResponse) response).getException();
393                        if (exception instanceof SecurityException) {
394                            error = new ErrorCondition(AmqpError.UNAUTHORIZED_ACCESS, exception.getMessage());
395                        } else if (exception instanceof InvalidSelectorException) {
396                            error = new ErrorCondition(AmqpError.INVALID_FIELD, exception.getMessage());
397                        } else {
398                            error = new ErrorCondition(AmqpError.INTERNAL_ERROR, exception.getMessage());
399                        }
400
401                        sender.close(error);
402                    } else {
403                        sender.open();
404                    }
405                    pumpProtonToSocket();
406                }
407            });
408
409        } catch (AmqpProtocolException e) {
410            sender.close(new ErrorCondition(Symbol.getSymbol(e.getSymbolicName()), e.getMessage()));
411        }
412    }
413
414    /**
415     * Send all pending work out to the remote peer.
416     */
417    public void pumpProtonToSocket() {
418        connection.pumpProtonToSocket();
419    }
420
421    public void registerSender(ConsumerId consumerId, AmqpSender sender) {
422        consumers.put(consumerId, sender);
423        connection.registerSender(consumerId, sender);
424    }
425
426    public void unregisterSender(ConsumerId consumerId) {
427        consumers.remove(consumerId);
428        connection.unregisterSender(consumerId);
429    }
430
431    public void enlist(TransactionId txId) {
432        if (!enlisted) {
433            connection.getTxCoordinator(txId).enlist(this);
434            enlisted = true;
435        }
436    }
437
438    //----- Configuration accessors ------------------------------------------//
439
440    public AmqpConnection getConnection() {
441        return connection;
442    }
443
444    public SessionId getSessionId() {
445        return sessionId;
446    }
447
448    public Session getEndpoint() {
449        return protonSession;
450    }
451
452    public long getMaxFrameSize() {
453        return connection.getMaxFrameSize();
454    }
455
456    //----- Internal Implementation ------------------------------------------//
457
458    private ConsumerId getNextConsumerId() {
459        return new ConsumerId(sessionId, nextConsumerId++);
460    }
461
462    private ProducerId getNextProducerId() {
463        return new ProducerId(sessionId, nextProducerId++);
464    }
465}