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}