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.toLong;
020
021import java.io.IOException;
022import java.util.LinkedList;
023import java.util.concurrent.atomic.AtomicInteger;
024
025import org.apache.activemq.broker.region.AbstractSubscription;
026import org.apache.activemq.command.ActiveMQDestination;
027import org.apache.activemq.command.ActiveMQMessage;
028import org.apache.activemq.command.ConsumerControl;
029import org.apache.activemq.command.ConsumerId;
030import org.apache.activemq.command.ConsumerInfo;
031import org.apache.activemq.command.ExceptionResponse;
032import org.apache.activemq.command.LocalTransactionId;
033import org.apache.activemq.command.MessageAck;
034import org.apache.activemq.command.MessageDispatch;
035import org.apache.activemq.command.MessagePull;
036import org.apache.activemq.command.RemoveInfo;
037import org.apache.activemq.command.RemoveSubscriptionInfo;
038import org.apache.activemq.command.Response;
039import org.apache.activemq.command.TransactionId;
040import org.apache.activemq.transport.amqp.AmqpProtocolConverter;
041import org.apache.activemq.transport.amqp.message.AutoOutboundTransformer;
042import org.apache.activemq.transport.amqp.message.EncodedMessage;
043import org.apache.activemq.transport.amqp.message.OutboundTransformer;
044import org.apache.qpid.proton.amqp.messaging.Accepted;
045import org.apache.qpid.proton.amqp.messaging.Modified;
046import org.apache.qpid.proton.amqp.messaging.Outcome;
047import org.apache.qpid.proton.amqp.messaging.Rejected;
048import org.apache.qpid.proton.amqp.messaging.Released;
049import org.apache.qpid.proton.amqp.transaction.TransactionalState;
050import org.apache.qpid.proton.amqp.transport.AmqpError;
051import org.apache.qpid.proton.amqp.transport.DeliveryState;
052import org.apache.qpid.proton.amqp.transport.ErrorCondition;
053import org.apache.qpid.proton.amqp.transport.ReceiverSettleMode;
054import org.apache.qpid.proton.amqp.transport.SenderSettleMode;
055import org.apache.qpid.proton.engine.Delivery;
056import org.apache.qpid.proton.engine.Link;
057import org.apache.qpid.proton.engine.Sender;
058import org.fusesource.hawtbuf.Buffer;
059import org.slf4j.Logger;
060import org.slf4j.LoggerFactory;
061
062/**
063 * An AmqpSender wraps the AMQP Sender end of a link from the remote peer
064 * which holds the corresponding Receiver which receives messages transfered
065 * across the link from the Broker.
066 *
067 * An AmqpSender is in turn a message consumer subscribed to some destination
068 * on the broker.  As messages are dispatched to this sender that are sent on
069 * to the remote Receiver end of the lin.
070 */
071public class AmqpSender extends AmqpAbstractLink<Sender> {
072
073    private static final Logger LOG = LoggerFactory.getLogger(AmqpSender.class);
074
075    private static final byte[] EMPTY_BYTE_ARRAY = new byte[] {};
076
077    private final OutboundTransformer outboundTransformer = new AutoOutboundTransformer();
078    private final AmqpTransferTagGenerator tagCache = new AmqpTransferTagGenerator();
079    private final LinkedList<MessageDispatch> outbound = new LinkedList<>();
080    private final LinkedList<Delivery> dispatchedInTx = new LinkedList<>();
081
082    private final ConsumerInfo consumerInfo;
083    private AbstractSubscription subscription;
084    private AtomicInteger prefetchExtension;
085    private int currentCreditRequest;
086    private int logicalDeliveryCount; // echoes prefetch extension but from protons perspective
087    private final boolean presettle;
088
089    private boolean draining;
090    private long lastDeliveredSequenceId;
091
092    private Buffer currentBuffer;
093    private Delivery currentDelivery;
094
095    /**
096     * Creates a new AmqpSender instance that manages the given Sender
097     *
098     * @param session
099     *        the AmqpSession object that is the parent of this instance.
100     * @param endpoint
101     *        the AMQP Sender instance that this class manages.
102     * @param consumerInfo
103     *        the ConsumerInfo instance that holds configuration for this sender.
104     */
105    public AmqpSender(AmqpSession session, Sender endpoint, ConsumerInfo consumerInfo) {
106        super(session, endpoint);
107
108        // We don't support second so enforce it as First and let remote decide what to do
109        this.endpoint.setReceiverSettleMode(ReceiverSettleMode.FIRST);
110
111        // Match what the sender mode is
112        this.endpoint.setSenderSettleMode(endpoint.getRemoteSenderSettleMode());
113
114        this.consumerInfo = consumerInfo;
115        this.presettle = getEndpoint().getSenderSettleMode() == SenderSettleMode.SETTLED;
116    }
117
118    @Override
119    public void open() {
120        if (!isClosed()) {
121            session.registerSender(getConsumerId(), this);
122            subscription = (AbstractSubscription)session.getConnection().lookupPrefetchSubscription(consumerInfo);
123            prefetchExtension = subscription.getPrefetchExtension();
124        }
125
126        super.open();
127    }
128
129    @Override
130    public void detach() {
131        if (!isClosed() && isOpened()) {
132            RemoveInfo removeCommand = new RemoveInfo(getConsumerId());
133            removeCommand.setLastDeliveredSequenceId(lastDeliveredSequenceId);
134
135            sendToActiveMQ(removeCommand, new ResponseHandler() {
136
137                @Override
138                public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
139                    session.unregisterSender(getConsumerId());
140                    AmqpSender.super.detach();
141                }
142            });
143        } else {
144            super.detach();
145        }
146    }
147
148    @Override
149    public void close() {
150        if (!isClosed() && isOpened()) {
151            RemoveInfo removeCommand = new RemoveInfo(getConsumerId());
152            removeCommand.setLastDeliveredSequenceId(lastDeliveredSequenceId);
153
154            sendToActiveMQ(removeCommand, new ResponseHandler() {
155
156                @Override
157                public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
158                    if (consumerInfo.isDurable()) {
159                        RemoveSubscriptionInfo rsi = new RemoveSubscriptionInfo();
160                        rsi.setConnectionId(session.getConnection().getConnectionId());
161                        rsi.setSubscriptionName(getEndpoint().getName());
162                        rsi.setClientId(session.getConnection().getClientId());
163
164                        sendToActiveMQ(rsi);
165                    }
166
167                    session.unregisterSender(getConsumerId());
168                    AmqpSender.super.close();
169                }
170            });
171        } else {
172            super.close();
173        }
174    }
175
176    @Override
177    public void flow() throws Exception {
178        Link endpoint = getEndpoint();
179        if (LOG.isTraceEnabled()) {
180            LOG.trace("Flow: draining={}, drain={} credit={}, currentCredit={}, senderDeliveryCount={} - Sub={}",
181                    draining, endpoint.getDrain(),
182                    endpoint.getCredit(), currentCreditRequest, logicalDeliveryCount, subscription);
183        }
184
185        final int endpointCredit = endpoint.getCredit();
186        if (endpoint.getDrain() && !draining) {
187
188            if (endpointCredit > 0) {
189                draining = true;
190
191                // Now request dispatch of the drain amount, we request immediate
192                // timeout and an completion message regardless so that we can know
193                // when we should marked the link as drained.
194                MessagePull pullRequest = new MessagePull();
195                pullRequest.setConsumerId(getConsumerId());
196                pullRequest.setDestination(getDestination());
197                pullRequest.setTimeout(-1);
198                pullRequest.setAlwaysSignalDone(true);
199                pullRequest.setQuantity(endpointCredit);
200
201                LOG.trace("Pull case -> consumer pull request quantity = {}", endpointCredit);
202
203                sendToActiveMQ(pullRequest);
204            } else {
205                LOG.trace("Pull case -> sending any Queued messages and marking drained");
206
207                pumpOutbound();
208                getEndpoint().drained();
209                session.pumpProtonToSocket();
210                currentCreditRequest = 0;
211                logicalDeliveryCount = 0;
212            }
213        } else if (endpointCredit >= 0) {
214
215            if (endpointCredit == 0 && currentCreditRequest != 0) {
216                prefetchExtension.set(0);
217                currentCreditRequest = 0;
218                logicalDeliveryCount = 0;
219                LOG.trace("Flow: credit 0 for sub:" + subscription);
220            } else {
221                int deltaToAdd = endpointCredit;
222                int logicalCredit = currentCreditRequest - logicalDeliveryCount;
223                if (logicalCredit > 0) {
224                    deltaToAdd -= logicalCredit;
225                } else {
226                    // reset delivery counter - dispatch from broker concurrent with credit=0
227                    // flow can go negative
228                    logicalDeliveryCount = 0;
229                }
230
231                if (deltaToAdd > 0) {
232                    currentCreditRequest = prefetchExtension.addAndGet(deltaToAdd);
233                    subscription.wakeupDestinationsForDispatch();
234                    // force dispatch of matched/pending for topics (pending messages accumulate
235                    // in the sub and are dispatched on update of prefetch)
236                    subscription.setPrefetchSize(0);
237                    LOG.trace("Flow: credit addition of {} for sub {}", deltaToAdd, subscription);
238                }
239            }
240        }
241    }
242
243    @Override
244    public void delivery(Delivery delivery) throws Exception {
245        MessageDispatch md = (MessageDispatch) delivery.getContext();
246        DeliveryState state = delivery.getRemoteState();
247
248        if (state instanceof TransactionalState) {
249            TransactionalState txState = (TransactionalState) state;
250            LOG.trace("onDelivery: TX delivery state = {}", state);
251            if (txState.getOutcome() != null) {
252                Outcome outcome = txState.getOutcome();
253                if (outcome instanceof Accepted) {
254                    TransactionId txId = new LocalTransactionId(session.getConnection().getConnectionId(), toLong(txState.getTxnId()));
255
256                    // Store the message sent in this TX we might need to re-send on rollback
257                    // and we need to ACK it on commit.
258                    session.enlist(txId);
259                    dispatchedInTx.addFirst(delivery);
260
261                    if (!delivery.remotelySettled()) {
262                        TransactionalState txAccepted = new TransactionalState();
263                        txAccepted.setOutcome(Accepted.getInstance());
264                        txAccepted.setTxnId(txState.getTxnId());
265
266                        delivery.disposition(txAccepted);
267                    }
268                }
269            }
270        } else {
271            if (state instanceof Accepted) {
272                LOG.trace("onDelivery: accepted state = {}", state);
273                if (!delivery.remotelySettled()) {
274                    delivery.disposition(new Accepted());
275                }
276                settle(delivery, MessageAck.INDIVIDUAL_ACK_TYPE);
277            } else if (state instanceof Rejected) {
278                // Rejection is a terminal outcome, we poison the message for dispatch to
279                // the DLQ.  If a custom redelivery policy is used on the broker the message
280                // can still be redelivered based on the configation of that policy.
281                LOG.trace("onDelivery: Rejected state = {}, message poisoned.", state);
282                settle(delivery, MessageAck.POISON_ACK_TYPE);
283            } else if (state instanceof Released) {
284                LOG.trace("onDelivery: Released state = {}", state);
285                // re-deliver && don't increment the counter.
286                settle(delivery, -1);
287            } else if (state instanceof Modified) {
288                Modified modified = (Modified) state;
289                if (Boolean.TRUE.equals(modified.getDeliveryFailed())) {
290                    // increment delivery counter..
291                    md.setRedeliveryCounter(md.getRedeliveryCounter() + 1);
292                }
293                LOG.trace("onDelivery: Modified state = {}, delivery count now {}", state, md.getRedeliveryCounter());
294                byte ackType = -1;
295                Boolean undeliverableHere = modified.getUndeliverableHere();
296                if (undeliverableHere != null && undeliverableHere) {
297                    // receiver does not want the message..
298                    // perhaps we should DLQ it?
299                    ackType = MessageAck.POISON_ACK_TYPE;
300                }
301                settle(delivery, ackType);
302            }
303        }
304
305        pumpOutbound();
306    }
307
308    @Override
309    public void commit(LocalTransactionId txnId) throws Exception {
310        if (!dispatchedInTx.isEmpty()) {
311            for (final Delivery delivery : dispatchedInTx) {
312                MessageDispatch dispatch = (MessageDispatch) delivery.getContext();
313
314                MessageAck pendingTxAck = new MessageAck(dispatch, MessageAck.INDIVIDUAL_ACK_TYPE, 1);
315                pendingTxAck.setFirstMessageId(dispatch.getMessage().getMessageId());
316                pendingTxAck.setTransactionId(txnId);
317
318                LOG.trace("Sending commit Ack to ActiveMQ: {}", pendingTxAck);
319
320                sendToActiveMQ(pendingTxAck, new ResponseHandler() {
321                    @Override
322                    public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
323                        if (response.isException()) {
324                            Throwable exception = ((ExceptionResponse) response).getException();
325                            exception.printStackTrace();
326                            getEndpoint().close();
327                        } else {
328                            delivery.settle();
329                        }
330                        session.pumpProtonToSocket();
331                    }
332                });
333            }
334
335            dispatchedInTx.clear();
336        }
337    }
338
339    @Override
340    public void rollback(LocalTransactionId txnId) throws Exception {
341        synchronized (outbound) {
342
343            LOG.trace("Rolling back {} messages for redelivery. ", dispatchedInTx.size());
344
345            for (Delivery delivery : dispatchedInTx) {
346                // Only settled deliveries should be re-dispatched, unsettled deliveries
347                // remain acquired on the remote end and can be accepted again in a new
348                // TX or released or rejected etc.
349                MessageDispatch dispatch = (MessageDispatch) delivery.getContext();
350                dispatch.getMessage().setTransactionId(null);
351
352                if (delivery.remotelySettled()) {
353                    dispatch.setRedeliveryCounter(dispatch.getRedeliveryCounter() + 1);
354                    outbound.addFirst(dispatch);
355                }
356            }
357
358            dispatchedInTx.clear();
359        }
360    }
361
362    /**
363     * Event point for incoming message from ActiveMQ on this Sender's
364     * corresponding subscription.
365     *
366     * @param dispatch
367     *        the MessageDispatch to process and send across the link.
368     *
369     * @throws Exception if an error occurs while encoding the message for send.
370     */
371    public void onMessageDispatch(MessageDispatch dispatch) throws Exception {
372        if (!isClosed()) {
373            // Lock to prevent stepping on TX redelivery
374            synchronized (outbound) {
375                outbound.addLast(dispatch);
376            }
377            pumpOutbound();
378            session.pumpProtonToSocket();
379        }
380    }
381
382    /**
383     * Called when the Broker sends a ConsumerControl command to the Consumer that
384     * this sender creates to obtain messages to dispatch via the sender for this
385     * end of the open link.
386     *
387     * @param control
388     *        The ConsumerControl command to process.
389     */
390    public void onConsumerControl(ConsumerControl control) {
391        if (control.isClose()) {
392            close(new ErrorCondition(AmqpError.INTERNAL_ERROR, "Receiver forcably closed"));
393            session.pumpProtonToSocket();
394        }
395    }
396
397    @Override
398    public String toString() {
399        return "AmqpSender {" + getConsumerId() + "}";
400    }
401
402    //----- Property getters and setters -------------------------------------//
403
404    public ConsumerId getConsumerId() {
405        return consumerInfo.getConsumerId();
406    }
407
408    @Override
409    public ActiveMQDestination getDestination() {
410        return consumerInfo.getDestination();
411    }
412
413    @Override
414    public void setDestination(ActiveMQDestination destination) {
415        consumerInfo.setDestination(destination);
416    }
417
418    //----- Internal Implementation ------------------------------------------//
419
420    public void pumpOutbound() throws Exception {
421        while (!isClosed()) {
422            while (currentBuffer != null) {
423                int sent = getEndpoint().send(currentBuffer.data, currentBuffer.offset, currentBuffer.length);
424                if (sent > 0) {
425                    currentBuffer.moveHead(sent);
426                    if (currentBuffer.length == 0) {
427                        if (presettle) {
428                            settle(currentDelivery, MessageAck.INDIVIDUAL_ACK_TYPE);
429                        } else {
430                            getEndpoint().advance();
431                        }
432                        currentBuffer = null;
433                        currentDelivery = null;
434                        logicalDeliveryCount++;
435                    }
436                } else {
437                    return;
438                }
439            }
440
441            if (outbound.isEmpty()) {
442                return;
443            }
444
445            final MessageDispatch md = outbound.removeFirst();
446            try {
447
448                ActiveMQMessage temp = null;
449                if (md.getMessage() != null) {
450                    temp = (ActiveMQMessage) md.getMessage().copy();
451                }
452
453                final ActiveMQMessage jms = temp;
454                if (jms == null) {
455                    LOG.trace("Sender:[{}] browse done.", getEndpoint().getName());
456                    // It's the end of browse signal in response to a MessagePull
457                    getEndpoint().drained();
458                    draining = false;
459                    currentCreditRequest = 0;
460                    logicalDeliveryCount = 0;
461                } else {
462                    if (LOG.isTraceEnabled()) {
463                        LOG.trace("Sender:[{}] msgId={} draining={}, drain={}, credit={}, remoteCredit={}, queued={}",
464                                  getEndpoint().getName(), jms.getJMSMessageID(), draining, getEndpoint().getDrain(),
465                                  getEndpoint().getCredit(), getEndpoint().getRemoteCredit(), getEndpoint().getQueued());
466                    }
467
468                    if (draining && getEndpoint().getCredit() == 0) {
469                        LOG.trace("Sender:[{}] browse complete.", getEndpoint().getName());
470                        getEndpoint().drained();
471                        draining = false;
472                        currentCreditRequest = 0;
473                        logicalDeliveryCount = 0;
474                    }
475
476                    jms.setRedeliveryCounter(md.getRedeliveryCounter());
477                    jms.setReadOnlyBody(true);
478                    final EncodedMessage amqp = outboundTransformer.transform(jms);
479                    if (amqp != null && amqp.getLength() > 0) {
480                        currentBuffer = new Buffer(amqp.getArray(), amqp.getArrayOffset(), amqp.getLength());
481                        if (presettle) {
482                            currentDelivery = getEndpoint().delivery(EMPTY_BYTE_ARRAY, 0, 0);
483                        } else {
484                            final byte[] tag = tagCache.getNextTag();
485                            currentDelivery = getEndpoint().delivery(tag, 0, tag.length);
486                        }
487                        currentDelivery.setContext(md);
488                        currentDelivery.setMessageFormat((int) amqp.getMessageFormat());
489                    } else {
490                        // TODO: message could not be generated what now?
491                    }
492                }
493            } catch (Exception e) {
494                LOG.warn("Error detected while flushing outbound messages: {}", e.getMessage());
495            }
496        }
497    }
498
499    private void settle(final Delivery delivery, final int ackType) throws Exception {
500        byte[] tag = delivery.getTag();
501        if (tag != null && tag.length > 0 && delivery.remotelySettled()) {
502            tagCache.returnTag(tag);
503        }
504
505        if (ackType == -1) {
506            // we are going to settle, but redeliver.. we we won't yet ack to ActiveMQ
507            delivery.settle();
508            onMessageDispatch((MessageDispatch) delivery.getContext());
509        } else {
510            MessageDispatch md = (MessageDispatch) delivery.getContext();
511            lastDeliveredSequenceId = md.getMessage().getMessageId().getBrokerSequenceId();
512            MessageAck ack = new MessageAck();
513            ack.setConsumerId(getConsumerId());
514            ack.setFirstMessageId(md.getMessage().getMessageId());
515            ack.setLastMessageId(md.getMessage().getMessageId());
516            ack.setMessageCount(1);
517            ack.setAckType((byte) ackType);
518            ack.setDestination(md.getDestination());
519            LOG.trace("Sending Ack to ActiveMQ: {}", ack);
520
521            sendToActiveMQ(ack, new ResponseHandler() {
522                @Override
523                public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
524                    if (response.isException()) {
525                        if (response.isException()) {
526                            Throwable exception = ((ExceptionResponse) response).getException();
527                            exception.printStackTrace();
528                            getEndpoint().close();
529                        }
530                    } else {
531                        delivery.settle();
532                    }
533                    session.pumpProtonToSocket();
534                }
535            });
536        }
537    }
538}