diff options
| author | Kenneth Anthony Giusti <kgiusti@apache.org> | 2011-08-08 20:50:59 +0000 |
|---|---|---|
| committer | Kenneth Anthony Giusti <kgiusti@apache.org> | 2011-08-08 20:50:59 +0000 |
| commit | 7ef72217880f040495b824324329189f9e38770e (patch) | |
| tree | ed73caa31d799b5cbdd83c837802b8b6b5ef21e2 /qpid/cpp/src | |
| parent | 45012dc9764465561ee1edbc4d3de4fac03c5b54 (diff) | |
| download | qpid-python-7ef72217880f040495b824324329189f9e38770e.tar.gz | |
QPID-3346: incorporate review input
git-svn-id: https://svn.apache.org/repos/asf/qpid/branches/qpid-3346@1155095 13f79535-47bb-0310-9956-ffa450edef68
Diffstat (limited to 'qpid/cpp/src')
| -rw-r--r-- | qpid/cpp/src/qpid/broker/DeliveryRecord.cpp | 11 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/DeliveryRecord.h | 7 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/LegacyLVQ.cpp | 4 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/Queue.cpp | 72 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/Queue.h | 8 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/QueueEvents.cpp | 2 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/QueueFlowLimit.cpp | 6 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/QueueFlowLimit.h | 2 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/QueueObserver.h | 6 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/SemanticState.cpp | 2 | ||||
| -rw-r--r-- | qpid/cpp/src/qpid/broker/ThresholdAlerts.h | 2 | ||||
| -rw-r--r-- | qpid/cpp/src/tests/QueueTest.cpp | 10 |
12 files changed, 66 insertions, 66 deletions
diff --git a/qpid/cpp/src/qpid/broker/DeliveryRecord.cpp b/qpid/cpp/src/qpid/broker/DeliveryRecord.cpp index 1b42c67edd..a82a54f077 100644 --- a/qpid/cpp/src/qpid/broker/DeliveryRecord.cpp +++ b/qpid/cpp/src/qpid/broker/DeliveryRecord.cpp @@ -152,14 +152,9 @@ uint32_t DeliveryRecord::getCredit() const return credit; } -void DeliveryRecord::acquire(SemanticState* const session, DeliveryIds& results) { - SemanticState::ConsumerImpl::shared_ptr consumer; - - if (!session->find( tag, consumer )) { - QPID_LOG(error, "Can't acquire message " << id.getValue() << ": original subscription no longer exists."); - } - - if (queue->acquire(msg, consumer)) { +void DeliveryRecord::acquire(DeliveryIds& results) +{ + if (queue->acquire(msg, tag)) { acquired = true; results.push_back(id); if (!acceptExpected) { diff --git a/qpid/cpp/src/qpid/broker/DeliveryRecord.h b/qpid/cpp/src/qpid/broker/DeliveryRecord.h index ba3e1d5cfb..4465825d7f 100644 --- a/qpid/cpp/src/qpid/broker/DeliveryRecord.h +++ b/qpid/cpp/src/qpid/broker/DeliveryRecord.h @@ -82,7 +82,7 @@ class DeliveryRecord void reject(); void cancel(const std::string& tag); void redeliver(SemanticState* const); - void acquire(SemanticState* const, DeliveryIds& results); + void acquire(DeliveryIds& results); void complete(); bool accept(TransactionContext* ctxt); // Returns isRedundant() bool setEnded(); // Returns isRedundant() @@ -117,14 +117,13 @@ inline bool operator<(const DeliveryRecord& a, const framing::SequenceNumber& b) struct AcquireFunctor { - SemanticState* session; DeliveryIds& results; - AcquireFunctor(SemanticState* _session, DeliveryIds& _results) : session(_session), results(_results) {} + AcquireFunctor(DeliveryIds& _results) : results(_results) {} void operator()(DeliveryRecord& record) { - record.acquire(session, results); + record.acquire(results); } }; diff --git a/qpid/cpp/src/qpid/broker/LegacyLVQ.cpp b/qpid/cpp/src/qpid/broker/LegacyLVQ.cpp index 7d9cb4c1a0..a811a86492 100644 --- a/qpid/cpp/src/qpid/broker/LegacyLVQ.cpp +++ b/qpid/cpp/src/qpid/broker/LegacyLVQ.cpp @@ -35,9 +35,7 @@ void LegacyLVQ::setNoBrowse(bool b) bool LegacyLVQ::remove(const framing::SequenceNumber& position, QueuedMessage& message) { Ordering::iterator i = messages.find(position); - if (i != messages.end() && - // @todo KAG: gsim? is a bug? message is a *return* value - we really shouldn't check ".payload" below: - i->second.payload == message.payload) { + if (i != messages.end() && i->second.payload == message.payload) { message = i->second; erase(i); return true; diff --git a/qpid/cpp/src/qpid/broker/Queue.cpp b/qpid/cpp/src/qpid/broker/Queue.cpp index c9cea9212a..68dd2ae125 100644 --- a/qpid/cpp/src/qpid/broker/Queue.cpp +++ b/qpid/cpp/src/qpid/broker/Queue.cpp @@ -92,19 +92,25 @@ const int ENQUEUE_AND_DEQUEUE=2; namespace qpid { namespace broker { -class MessageSelector +class MessageAllocator { protected: Queue *queue; public: - MessageSelector( Queue *q ) : queue(q) {} - virtual ~MessageSelector() {}; + MessageAllocator( Queue *q ) : queue(q) {} + virtual ~MessageAllocator() {}; // assumes caller holds messageLock virtual bool nextMessage( Consumer::shared_ptr c, QueuedMessage& next, const Mutex::ScopedLock&); - virtual bool canAcquire(Consumer::shared_ptr consumer, const QueuedMessage& qm, - const Mutex::ScopedLock&); + /** acquire a message previously browsed via nextMessage(). assume messageLock held + * @param consumer name of consumer that is attempting to acquire the message + * @param qm the message to be acquired + * @param messageLock - ensures caller is holding it! + * @returns true if acquire is successful, false if acquire failed. + */ + virtual bool canAcquire( const std::string& consumer, const QueuedMessage& qm, + const Mutex::ScopedLock&); }; }} @@ -134,7 +140,7 @@ Queue::Queue(const string& _name, bool _autodelete, deleted(false), barrier(*this), autoDeleteTimeout(0), - selector(new MessageSelector( this )) // KAG TODO: FIX!! + allocator(new MessageAllocator( this )) // KAG TODO: FIX!! { if (parent != 0 && broker != 0) { ManagementAgent* agent = broker->getManagementAgent(); @@ -269,13 +275,13 @@ bool Queue::acquireMessageAt(const SequenceNumber& position, QueuedMessage& mess } } -bool Queue::acquire(const QueuedMessage& msg, Consumer::shared_ptr c) +bool Queue::acquire(const QueuedMessage& msg, const std::string& consumer) { Mutex::ScopedLock locker(messageLock); assertClusterSafe(); - QPID_LOG(debug, c->getName() << " attempting to acquire message at " << msg.position); + QPID_LOG(debug, consumer << " attempting to acquire message at " << msg.position); - if (!selector->canAcquire( c, msg, locker )) { + if (!allocator->canAcquire( consumer, msg, locker )) { QPID_LOG(debug, "Not permitted to acquire msg at " << msg.position << " from '" << name); return false; } @@ -327,7 +333,7 @@ Queue::ConsumeCode Queue::consumeNextMessage(QueuedMessage& m, Consumer::shared_ Mutex::ScopedLock locker(messageLock); QueuedMessage msg; - if (!selector->nextMessage(c, msg, locker)) { // no next available + if (!allocator->nextMessage(c, msg, locker)) { // no next available QPID_LOG(debug, "No messages available to dispatch to consumer " << c->getName() << " on queue '" << name << "'"); listeners.addListener(c); @@ -346,8 +352,10 @@ Queue::ConsumeCode Queue::consumeNextMessage(QueuedMessage& m, Consumer::shared_ if (c->filter(msg.payload)) { if (c->accept(msg.payload)) { - acquire( msg.position, m ); - c->position = msg.position; + bool ok = acquire( msg.position, msg ); + (void) ok; assert(ok); + m = msg; + c->position = m.position; return CONSUMED; } else { //message(s) are available but consumer hasn't got enough credit @@ -369,7 +377,7 @@ bool Queue::browseNextMessage(QueuedMessage& m, Consumer::shared_ptr c) Mutex::ScopedLock locker(messageLock); QueuedMessage msg; - if (!selector->nextMessage(c, msg, locker)) { // no next available + if (!allocator->nextMessage(c, msg, locker)) { // no next available QPID_LOG(debug, "No browsable messages available for consumer " << c->getName() << " on queue '" << name << "'"); listeners.addListener(c); @@ -484,7 +492,7 @@ QueuedMessage Queue::get(){ Mutex::ScopedLock locker(messageLock); QueuedMessage msg(this); if (messages->pop(msg)) - consumed( msg ); + acquired( msg ); return msg; } @@ -520,7 +528,7 @@ void Queue::purgeExpired(qpid::sys::Duration lapse) i != expired.end(); ++i) { { Mutex::ScopedLock locker(messageLock); - consumed( *i ); // expects messageLock held + acquired( *i ); // expects messageLock held } dequeue( 0, *i ); } @@ -596,7 +604,7 @@ void Queue::pop() assertClusterSafe(); QueuedMessage msg; if (messages->pop(msg)) { - consumed( msg ); // mark it removed + acquired( msg ); // mark it removed ++dequeueSincePurge; } } @@ -605,7 +613,7 @@ void Queue::pop() bool Queue::acquire(const qpid::framing::SequenceNumber& position, QueuedMessage& msg ) { if (messages->remove(position, msg)) { - consumed( msg ); + acquired( msg ); return true; } return false; @@ -627,7 +635,7 @@ void Queue::push(boost::intrusive_ptr<Message>& msg, bool isRecovery){ } copy.notify(); if (dequeueRequired) { - consumed( removed ); // tell observers + acquired( removed ); // tell observers if (isRecovery) { //can't issue new requests for the store until //recovery is complete @@ -823,11 +831,11 @@ void Queue::dequeued(const QueuedMessage& msg) /** updates queue observers when a message has become unavailable for transfer, * expects messageLock to be held */ -void Queue::consumed(const QueuedMessage& msg) +void Queue::acquired(const QueuedMessage& msg) { for (Observers::const_iterator i = observers.begin(); i != observers.end(); ++i) { try{ - (*i)->consumed(msg); + (*i)->acquired(msg); } catch (const std::exception& e) { QPID_LOG(warning, "Exception on notification of message removal for queue " << getName() << ": " << e.what()); } @@ -1349,7 +1357,7 @@ void Queue::UsageBarrier::destroy() // KAG TBD: flesh out... -class MessageGroupManager : public QueueObserver, public MessageSelector +class MessageGroupManager : public QueueObserver, public MessageAllocator { const std::string groupIdHeader; // msg header holding group identifier struct GroupState { @@ -1366,7 +1374,7 @@ class MessageGroupManager : public QueueObserver, public MessageSelector public: MessageGroupManager(const std::string& header, Queue *q ) - : QueueObserver(), MessageSelector(q), groupIdHeader( header ) {} + : QueueObserver(), MessageAllocator(q), groupIdHeader( header ) {} void enqueued( const QueuedMessage& qm ); void removed( const QueuedMessage& qm ); void requeued( const QueuedMessage& qm ); @@ -1375,7 +1383,7 @@ class MessageGroupManager : public QueueObserver, public MessageSelector void consumerRemoved( const Consumer& ); bool nextMessage( Consumer::shared_ptr c, QueuedMessage& next, const Mutex::ScopedLock&); - bool canAcquire(Consumer::shared_ptr consumer, const QueuedMessage& msg, + bool canAcquire(const std::string& consumer, const QueuedMessage& msg, const Mutex::ScopedLock&); }; @@ -1469,11 +1477,11 @@ bool MessageGroupManager::nextMessage( Consumer::shared_ptr c, QueuedMessage& ne const Mutex::ScopedLock& l) { // KAG TODO: FIX!!! - return MessageSelector::nextMessage( c, next, l ); + return MessageAllocator::nextMessage( c, next, l ); } -bool MessageGroupManager::canAcquire(Consumer::shared_ptr consumer, const QueuedMessage& qm, +bool MessageGroupManager::canAcquire(const std::string& consumer, const QueuedMessage& qm, const Mutex::ScopedLock&) { std::string group( getGroupId(qm, groupIdHeader) ); @@ -1482,18 +1490,18 @@ bool MessageGroupManager::canAcquire(Consumer::shared_ptr consumer, const Queued GroupState& state( gs->second ); if (state.owner.empty()) { - state.owner = consumer->getName(); + state.owner = consumer; return true; } - return state.owner == consumer->getName(); + return state.owner == consumer; } -// default selector - requires messageLock to be held by caller! -bool MessageSelector::nextMessage( Consumer::shared_ptr c, QueuedMessage& next, +// default allocator - requires messageLock to be held by caller! +bool MessageAllocator::nextMessage( Consumer::shared_ptr c, QueuedMessage& next, const Mutex::ScopedLock& /*just to enforce locking*/) { Messages& messages(queue->getMessages()); @@ -1510,9 +1518,9 @@ bool MessageSelector::nextMessage( Consumer::shared_ptr c, QueuedMessage& next, } -// default selector - requires messageLock to be held by caller! -bool MessageSelector::canAcquire(Consumer::shared_ptr, const QueuedMessage&, - const Mutex::ScopedLock& /*just to enforce locking*/) +// default allocator - requires messageLock to be held by caller! +bool MessageAllocator::canAcquire(const std::string&, const QueuedMessage&, + const Mutex::ScopedLock& /*just to enforce locking*/) { return true; // always give permission to acquire } diff --git a/qpid/cpp/src/qpid/broker/Queue.h b/qpid/cpp/src/qpid/broker/Queue.h index a6bb0d6915..5c895ceb1b 100644 --- a/qpid/cpp/src/qpid/broker/Queue.h +++ b/qpid/cpp/src/qpid/broker/Queue.h @@ -59,7 +59,7 @@ class MessageStore; class QueueEvents; class QueueRegistry; class TransactionContext; -class MessageSelector; +class MessageAllocator; /** * The brokers representation of an amqp queue. Messages are @@ -129,7 +129,7 @@ class Queue : public boost::enable_shared_from_this<Queue>, UsageBarrier barrier; int autoDeleteTimeout; boost::intrusive_ptr<qpid::sys::TimerTask> autoDeleteTask; - std::auto_ptr<MessageSelector> selector; + std::auto_ptr<MessageAllocator> allocator; void push(boost::intrusive_ptr<Message>& msg, bool isRecovery=false); void setPolicy(std::auto_ptr<QueuePolicy> policy); @@ -144,7 +144,7 @@ class Queue : public boost::enable_shared_from_this<Queue>, /** update queue observers with new message state */ void enqueued(const QueuedMessage& msg); - void consumed(const QueuedMessage& msg); + void acquired(const QueuedMessage& msg); void dequeued(const QueuedMessage& msg); /** modify the Queue's message container - assumes messageLock held */ @@ -204,7 +204,7 @@ class Queue : public boost::enable_shared_from_this<Queue>, * @param msg - message to be acquired. * @return false if message is no longer available for acquire. */ - QPID_BROKER_EXTERN bool acquire(const QueuedMessage& msg, const Consumer::shared_ptr c); + QPID_BROKER_EXTERN bool acquire(const QueuedMessage& msg, const std::string& consumer); /** * Used to configure a new queue and create a persistent record diff --git a/qpid/cpp/src/qpid/broker/QueueEvents.cpp b/qpid/cpp/src/qpid/broker/QueueEvents.cpp index 764faf5fd7..c66bdabf0f 100644 --- a/qpid/cpp/src/qpid/broker/QueueEvents.cpp +++ b/qpid/cpp/src/qpid/broker/QueueEvents.cpp @@ -130,7 +130,7 @@ class EventGenerator : public QueueObserver if (!enqueueOnly) manager.dequeued(m); } - void consumed(const QueuedMessage&) {}; + void acquired(const QueuedMessage&) {}; void requeued(const QueuedMessage&) {}; private: diff --git a/qpid/cpp/src/qpid/broker/QueueFlowLimit.cpp b/qpid/cpp/src/qpid/broker/QueueFlowLimit.cpp index db18325c78..f15bb45c01 100644 --- a/qpid/cpp/src/qpid/broker/QueueFlowLimit.cpp +++ b/qpid/cpp/src/qpid/broker/QueueFlowLimit.cpp @@ -378,11 +378,11 @@ void QueueFlowLimit::setState(const qpid::framing::FieldTable& state) fcmsg.add(first, last); for (SequenceNumber seq = first; seq <= last; ++seq) { QueuedMessage msg; - bool found = queue->find(seq, msg); // fyi: msg.payload may be null if msg is delivered & unacked - (void) found; assert(found); // avoid unused variable warning when NDEBUG set + queue->find(seq, msg); // fyi: may not be found if msg is acquired & unacked bool unique; unique = index.insert(std::pair<framing::SequenceNumber, boost::intrusive_ptr<Message> >(seq, msg.payload)).second; - (void) unique; assert(unique); // ditto NDEBUG warning + // Like this to avoid tripping up unused variable warning when NDEBUG set + if (!unique) assert(unique); } } } diff --git a/qpid/cpp/src/qpid/broker/QueueFlowLimit.h b/qpid/cpp/src/qpid/broker/QueueFlowLimit.h index 3d4b31bb69..ad8a2720ef 100644 --- a/qpid/cpp/src/qpid/broker/QueueFlowLimit.h +++ b/qpid/cpp/src/qpid/broker/QueueFlowLimit.h @@ -85,7 +85,7 @@ class Broker; /** the queue has removed QueuedMessage. Returns true if flow state changes */ QPID_BROKER_EXTERN void dequeued(const QueuedMessage&); /** ignored */ - QPID_BROKER_EXTERN void consumed(const QueuedMessage&) {}; + QPID_BROKER_EXTERN void acquired(const QueuedMessage&) {}; QPID_BROKER_EXTERN void requeued(const QueuedMessage&) {}; /** for clustering: */ diff --git a/qpid/cpp/src/qpid/broker/QueueObserver.h b/qpid/cpp/src/qpid/broker/QueueObserver.h index 9c3c186f23..a2e8aa5aca 100644 --- a/qpid/cpp/src/qpid/broker/QueueObserver.h +++ b/qpid/cpp/src/qpid/broker/QueueObserver.h @@ -45,7 +45,7 @@ class Consumer; * "Enqueued" - the message is "Available" - on the queue for transfer to any consumer * (e.g. browse or acquire) * - * "Consumed" - the message is "Locked" - a consumer has claimed exclusive access to it. + * "Acquired" - the message is "Locked" - a consumer has claimed exclusive access to it. * It is no longer available for other consumers to browse or acquire, but it is not yet * considered dequeued as it may be requeued by the consumer. * @@ -62,9 +62,9 @@ class QueueObserver // note: the Queue will hold the messageLock while calling these methods! virtual void enqueued(const QueuedMessage&) = 0; - virtual void consumed(const QueuedMessage&) = 0; - virtual void requeued(const QueuedMessage&) = 0; virtual void dequeued(const QueuedMessage&) = 0; + virtual void acquired(const QueuedMessage&) = 0; + virtual void requeued(const QueuedMessage&) = 0; virtual void consumerAdded( const Consumer& ) {}; virtual void consumerRemoved( const Consumer& ) {}; private: diff --git a/qpid/cpp/src/qpid/broker/SemanticState.cpp b/qpid/cpp/src/qpid/broker/SemanticState.cpp index c8f77ba64e..24d6607ac1 100644 --- a/qpid/cpp/src/qpid/broker/SemanticState.cpp +++ b/qpid/cpp/src/qpid/broker/SemanticState.cpp @@ -691,7 +691,7 @@ AckRange SemanticState::findRange(DeliveryId first, DeliveryId last) void SemanticState::acquire(DeliveryId first, DeliveryId last, DeliveryIds& acquired) { AckRange range = findRange(first, last); - for_each(range.start, range.end, AcquireFunctor(this, acquired)); + for_each(range.start, range.end, AcquireFunctor(acquired)); } void SemanticState::release(DeliveryId first, DeliveryId last, bool setRedelivered) diff --git a/qpid/cpp/src/qpid/broker/ThresholdAlerts.h b/qpid/cpp/src/qpid/broker/ThresholdAlerts.h index c27c97d6f5..2b4a46b736 100644 --- a/qpid/cpp/src/qpid/broker/ThresholdAlerts.h +++ b/qpid/cpp/src/qpid/broker/ThresholdAlerts.h @@ -50,7 +50,7 @@ class ThresholdAlerts : public QueueObserver const long repeatInterval); void enqueued(const QueuedMessage&); void dequeued(const QueuedMessage&); - void consumed(const QueuedMessage&) {}; + void acquired(const QueuedMessage&) {}; void requeued(const QueuedMessage&) {}; static void observe(Queue& queue, qpid::management::ManagementAgent& agent, diff --git a/qpid/cpp/src/tests/QueueTest.cpp b/qpid/cpp/src/tests/QueueTest.cpp index 1f1eb16af5..0eb3c9d194 100644 --- a/qpid/cpp/src/tests/QueueTest.cpp +++ b/qpid/cpp/src/tests/QueueTest.cpp @@ -331,7 +331,7 @@ QPID_AUTO_TEST_CASE(testSearch){ BOOST_CHECK_EQUAL(seq.getValue(), qm.position.getValue()); - queue->acquire(qm, c1); + queue->acquire(qm, c1->getName()); BOOST_CHECK_EQUAL(queue->getMessageCount(), 2u); SequenceNumber seq1(3); QueuedMessage qm1; @@ -557,11 +557,11 @@ QPID_AUTO_TEST_CASE(testLVQAcquire){ QueuedMessage qmsg3(queue.get(), 0, sequence1); TestConsumer::shared_ptr dummy(new TestConsumer()); - BOOST_CHECK(!queue->acquire(qmsg, dummy)); - BOOST_CHECK(queue->acquire(qmsg2, dummy)); + BOOST_CHECK(!queue->acquire(qmsg, dummy->getName())); + BOOST_CHECK(queue->acquire(qmsg2, dummy->getName())); // Acquire the massage again to test failure case. - BOOST_CHECK(!queue->acquire(qmsg2, dummy)); - BOOST_CHECK(!queue->acquire(qmsg3, dummy)); + BOOST_CHECK(!queue->acquire(qmsg2, dummy->getName())); + BOOST_CHECK(!queue->acquire(qmsg3, dummy->getName())); BOOST_CHECK_EQUAL(queue->getMessageCount(), 2u); |
