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}