diff options
| author | Robert Godfrey <rgodfrey@apache.org> | 2014-02-01 15:40:47 +0000 |
|---|---|---|
| committer | Robert Godfrey <rgodfrey@apache.org> | 2014-02-01 15:40:47 +0000 |
| commit | 6823d23dbeca328f4e860538a52015bc9313a6db (patch) | |
| tree | 8af7823f6e4ac169835909ed6babf38d4ad42e57 /qpid/java/broker-core | |
| parent | ade50f17b8ffea099f8fffaaf283b2412f393bce (diff) | |
| download | qpid-python-6823d23dbeca328f4e860538a52015bc9313a6db.tar.gz | |
QPID-5504 : Moving routing to Exchange from session classes
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk@1563431 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/java/broker-core')
8 files changed, 158 insertions, 176 deletions
diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java index b00d98637e..6a959df440 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/AbstractExchange.java @@ -33,12 +33,14 @@ import org.apache.qpid.server.logging.messages.ExchangeMessages; import org.apache.qpid.server.logging.subjects.BindingLogSubject; import org.apache.qpid.server.logging.subjects.ExchangeLogSubject; import org.apache.qpid.server.message.InstanceProperties; +import org.apache.qpid.server.message.MessageReference; import org.apache.qpid.server.message.ServerMessage; import org.apache.qpid.server.model.UUIDGenerator; import org.apache.qpid.server.plugin.ExchangeType; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.queue.BaseQueue; import org.apache.qpid.server.store.DurableConfigurationStoreHelper; +import org.apache.qpid.server.txn.ServerTransaction; import org.apache.qpid.server.virtualhost.VirtualHost; import java.util.Collection; @@ -374,9 +376,9 @@ public abstract class AbstractExchange implements Exchange return getBindings().size(); } - @Override - public final List<? extends BaseQueue> route(final ServerMessage message, - final InstanceProperties instanceProperties) + + final List<? extends BaseQueue> route(final ServerMessage message, + final InstanceProperties instanceProperties) { _receivedMessageCount.incrementAndGet(); _receivedMessageSize.addAndGet(message.getSize()); @@ -416,6 +418,59 @@ public abstract class AbstractExchange implements Exchange return queues; } + public final int send(final ServerMessage message, + final InstanceProperties instanceProperties, + final ServerTransaction txn, + final BaseQueue.PostEnqueueAction postEnqueueAction) + { + List<? extends BaseQueue> queues = route(message, instanceProperties); + + if(queues == null || queues.isEmpty()) + { + Exchange altExchange = getAlternateExchange(); + if(altExchange != null) + { + return altExchange.send(message, instanceProperties, txn, postEnqueueAction); + } + else + { + return 0; + } + } + else + { + final BaseQueue[] baseQueues = queues.toArray(new BaseQueue[queues.size()]); + + txn.enqueue(queues,message, new ServerTransaction.Action() + { + MessageReference _reference = message.newReference(); + + public void postCommit() + { + for(int i = 0; i < baseQueues.length; i++) + { + try + { + baseQueues[i].enqueue(message, postEnqueueAction); + } + catch (AMQException e) + { + // TODO + throw new RuntimeException(e); + } + } + _reference.release(); + } + + public void onRollback() + { + _reference.release(); + } + }); + return queues.size(); + } + } + protected abstract List<? extends BaseQueue> doRoute(final ServerMessage message, final InstanceProperties instanceProperties); @@ -679,4 +734,6 @@ public abstract class AbstractExchange implements Exchange public void onClose(Exchange exchange) throws AMQSecurityException, AMQInternalException; } + + } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/DefaultExchange.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/DefaultExchange.java index e2582019cd..71d0f8b4dd 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/DefaultExchange.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/DefaultExchange.java @@ -36,11 +36,14 @@ import org.apache.qpid.server.logging.LogSubject; import org.apache.qpid.server.logging.actors.CurrentActor; import org.apache.qpid.server.logging.messages.ExchangeMessages; import org.apache.qpid.server.message.InstanceProperties; +import org.apache.qpid.server.message.MessageReference; import org.apache.qpid.server.message.ServerMessage; import org.apache.qpid.server.model.UUIDGenerator; import org.apache.qpid.server.plugin.ExchangeType; import org.apache.qpid.server.queue.AMQQueue; +import org.apache.qpid.server.queue.BaseQueue; import org.apache.qpid.server.queue.QueueRegistry; +import org.apache.qpid.server.txn.ServerTransaction; import org.apache.qpid.server.virtualhost.VirtualHost; public class DefaultExchange implements Exchange @@ -204,22 +207,6 @@ public class DefaultExchange implements Exchange } @Override - public List<AMQQueue> route(ServerMessage message, final InstanceProperties instanceProperties) - { - AMQQueue q = _virtualHost.getQueue(message.getRoutingKey()); - if(q == null) - { - List<AMQQueue> noQueues = Collections.emptyList(); - return noQueues; - } - else - { - return Collections.singletonList(q); - } - - } - - @Override public boolean isBound(AMQQueue queue) { return _virtualHost.getQueue(queue.getName()) == queue; @@ -343,4 +330,47 @@ public class DefaultExchange implements Exchange { return _id; } + + public final int send(final ServerMessage message, + final InstanceProperties instanceProperties, + final ServerTransaction txn, + final BaseQueue.PostEnqueueAction postEnqueueAction) + { + final AMQQueue q = _virtualHost.getQueue(message.getRoutingKey()); + if(q == null) + { + return 0; + } + else + { + txn.enqueue(q,message, new ServerTransaction.Action() + { + MessageReference _reference = message.newReference(); + + public void postCommit() + { + try + { + q.enqueue(message, postEnqueueAction); + } + catch (AMQException e) + { + // TODO + throw new RuntimeException(e); + } + finally + { + _reference.release(); + } + } + + public void onRollback() + { + _reference.release(); + } + }); + return 1; + } + } + } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/Exchange.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/Exchange.java index 78455c9261..18e912e972 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/Exchange.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/exchange/Exchange.java @@ -29,6 +29,7 @@ import org.apache.qpid.server.message.ServerMessage; import org.apache.qpid.server.plugin.ExchangeType; import org.apache.qpid.server.queue.AMQQueue; import org.apache.qpid.server.queue.BaseQueue; +import org.apache.qpid.server.txn.ServerTransaction; import org.apache.qpid.server.virtualhost.VirtualHost; import java.util.Collection; @@ -94,13 +95,17 @@ public interface Exchange extends ExchangeReferrer void close() throws AMQException; /** - * Returns a list of queues to which to route this message. If there are - * no queues the empty list must be returned. - * - * @return list of queues to which to route the message. + * Routes a message + * @param message the message to be routed + * @param instanceProperties the instance properties + * @param txn the transaction to enqueue within + * @param postEnqueueAction action to perform on the result of every enqueue (may be null) + * @return the number of queues in which the message was enqueued performed */ - List<? extends BaseQueue> route(ServerMessage message, final InstanceProperties instanceProperties); - + int send(ServerMessage message, + InstanceProperties instanceProperties, + ServerTransaction txn, + BaseQueue.PostEnqueueAction postEnqueueAction); /** * Determines whether a message would be isBound to a particular queue using a specific routing key and arguments diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueEntry.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueEntry.java index 80ccbe1649..2aa1d1f473 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueEntry.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueEntry.java @@ -22,11 +22,11 @@ package org.apache.qpid.server.queue; import org.apache.qpid.AMQException; import org.apache.qpid.server.filter.Filterable; -import org.apache.qpid.server.message.InstanceProperties; -import org.apache.qpid.server.message.ServerMessage; +import org.apache.qpid.server.message.MessageInstance; import org.apache.qpid.server.subscription.Subscription; +import org.apache.qpid.server.txn.ServerTransaction; -public interface QueueEntry extends Comparable<QueueEntry> +public interface QueueEntry extends MessageInstance, Comparable<QueueEntry> { @@ -177,26 +177,17 @@ public interface QueueEntry extends Comparable<QueueEntry> AMQQueue getQueue(); - ServerMessage getMessage(); - long getSize(); boolean getDeliveredToConsumer(); boolean expired() throws AMQException; - boolean isAvailable(); - - boolean isAcquired(); - - boolean acquire(); boolean acquire(Subscription sub); boolean acquiredBySubscription(); boolean isAcquiredBy(Subscription subscription); - void release(); - void setRedelivered(); boolean isRedelivered(); @@ -207,16 +198,7 @@ public interface QueueEntry extends Comparable<QueueEntry> boolean isRejectedBy(long subscriptionId); - void delete(); - - /** - * Returns true if entry is either DEQUED or DELETED state. - * - * @return true if entry is either DEQUED or DELETED state - */ - boolean isDeleted(); - - void routeToAlternate(); + int routeToAlternate(final BaseQueue.PostEnqueueAction action, ServerTransaction txn); boolean isQueueDeleted(); @@ -241,5 +223,4 @@ public interface QueueEntry extends Comparable<QueueEntry> Filterable asFilterable(); - InstanceProperties getInstanceProperties(); } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueEntryImpl.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueEntryImpl.java index ed61f1acf6..461d493437 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueEntryImpl.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/QueueEntryImpl.java @@ -34,7 +34,6 @@ import org.apache.qpid.server.txn.ServerTransaction; import java.util.EnumMap; import java.util.HashSet; -import java.util.List; import java.util.Set; import java.util.concurrent.CopyOnWriteArraySet; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; @@ -250,7 +249,7 @@ public abstract class QueueEntryImpl implements QueueEntry } else if(acquire()) { - routeToAlternate(); + routeToAlternate(null, null); } } @@ -368,65 +367,43 @@ public abstract class QueueEntryImpl implements QueueEntry dispose(); } - public void routeToAlternate() + public int routeToAlternate(final BaseQueue.PostEnqueueAction action, ServerTransaction txn) { final AMQQueue currentQueue = getQueue(); Exchange alternateExchange = currentQueue.getAlternateExchange(); - + boolean autocommit = txn == null; if (alternateExchange != null) { - List<? extends BaseQueue> queues = alternateExchange.route(getMessage(), getInstanceProperties()); - final ServerMessage message = getMessage(); - if ((queues == null || queues.size() == 0) && alternateExchange.getAlternateExchange() != null) + if(autocommit) { - queues = alternateExchange.getAlternateExchange().route(getMessage(), getInstanceProperties()); + txn = new LocalTransaction(getQueue().getVirtualHost().getMessageStore()); } + int enqueues = alternateExchange.send(getMessage(), getInstanceProperties(), txn, action); - - if (queues != null && queues.size() != 0) + txn.dequeue(currentQueue, getMessage(), new ServerTransaction.Action() { - final List<? extends BaseQueue> rerouteQueues = queues; - ServerTransaction txn = new LocalTransaction(getQueue().getVirtualHost().getMessageStore()); - - txn.enqueue(rerouteQueues, message, new ServerTransaction.Action() + public void postCommit() { - public void postCommit() - { - try - { - for (BaseQueue queue : rerouteQueues) - { - queue.enqueue(message); - } - } - catch (AMQException e) - { - throw new RuntimeException(e); - } - } - - public void onRollback() - { - - } - }); - - txn.dequeue(currentQueue, message, new ServerTransaction.Action() - { - public void postCommit() - { - delete(); - } + delete(); + } - public void onRollback() - { + public void onRollback() + { - } - }); + } + }); + if(autocommit) + { txn.commit(); } + return enqueues; + + } + else + { + return 0; } } diff --git a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java index d63d1946d3..87d11a892e 100644 --- a/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java +++ b/qpid/java/broker-core/src/main/java/org/apache/qpid/server/queue/SimpleAMQQueue.java @@ -1355,93 +1355,25 @@ public class SimpleAMQQueue implements AMQQueue, Subscription.StateListener, Mes ServerTransaction txn = new LocalTransaction(getVirtualHost().getMessageStore()); - if(_alternateExchange != null) + + for(final QueueEntry entry : entries) { + // TODO log requeues with a post enqueue action + int requeues = entry.routeToAlternate(null, txn); - for(final QueueEntry entry : entries) + if(requeues == 0) { - - List<? extends BaseQueue> queues = _alternateExchange.route(entry.getMessage(), entry.getInstanceProperties()); - if((queues == null || queues.size() == 0) && _alternateExchange.getAlternateExchange() != null) - { - queues = _alternateExchange.getAlternateExchange().route(entry.getMessage(), entry.getInstanceProperties()); - } - - final ServerMessage message = entry.getMessage(); - if(queues != null && queues.size() != 0) - { - final List<? extends BaseQueue> rerouteQueues = queues; - txn.enqueue(rerouteQueues, entry.getMessage(), - new ServerTransaction.Action() - { - - public void postCommit() - { - try - { - for(BaseQueue queue : rerouteQueues) - { - queue.enqueue(message); - } - } - catch (AMQException e) - { - throw new RuntimeException(e); - } - - } - - public void onRollback() - { - - } - }); - txn.dequeue(this, entry.getMessage(), - new ServerTransaction.Action() - { - - public void postCommit() - { - entry.delete(); - } - - public void onRollback() - { - } - }); - } - + // TODO log discard } - - _alternateExchange.removeReference(this); } - else - { - // TODO log discard - - for(final QueueEntry entry : entries) - { - final ServerMessage message = entry.getMessage(); - if(message != null) - { - txn.dequeue(this, message, - new ServerTransaction.Action() - { - public void postCommit() - { - entry.delete(); - } + txn.commit(); - public void onRollback() - { - } - }); - } - } + if(_alternateExchange != null) + { + _alternateExchange.removeReference(this); } - txn.commit(); for (Task task : _deleteTaskList) { diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/TopicExchangeTest.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/TopicExchangeTest.java index 7bd525c90f..764549626a 100644 --- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/TopicExchangeTest.java +++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/exchange/TopicExchangeTest.java @@ -312,10 +312,9 @@ public class TopicExchangeTest extends QpidTestCase private int routeMessage(String routingKey, long messageNumber) throws AMQException { - ServerMessage serverMessage = mock(ServerMessage.class); - when(serverMessage.getRoutingKey()).thenReturn(routingKey); - List<? extends BaseQueue> queues = _exchange.route(serverMessage, InstanceProperties.EMPTY); ServerMessage message = mock(ServerMessage.class); + when(message.getRoutingKey()).thenReturn(routingKey); + List<? extends BaseQueue> queues = _exchange.route(message, InstanceProperties.EMPTY); MessageReference ref = mock(MessageReference.class); when(ref.getMessage()).thenReturn(message); when(message.newReference()).thenReturn(ref); diff --git a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/queue/MockQueueEntry.java b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/queue/MockQueueEntry.java index 2e3231e208..d3c866f747 100644 --- a/qpid/java/broker-core/src/test/java/org/apache/qpid/server/queue/MockQueueEntry.java +++ b/qpid/java/broker-core/src/test/java/org/apache/qpid/server/queue/MockQueueEntry.java @@ -26,6 +26,7 @@ import org.apache.qpid.server.message.AMQMessageHeader; import org.apache.qpid.server.message.InstanceProperties; import org.apache.qpid.server.message.ServerMessage; import org.apache.qpid.server.subscription.Subscription; +import org.apache.qpid.server.txn.ServerTransaction; public class MockQueueEntry implements QueueEntry { @@ -62,9 +63,9 @@ public class MockQueueEntry implements QueueEntry } - public void routeToAlternate() + public int routeToAlternate(final BaseQueue.PostEnqueueAction action, final ServerTransaction txn) { - + return 0; } public boolean expired() throws AMQException |
