summaryrefslogtreecommitdiff
path: root/qpid/cpp/src
diff options
context:
space:
mode:
authorKenneth Anthony Giusti <kgiusti@apache.org>2011-08-08 20:50:59 +0000
committerKenneth Anthony Giusti <kgiusti@apache.org>2011-08-08 20:50:59 +0000
commit7ef72217880f040495b824324329189f9e38770e (patch)
treeed73caa31d799b5cbdd83c837802b8b6b5ef21e2 /qpid/cpp/src
parent45012dc9764465561ee1edbc4d3de4fac03c5b54 (diff)
downloadqpid-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.cpp11
-rw-r--r--qpid/cpp/src/qpid/broker/DeliveryRecord.h7
-rw-r--r--qpid/cpp/src/qpid/broker/LegacyLVQ.cpp4
-rw-r--r--qpid/cpp/src/qpid/broker/Queue.cpp72
-rw-r--r--qpid/cpp/src/qpid/broker/Queue.h8
-rw-r--r--qpid/cpp/src/qpid/broker/QueueEvents.cpp2
-rw-r--r--qpid/cpp/src/qpid/broker/QueueFlowLimit.cpp6
-rw-r--r--qpid/cpp/src/qpid/broker/QueueFlowLimit.h2
-rw-r--r--qpid/cpp/src/qpid/broker/QueueObserver.h6
-rw-r--r--qpid/cpp/src/qpid/broker/SemanticState.cpp2
-rw-r--r--qpid/cpp/src/qpid/broker/ThresholdAlerts.h2
-rw-r--r--qpid/cpp/src/tests/QueueTest.cpp10
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);