diff options
| author | Alan Conway <aconway@apache.org> | 2010-10-27 18:01:27 +0000 |
|---|---|---|
| committer | Alan Conway <aconway@apache.org> | 2010-10-27 18:01:27 +0000 |
| commit | 326dddd0d0d48401d14ca93044b3fc0e35ad87d9 (patch) | |
| tree | 019a45480d8cdf832f62d7176b7a10a5d0971535 /cpp/src/qpid/broker | |
| parent | aae11121cfcf891b2365241141f9ab9cb47d3024 (diff) | |
| download | qpid-python-326dddd0d0d48401d14ca93044b3fc0e35ad87d9.tar.gz | |
Revert experimental cluster code, too close to 0.8 release.
Reverts revisions:
r1023966 "Introduce broker::Cluster interface."
r1024275 "Fix compile error: outline set/getCluster fucntions on Broker."
r1027210 "New cluster: core framework and initial implementation of enqueue logic."
git-svn-id: https://svn.apache.org/repos/asf/qpid/trunk/qpid@1028055 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'cpp/src/qpid/broker')
| -rw-r--r-- | cpp/src/qpid/broker/Broker.cpp | 6 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/Broker.h | 5 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/Cluster.h | 110 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/DeliveryRecord.cpp | 19 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/Exchange.cpp | 16 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/NullCluster.h | 66 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/Queue.cpp | 137 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/Queue.h | 7 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/QueuedMessage.h | 3 | ||||
| -rw-r--r-- | cpp/src/qpid/broker/SemanticState.cpp | 7 |
10 files changed, 56 insertions, 320 deletions
diff --git a/cpp/src/qpid/broker/Broker.cpp b/cpp/src/qpid/broker/Broker.cpp index c93949e33f..33364e48df 100644 --- a/cpp/src/qpid/broker/Broker.cpp +++ b/cpp/src/qpid/broker/Broker.cpp @@ -24,7 +24,6 @@ #include "qpid/broker/FanOutExchange.h" #include "qpid/broker/HeadersExchange.h" #include "qpid/broker/MessageStoreModule.h" -#include "qpid/broker/NullCluster.h" #include "qpid/broker/NullMessageStore.h" #include "qpid/broker/RecoveryManagerImpl.h" #include "qpid/broker/SaslAuthenticator.h" @@ -147,7 +146,6 @@ Broker::Broker(const Broker::Options& conf) : conf.qmf2Support) : 0), store(new NullMessageStore), - cluster(new NullCluster), acl(0), dataDir(conf.noDataDir ? std::string() : conf.dataDir), queues(this), @@ -512,9 +510,5 @@ void Broker::setClusterTimer(std::auto_ptr<sys::Timer> t) { const std::string Broker::TCP_TRANSPORT("tcp"); -void Broker::setCluster(std::auto_ptr<Cluster> c) { cluster = c; } - -Cluster& Broker::getCluster() { return *cluster; } - }} // namespace qpid::broker diff --git a/cpp/src/qpid/broker/Broker.h b/cpp/src/qpid/broker/Broker.h index d589b15f19..6636b5d912 100644 --- a/cpp/src/qpid/broker/Broker.h +++ b/cpp/src/qpid/broker/Broker.h @@ -70,7 +70,6 @@ namespace broker { class ExpiryPolicy; class Message; -class Cluster; static const uint16_t DEFAULT_PORT=5672; @@ -154,7 +153,6 @@ public: std::auto_ptr<management::ManagementAgent> managementAgent; ProtocolFactoryMap protocolFactories; std::auto_ptr<MessageStore> store; - std::auto_ptr<Cluster> cluster; AclModule* acl; DataDir dataDir; @@ -275,9 +273,6 @@ public: void setClusterUpdatee(bool set) { clusterUpdatee = set; } bool isClusterUpdatee() const { return clusterUpdatee; } - QPID_BROKER_EXTERN void setCluster(std::auto_ptr<Cluster> c); - QPID_BROKER_EXTERN Cluster& getCluster(); - management::ManagementAgent* getManagementAgent() { return managementAgent.get(); } ConnectionCounter& getConnectionCounter() {return connectionCounter;} diff --git a/cpp/src/qpid/broker/Cluster.h b/cpp/src/qpid/broker/Cluster.h deleted file mode 100644 index 4dabd98eab..0000000000 --- a/cpp/src/qpid/broker/Cluster.h +++ /dev/null @@ -1,110 +0,0 @@ -#ifndef QPID_BROKER_CLUSTER_H -#define QPID_BROKER_CLUSTER_H - -/* - * - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - * - */ - -#include <boost/intrusive_ptr.hpp> - -namespace qpid { - -namespace framing { -class FieldTable; -} - -namespace broker { - -class Message; -struct QueuedMessage; -class Queue; -class Exchange; - -/** - * NOTE: this is part of an experimental cluster implementation that is not - * yet fully functional. The original cluster implementation remains in place. - * See ../cluster/new-cluster-design.txt - * - * Interface for cluster implementations. Functions on this interface are - * called at relevant points in the Broker's processing. - */ -class Cluster -{ - public: - virtual ~Cluster() {} - - // Messages - - /** In Exchange::route, before the message is enqueued. */ - virtual void routing(const boost::intrusive_ptr<Message>&) = 0; - - /** A message is delivered to a queue. - * Called before actually pushing the message to the queue. - *@return If true the message should be pushed to the queue now. - * otherwise the cluster code will push the message when it is replicated. - */ - virtual bool enqueue(Queue& queue, const boost::intrusive_ptr<Message>&) = 0; - - /** In Exchange::route, after all enqueues for the message. */ - virtual void routed(const boost::intrusive_ptr<Message>&) = 0; - - /** A message is acquired by a local consumer, it is unavailable to replicas. */ - virtual void acquire(const QueuedMessage&) = 0; - /** A locally-acquired message is accepted, it is removed from all replicas. */ - virtual void accept(const QueuedMessage&) = 0; - - /** A locally-acquired message is rejected, and may be re-routed. */ - virtual void reject(const QueuedMessage&) = 0; - /** Re-routing (if any) is complete for a rejected message. */ - virtual void rejected(const QueuedMessage&) = 0; - - /** A locally-acquired message is released by the consumer and re-queued. */ - virtual void release(const QueuedMessage&) = 0; - - /** A message is removed from the queue. It could have been - * accepted, rejected or dropped for other reasons e.g. expired or - * replaced on an LVQ. - */ - virtual void drop(const QueuedMessage&) = 0; - - // Consumers - - /** A new consumer subscribes to a queue. */ - virtual void consume(const Queue&, size_t consumerCount) = 0; - /** A consumer cancels its subscription to a queue */ - virtual void cancel(const Queue&, size_t consumerCount) = 0; - - // Wiring - - /** A queue is created */ - virtual void create(const Queue&) = 0; - /** A queue is destroyed */ - virtual void destroy(const Queue&) = 0; - /** An exchange is created */ - virtual void create(const Exchange&) = 0; - /** An exchange is destroyed */ - virtual void destroy(const Exchange&) = 0; - /** A binding is created */ - virtual void bind(const Queue&, const Exchange&, const std::string& key, const framing::FieldTable& args) = 0; -}; - -}} // namespace qpid::broker - -#endif /*!QPID_BROKER_CLUSTER_H*/ diff --git a/cpp/src/qpid/broker/DeliveryRecord.cpp b/cpp/src/qpid/broker/DeliveryRecord.cpp index 315b1af2a8..9443eb6ea5 100644 --- a/cpp/src/qpid/broker/DeliveryRecord.cpp +++ b/cpp/src/qpid/broker/DeliveryRecord.cpp @@ -112,7 +112,7 @@ void DeliveryRecord::complete() { bool DeliveryRecord::accept(TransactionContext* ctxt) { if (acquired && !ended) { - queue->accept(ctxt, msg); + queue->dequeue(ctxt, msg); setEnded(); QPID_LOG(debug, "Accepted " << id); } @@ -130,8 +130,19 @@ void DeliveryRecord::committed() const{ } void DeliveryRecord::reject() -{ - queue->reject(msg); +{ + Exchange::shared_ptr alternate = queue->getAlternateExchange(); + if (alternate) { + DeliverableMessage delivery(msg.payload); + alternate->route(delivery, msg.payload->getRoutingKey(), msg.payload->getApplicationHeaders()); + QPID_LOG(info, "Routed rejected message from " << queue->getName() << " to " + << alternate->getName()); + } else { + //just drop it + QPID_LOG(info, "Dropping rejected message from " << queue->getName()); + } + + dequeue(); } uint32_t DeliveryRecord::getCredit() const @@ -145,7 +156,7 @@ void DeliveryRecord::acquire(DeliveryIds& results) { results.push_back(id); if (!acceptExpected) { if (ended) { QPID_LOG(error, "Can't dequeue ended message"); } - else { queue->accept(0, msg); setEnded(); } + else { queue->dequeue(0, msg); setEnded(); } } } else { QPID_LOG(info, "Message already acquired " << id.getValue()); diff --git a/cpp/src/qpid/broker/Exchange.cpp b/cpp/src/qpid/broker/Exchange.cpp index b499171418..d143471559 100644 --- a/cpp/src/qpid/broker/Exchange.cpp +++ b/cpp/src/qpid/broker/Exchange.cpp @@ -23,7 +23,6 @@ #include "qpid/broker/ExchangeRegistry.h" #include "qpid/broker/FedOps.h" #include "qpid/broker/Broker.h" -#include "qpid/broker/Cluster.h" #include "qpid/management/ManagementAgent.h" #include "qpid/broker/Queue.h" #include "qpid/log/Statement.h" @@ -71,23 +70,10 @@ Exchange::PreRoute::~PreRoute(){ } } -// Bracket a scope with calls to Cluster::routing and Cluster::routed -struct ScopedClusterRouting { - Broker* broker; - boost::intrusive_ptr<Message> message; - ScopedClusterRouting(Broker* b, boost::intrusive_ptr<Message> m) - : broker(b), message(m) { - if (broker) broker->getCluster().routing(message); - } - ~ScopedClusterRouting() { - if (broker) broker->getCluster().routed(message); - } -}; - void Exchange::doRoute(Deliverable& msg, ConstBindingList b) { - ScopedClusterRouting scr(broker, &msg.getMessage()); int count = 0; + if (b.get()) { // Block the content release if the message is transient AND there is more than one binding if (!msg.getMessage().isPersistent() && b->size() > 1) { diff --git a/cpp/src/qpid/broker/NullCluster.h b/cpp/src/qpid/broker/NullCluster.h deleted file mode 100644 index 0e11ceef27..0000000000 --- a/cpp/src/qpid/broker/NullCluster.h +++ /dev/null @@ -1,66 +0,0 @@ -#ifndef QPID_BROKER_NULLCLUSTER_H -#define QPID_BROKER_NULLCLUSTER_H - -/* - * - * Licensed to the Apache Software Foundation (ASF) under one - * or more contributor license agreements. See the NOTICE file - * distributed with this work for additional information - * regarding copyright ownership. The ASF licenses this file - * to you under the Apache License, Version 2.0 (the - * "License"); you may not use this file except in compliance - * with the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, - * software distributed under the License is distributed on an - * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY - * KIND, either express or implied. See the License for the - * specific language governing permissions and limitations - * under the License. - * - */ - -#include <qpid/broker/Cluster.h> - -namespace qpid { -namespace broker { - -/** - * No-op implementation of Cluster interface, installed by broker when - * no cluster plug-in is present or clustering is disabled. - */ -class NullCluster : public Cluster -{ - public: - - // Messages - - virtual void routing(const boost::intrusive_ptr<Message>&) {} - virtual bool enqueue(Queue&, const boost::intrusive_ptr<Message>&) { return true; } - virtual void routed(const boost::intrusive_ptr<Message>&) {} - virtual void acquire(const QueuedMessage&) {} - virtual void accept(const QueuedMessage&) {} - virtual void reject(const QueuedMessage&) {} - virtual void rejected(const QueuedMessage&) {} - virtual void release(const QueuedMessage&) {} - virtual void drop(const QueuedMessage&) {} - - // Consumers - - virtual void consume(const Queue&, size_t) {} - virtual void cancel(const Queue&, size_t) {} - - // Wiring - - virtual void create(const Queue&) {} - virtual void destroy(const Queue&) {} - virtual void create(const Exchange&) {} - virtual void destroy(const Exchange&) {} - virtual void bind(const Queue&, const Exchange&, const std::string&, const framing::FieldTable&) {} -}; - -}} // namespace qpid::broker - -#endif diff --git a/cpp/src/qpid/broker/Queue.cpp b/cpp/src/qpid/broker/Queue.cpp index c530e9cd51..e59857462c 100644 --- a/cpp/src/qpid/broker/Queue.cpp +++ b/cpp/src/qpid/broker/Queue.cpp @@ -20,7 +20,6 @@ */ #include "qpid/broker/Broker.h" -#include "qpid/broker/Cluster.h" #include "qpid/broker/Queue.h" #include "qpid/broker/QueueEvents.h" #include "qpid/broker/Exchange.h" @@ -146,10 +145,6 @@ void Queue::deliver(boost::intrusive_ptr<Message> msg){ // Check for deferred delivery in a cluster. if (broker && broker->deferDelivery(name, msg)) return; - // Same thing but for the new cluster interface. - if (broker && !broker->getCluster().enqueue(*this, msg)) - return; - if (msg->isImmediate() && getConsumerCount() == 0) { if (alternateExchange) { DeliverableMessage deliverable(msg); @@ -169,6 +164,7 @@ void Queue::deliver(boost::intrusive_ptr<Message> msg){ }else { push(msg); } + mgntEnqStats(msg); QPID_LOG(debug, "Message " << msg << " enqueued on " << name); } } @@ -202,6 +198,7 @@ void Queue::recover(boost::intrusive_ptr<Message>& msg){ void Queue::process(boost::intrusive_ptr<Message>& msg){ push(msg); + mgntEnqStats(msg); if (mgmtObject != 0){ mgmtObject->inc_msgTxnEnqueues (); mgmtObject->inc_byteTxnEnqueues (msg->contentSize ()); @@ -227,7 +224,6 @@ void Queue::requeue(const QueuedMessage& msg){ } } } - if (broker) broker->getCluster().release(msg); copy.notify(); } @@ -240,22 +236,8 @@ void Queue::clearLVQIndex(const QueuedMessage& msg){ } } -// Inform the cluster of an acquired message on exit from a function -// that does the acquiring. The calling function should set qmsg -// to the acquired message. -struct ClusterAcquireOnExit { - Broker* broker; - QueuedMessage qmsg; - ClusterAcquireOnExit(Broker* b) : broker(b) {} - ~ClusterAcquireOnExit() { - if (broker && qmsg.queue) broker->getCluster().acquire(qmsg); - } -}; - bool Queue::acquireMessageAt(const SequenceNumber& position, QueuedMessage& message) { - ClusterAcquireOnExit willAcquire(broker); - Mutex::ScopedLock locker(messageLock); assertClusterSafe(); QPID_LOG(debug, "Attempting to acquire message at " << position); @@ -266,18 +248,16 @@ bool Queue::acquireMessageAt(const SequenceNumber& position, QueuedMessage& mess if (lastValueQueue) { clearLVQIndex(*i); } - QPID_LOG(debug, "Acquired message at " << i->position << " from " << name); - willAcquire.qmsg = *i; + QPID_LOG(debug, + "Acquired message at " << i->position << " from " << name); messages.erase(i); return true; - } + } QPID_LOG(debug, "Could not acquire message at " << position << " from " << name << "; no message at that position"); return false; } bool Queue::acquire(const QueuedMessage& msg) { - ClusterAcquireOnExit acquire(broker); - Mutex::ScopedLock locker(messageLock); assertClusterSafe(); @@ -285,17 +265,16 @@ bool Queue::acquire(const QueuedMessage& msg) { Messages::iterator i = findAt(msg.position); if ((i != messages.end() && i->position == msg.position) && // note that in some cases payload not be set (!lastValueQueue || - (lastValueQueue && msg.payload.get() == checkLvqReplace(*i).payload.get()) ) // note this is safe for no payload set 0==0 - ) { + (lastValueQueue && msg.payload.get() == checkLvqReplace(*i).payload.get()) ) // note this is safe for no payload set 0==0 + ) { clearLVQIndex(msg); QPID_LOG(debug, "Match found, acquire succeeded: " << i->position << " == " << msg.position); - acquire.qmsg = *i; messages.erase(i); return true; - } + } QPID_LOG(debug, "Acquire failed for " << msg.position); return false; @@ -335,8 +314,6 @@ bool Queue::getNextMessage(QueuedMessage& m, Consumer::shared_ptr c) Queue::ConsumeCode Queue::consumeNextMessage(QueuedMessage& m, Consumer::shared_ptr c) { while (true) { - ClusterAcquireOnExit willAcquire(broker); // Outside the lock - Mutex::ScopedLock locker(messageLock); if (messages.empty()) { QPID_LOG(debug, "No messages to dispatch on queue '" << name << "'"); @@ -353,7 +330,6 @@ Queue::ConsumeCode Queue::consumeNextMessage(QueuedMessage& m, Consumer::shared_ if (c->filter(msg.payload)) { if (c->accept(msg.payload)) { m = msg; - willAcquire.qmsg = msg; popMsg(msg); return CONSUMED; } else { @@ -475,51 +451,40 @@ QueuedMessage Queue::find(SequenceNumber pos) const { return QueuedMessage(); } -void Queue::consume(Consumer::shared_ptr c, bool requestExclusive) { +void Queue::consume(Consumer::shared_ptr c, bool requestExclusive){ assertClusterSafe(); - size_t consumers; - { - Mutex::ScopedLock locker(consumerLock); - if(exclusive) { + Mutex::ScopedLock locker(consumerLock); + if(exclusive) { + throw ResourceLockedException( + QPID_MSG("Queue " << getName() << " has an exclusive consumer. No more consumers allowed.")); + } else if(requestExclusive) { + if(consumerCount) { throw ResourceLockedException( - QPID_MSG("Queue " << getName() << " has an exclusive consumer. No more consumers allowed.")); - } else if(requestExclusive) { - if(consumerCount) { - throw ResourceLockedException( - QPID_MSG("Queue " << getName() << " already has consumers. Exclusive access denied.")); - } else { - exclusive = c->getSession(); - } + QPID_MSG("Queue " << getName() << " already has consumers. Exclusive access denied.")); + } else { + exclusive = c->getSession(); } - consumers = ++consumerCount; - if (mgmtObject != 0) - mgmtObject->inc_consumerCount (); } - if (broker) broker->getCluster().consume(*this, consumers); + consumerCount++; + if (mgmtObject != 0) + mgmtObject->inc_consumerCount (); } void Queue::cancel(Consumer::shared_ptr c){ removeListener(c); - size_t consumers; - { - Mutex::ScopedLock locker(consumerLock); - consumers = --consumerCount; - if(exclusive) exclusive = 0; - if (mgmtObject != 0) - mgmtObject->dec_consumerCount (); - } - if (broker) broker->getCluster().cancel(*this, consumers); + Mutex::ScopedLock locker(consumerLock); + consumerCount--; + if(exclusive) exclusive = 0; + if (mgmtObject != 0) + mgmtObject->dec_consumerCount (); } QueuedMessage Queue::get(){ - ClusterAcquireOnExit acquire(broker); // Outside lock - Mutex::ScopedLock locker(messageLock); QueuedMessage msg(this); if(!messages.empty()){ msg = getFront(); - acquire.qmsg = msg; popMsg(msg); } return msg; @@ -644,12 +609,10 @@ void Queue::popMsg(QueuedMessage& qmsg) void Queue::push(boost::intrusive_ptr<Message>& msg, bool isRecovery){ assertClusterSafe(); - if (!isRecovery) mgntEnqStats(msg); - QueuedMessage qm; QueueListeners::NotificationSet copy; { Mutex::ScopedLock locker(messageLock); - qm = QueuedMessage(this, msg, ++sequence); + QueuedMessage qm(this, msg, ++sequence); if (insertSeqNo) msg->getOrInsertHeaders().setInt64(seqNoKey, sequence); LVQ::iterator i; @@ -666,14 +629,12 @@ void Queue::push(boost::intrusive_ptr<Message>& msg, bool isRecovery){ boost::intrusive_ptr<Message> old = i->second->getReplacementMessage(this); if (!old) old = i->second; i->second->setReplacementMessage(msg,this); - // FIXME aconway 2010-10-15: it is incorrect to use qm.position below - // should be using the position of the message being replaced. if (isRecovery) { //can't issue new requests for the store until //recovery is complete pendingDequeues.push_back(QueuedMessage(qm.queue, old, qm.position)); } else { - Mutex::ScopedUnlock u(messageLock); + Mutex::ScopedUnlock u(messageLock); dequeue(0, QueuedMessage(qm.queue, old, qm.position)); } } @@ -831,48 +792,19 @@ void Queue::enqueueAborted(boost::intrusive_ptr<Message> msg) if (policy.get()) policy->enqueueAborted(msg); } -void Queue::accept(TransactionContext* ctxt, const QueuedMessage& msg) { - if (broker) broker->getCluster().accept(msg); - dequeue(ctxt, msg); -} - -struct ScopedClusterReject { - Broker* broker; - const QueuedMessage& qmsg; - ScopedClusterReject(Broker* b, const QueuedMessage& m) : broker(b), qmsg(m) { - if (broker) broker->getCluster().reject(qmsg); - } - ~ScopedClusterReject() { - if (broker) broker->getCluster().rejected(qmsg); - } -}; - -void Queue::reject(const QueuedMessage &msg) { - ScopedClusterReject scr(broker, msg); - Exchange::shared_ptr alternate = getAlternateExchange(); - if (alternate) { - DeliverableMessage delivery(msg.payload); - alternate->route(delivery, msg.payload->getRoutingKey(), msg.payload->getApplicationHeaders()); - QPID_LOG(info, "Routed rejected message from " << getName() << " to " - << alternate->getName()); - } else { - //just drop it - QPID_LOG(info, "Dropping rejected message from " << getName()); - } - dequeue(0, msg); -} - // return true if store exists, bool Queue::dequeue(TransactionContext* ctxt, const QueuedMessage& msg) { ScopedUse u(barrier); if (!u.acquired) return false; + { Mutex::ScopedLock locker(messageLock); if (!isEnqueued(msg)) return false; - if (!ctxt) dequeued(msg); + if (!ctxt) { + dequeued(msg); + } } - if (!ctxt && broker) broker->getCluster().drop(msg); // Outside lock // This check prevents messages which have been forced persistent on one queue from dequeuing // from another on which no forcing has taken place and thus causing a store error. bool fp = msg.payload->isForcedPersistent(); @@ -889,7 +821,6 @@ bool Queue::dequeue(TransactionContext* ctxt, const QueuedMessage& msg) void Queue::dequeueCommitted(const QueuedMessage& msg) { - if (broker) broker->getCluster().drop(msg); // Outside lock Mutex::ScopedLock locker(messageLock); dequeued(msg); if (mgmtObject != 0) { @@ -915,8 +846,6 @@ void Queue::popAndDequeue() */ void Queue::dequeued(const QueuedMessage& msg) { - // Note: Cluster::drop does only local book-keeping, no multicast - // So OK to call here with lock held. if (policy.get()) policy->dequeued(msg); mgntDeqStats(msg.payload); if (eventMode == ENQUEUE_AND_DEQUEUE && eventMgr) { @@ -932,7 +861,6 @@ void Queue::create(const FieldTable& _settings) store->create(*this, _settings); } configure(_settings); - if (broker) broker->getCluster().create(*this); } void Queue::configure(const FieldTable& _settings, bool recovering) @@ -1006,7 +934,6 @@ void Queue::destroy() store->destroy(*this); store = 0;//ensure we make no more calls to the store for this queue } - if (broker) broker->getCluster().destroy(*this); } void Queue::notifyDeleted() diff --git a/cpp/src/qpid/broker/Queue.h b/cpp/src/qpid/broker/Queue.h index 572f3dc0e2..96c79d1b92 100644 --- a/cpp/src/qpid/broker/Queue.h +++ b/cpp/src/qpid/broker/Queue.h @@ -259,13 +259,6 @@ class Queue : public boost::enable_shared_from_this<Queue>, bool enqueue(TransactionContext* ctxt, boost::intrusive_ptr<Message>& msg, bool suppressPolicyCheck = false); void enqueueAborted(boost::intrusive_ptr<Message> msg); - - /** Message acknowledged, dequeue it. */ - QPID_BROKER_EXTERN void accept(TransactionContext* ctxt, const QueuedMessage &msg); - - /** Message rejected, dequeue it and re-route to alternate exchange if necessary. */ - QPID_BROKER_EXTERN void reject(const QueuedMessage &msg); - /** * dequeue from store (only done once messages is acknowledged) */ diff --git a/cpp/src/qpid/broker/QueuedMessage.h b/cpp/src/qpid/broker/QueuedMessage.h index 8cf73bda52..35e48b11f3 100644 --- a/cpp/src/qpid/broker/QueuedMessage.h +++ b/cpp/src/qpid/broker/QueuedMessage.h @@ -34,9 +34,10 @@ struct QueuedMessage framing::SequenceNumber position; Queue* queue; - QueuedMessage(Queue* q=0) : position(0), queue(q) {} + QueuedMessage() : queue(0) {} QueuedMessage(Queue* q, boost::intrusive_ptr<Message> msg, framing::SequenceNumber sn) : payload(msg), position(sn), queue(q) {} + QueuedMessage(Queue* q) : queue(q) {} }; inline bool operator<(const QueuedMessage& a, const QueuedMessage& b) { return a.position < b.position; } diff --git a/cpp/src/qpid/broker/SemanticState.cpp b/cpp/src/qpid/broker/SemanticState.cpp index f393879c16..c91cfba2f8 100644 --- a/cpp/src/qpid/broker/SemanticState.cpp +++ b/cpp/src/qpid/broker/SemanticState.cpp @@ -333,7 +333,7 @@ bool SemanticState::ConsumerImpl::deliver(QueuedMessage& msg) parent->record(record); } if (acquire && !ackExpected) { - queue->accept(0, msg); + queue->dequeue(0, msg); } if (mgmtObject) { mgmtObject->inc_delivered(); } return true; @@ -347,6 +347,11 @@ bool SemanticState::ConsumerImpl::filter(intrusive_ptr<Message>) bool SemanticState::ConsumerImpl::accept(intrusive_ptr<Message> msg) { assertClusterSafe(); + // FIXME aconway 2009-06-08: if we have byte & message credit but + // checkCredit fails because the message is to big, we should + // remain on queue's listener list for possible smaller messages + // in future. + // blocked = !(filter(msg) && checkCredit(msg)); return !blocked; } |
