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}