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;
022
023import javax.jms.Destination;
024import javax.jms.ResourceAllocationException;
025
026import org.apache.activemq.command.ActiveMQDestination;
027import org.apache.activemq.command.ActiveMQMessage;
028import org.apache.activemq.command.ExceptionResponse;
029import org.apache.activemq.command.LocalTransactionId;
030import org.apache.activemq.command.MessageId;
031import org.apache.activemq.command.ProducerId;
032import org.apache.activemq.command.ProducerInfo;
033import org.apache.activemq.command.RemoveInfo;
034import org.apache.activemq.command.Response;
035import org.apache.activemq.command.TransactionId;
036import org.apache.activemq.transport.amqp.AmqpProtocolConverter;
037import org.apache.activemq.transport.amqp.message.AMQPNativeInboundTransformer;
038import org.apache.activemq.transport.amqp.message.AMQPRawInboundTransformer;
039import org.apache.activemq.transport.amqp.message.EncodedMessage;
040import org.apache.activemq.transport.amqp.message.InboundTransformer;
041import org.apache.activemq.transport.amqp.message.JMSMappingInboundTransformer;
042import org.apache.activemq.util.LongSequenceGenerator;
043import org.apache.qpid.proton.amqp.Symbol;
044import org.apache.qpid.proton.amqp.messaging.Accepted;
045import org.apache.qpid.proton.amqp.messaging.Rejected;
046import org.apache.qpid.proton.amqp.transaction.TransactionalState;
047import org.apache.qpid.proton.amqp.transport.AmqpError;
048import org.apache.qpid.proton.amqp.transport.DeliveryState;
049import org.apache.qpid.proton.amqp.transport.ErrorCondition;
050import org.apache.qpid.proton.engine.Delivery;
051import org.apache.qpid.proton.engine.Receiver;
052import org.fusesource.hawtbuf.Buffer;
053import org.slf4j.Logger;
054import org.slf4j.LoggerFactory;
055
056/**
057 * An AmqpReceiver wraps the AMQP Receiver end of a link from the remote peer
058 * which holds the corresponding Sender which transfers message accross the
059 * link.  The AmqpReceiver handles all incoming deliveries by converting them
060 * or wrapping them into an ActiveMQ message object and forwarding that message
061 * on to the appropriate ActiveMQ Destination.
062 */
063public class AmqpReceiver extends AmqpAbstractReceiver {
064
065    private static final Logger LOG = LoggerFactory.getLogger(AmqpReceiver.class);
066
067    private final ProducerInfo producerInfo;
068    private final LongSequenceGenerator messageIdGenerator = new LongSequenceGenerator();
069
070    private InboundTransformer inboundTransformer;
071
072    private int sendsInFlight;
073
074    /**
075     * Create a new instance of an AmqpReceiver
076     *
077     * @param session
078     *        the Session that is the parent of this AmqpReceiver instance.
079     * @param endpoint
080     *        the AMQP receiver endpoint that the class manages.
081     * @param producerInfo
082     *        the ProducerInfo instance that contains this sender's configuration.
083     */
084    public AmqpReceiver(AmqpSession session, Receiver endpoint, ProducerInfo producerInfo) {
085        super(session, endpoint);
086
087        this.producerInfo = producerInfo;
088    }
089
090    @Override
091    public void close() {
092        if (!isClosed() && isOpened()) {
093            sendToActiveMQ(new RemoveInfo(getProducerId()), new ResponseHandler() {
094
095                @Override
096                public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
097                    AmqpReceiver.super.close();
098                }
099            });
100        } else {
101            super.close();
102        }
103    }
104
105    //----- Configuration accessors ------------------------------------------//
106
107    /**
108     * @return the ActiveMQ ProducerId used to register this Receiver on the Broker.
109     */
110    public ProducerId getProducerId() {
111        return producerInfo.getProducerId();
112    }
113
114    @Override
115    public ActiveMQDestination getDestination() {
116        return producerInfo.getDestination();
117    }
118
119    @Override
120    public void setDestination(ActiveMQDestination destination) {
121        producerInfo.setDestination(destination);
122    }
123
124    /**
125     * If the Sender that initiated this Receiver endpoint did not define an address
126     * then it is using anonymous mode and message are to be routed to the address
127     * that is defined in the AMQP message 'To' field.
128     *
129     * @return true if this Receiver should operate in anonymous mode.
130     */
131    public boolean isAnonymous() {
132        return producerInfo.getDestination() == null;
133    }
134
135    //----- Internal Implementation ------------------------------------------//
136
137    protected InboundTransformer getTransformer() {
138        if (inboundTransformer == null) {
139            String transformer = session.getConnection().getConfiguredTransformer();
140            if (transformer.equalsIgnoreCase(InboundTransformer.TRANSFORMER_JMS)) {
141                inboundTransformer = new JMSMappingInboundTransformer();
142            } else if (transformer.equalsIgnoreCase(InboundTransformer.TRANSFORMER_NATIVE)) {
143                inboundTransformer = new AMQPNativeInboundTransformer();
144            } else if (transformer.equalsIgnoreCase(InboundTransformer.TRANSFORMER_RAW)) {
145                inboundTransformer = new AMQPRawInboundTransformer();
146            } else {
147                LOG.warn("Unknown transformer type {} using native one instead", transformer);
148                inboundTransformer = new AMQPNativeInboundTransformer();
149            }
150        }
151        return inboundTransformer;
152    }
153
154    @Override
155    protected void processDelivery(final Delivery delivery, Buffer deliveryBytes) throws Exception {
156        if (!isClosed()) {
157            EncodedMessage em = new EncodedMessage(delivery.getMessageFormat(), deliveryBytes.data, deliveryBytes.offset, deliveryBytes.length);
158
159            InboundTransformer transformer = getTransformer();
160            ActiveMQMessage message = transformer.transform(em);
161
162            current = null;
163
164            if (isAnonymous()) {
165                Destination toDestination = message.getJMSDestination();
166                if (toDestination == null || !(toDestination instanceof ActiveMQDestination)) {
167                    Rejected rejected = new Rejected();
168                    ErrorCondition condition = new ErrorCondition();
169                    condition.setCondition(Symbol.valueOf("failed"));
170                    condition.setDescription("Missing to field for message sent to an anonymous producer");
171                    rejected.setError(condition);
172                    delivery.disposition(rejected);
173                    return;
174                }
175            } else {
176                message.setJMSDestination(getDestination());
177            }
178
179            message.setProducerId(getProducerId());
180
181            // Always override the AMQP client's MessageId with our own.  Preserve
182            // the original in the TextView property for later Ack.
183            MessageId messageId = new MessageId(getProducerId(), messageIdGenerator.getNextSequenceId());
184
185            MessageId amqpMessageId = message.getMessageId();
186            if (amqpMessageId != null) {
187                if (amqpMessageId.getTextView() != null) {
188                    messageId.setTextView(amqpMessageId.getTextView());
189                } else {
190                    messageId.setTextView(amqpMessageId.toString());
191                }
192            }
193
194            message.setMessageId(messageId);
195
196            LOG.trace("Inbound Message:{} from Producer:{}",
197                      message.getMessageId(), getProducerId() + ":" + messageId.getProducerSequenceId());
198
199            final DeliveryState remoteState = delivery.getRemoteState();
200            if (remoteState != null && remoteState instanceof TransactionalState) {
201                TransactionalState txState = (TransactionalState) remoteState;
202                TransactionId txId = new LocalTransactionId(session.getConnection().getConnectionId(), toLong(txState.getTxnId()));
203                session.enlist(txId);
204                message.setTransactionId(txId);
205            }
206
207            message.onSend();
208
209            sendsInFlight++;
210
211            sendToActiveMQ(message, createResponseHandler(delivery));
212        }
213    }
214
215    private ResponseHandler createResponseHandler(final Delivery delivery) {
216        return new ResponseHandler() {
217
218            @Override
219            public void onResponse(AmqpProtocolConverter converter, Response response) throws IOException {
220                if (!delivery.remotelySettled()) {
221                    if (response.isException()) {
222                        ExceptionResponse error = (ExceptionResponse) response;
223                        Rejected rejected = new Rejected();
224                        ErrorCondition condition = new ErrorCondition();
225
226                        if (error.getException() instanceof SecurityException) {
227                            condition.setCondition(AmqpError.UNAUTHORIZED_ACCESS);
228                        } else if (error.getException() instanceof ResourceAllocationException) {
229                            condition.setCondition(AmqpError.RESOURCE_LIMIT_EXCEEDED);
230                        } else {
231                            condition.setCondition(Symbol.valueOf("failed"));
232                        }
233
234                        condition.setDescription(error.getException().getMessage());
235                        rejected.setError(condition);
236                        delivery.disposition(rejected);
237                    } else {
238                        final DeliveryState remoteState = delivery.getRemoteState();
239                        if (remoteState != null && remoteState instanceof TransactionalState) {
240                            TransactionalState txAccepted = new TransactionalState();
241                            txAccepted.setOutcome(Accepted.getInstance());
242                            txAccepted.setTxnId(((TransactionalState) remoteState).getTxnId());
243
244                            delivery.disposition(txAccepted);
245                        } else {
246                            delivery.disposition(Accepted.getInstance());
247                        }
248                    }
249                }
250
251                if (getEndpoint().getCredit() + --sendsInFlight <= (getConfiguredReceiverCredit() * .3)) {
252                    LOG.trace("Sending more credit ({}) to producer: {}", getConfiguredReceiverCredit() * .7, getProducerId());
253                    getEndpoint().flow((int) (getConfiguredReceiverCredit() * .7));
254                }
255
256                delivery.settle();
257                session.pumpProtonToSocket();
258            }
259        };
260    }
261}