From 9d6e5ee84b3c3a53dfe78d3b5b74986495e7abee Mon Sep 17 00:00:00 2001 From: Alan Conway Date: Thu, 23 Oct 2008 16:21:56 +0000 Subject: Minor changes to provide access for cluster to replicate delivery records. - broker::Queue: find message by position, set position. - broker::SemanticState: make record() public, add eachUnacked(), fix typo "NotifyEnabld" - broker::DeliveryRecord: added more public accessors git-svn-id: https://svn.apache.org/repos/asf/incubator/qpid/trunk/qpid@707406 13f79535-47bb-0310-9956-ffa450edef68 --- cpp/src/qpid/broker/Queue.cpp | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) (limited to 'cpp/src/qpid/broker/Queue.cpp') diff --git a/cpp/src/qpid/broker/Queue.cpp b/cpp/src/qpid/broker/Queue.cpp index 1f508a1cc7..52404c826c 100644 --- a/cpp/src/qpid/broker/Queue.cpp +++ b/cpp/src/qpid/broker/Queue.cpp @@ -365,6 +365,25 @@ bool Queue::seek(QueuedMessage& msg, Consumer::shared_ptr c) { return false; } +namespace { +struct PositionEquals { + SequenceNumber pos; + PositionEquals(SequenceNumber p) : pos(p) {} + bool operator()(const QueuedMessage& msg) const { return msg.position == pos; } +}; +}// namespace + +bool Queue::find(QueuedMessage& msg, SequenceNumber pos) const { + Mutex::ScopedLock locker(messageLock); + Messages::const_iterator i = std::find_if(messages.begin(), messages.end(), PositionEquals(pos)); + if (i == messages.end()) + return false; + else { + msg = *i; + return true; + } +} + void Queue::consume(Consumer::shared_ptr c, bool requestExclusive){ Mutex::ScopedLock locker(consumerLock); if(exclusive) { @@ -827,3 +846,11 @@ Manageable::status_t Queue::ManagementMethod (uint32_t methodId, Args& args, str return status; } + +void Queue::setPosition(SequenceNumber n) { + if (n <= sequence) + throw InvalidArgumentException(QPID_MSG("Invalid position " << n << " < " << sequence + << " for queue " << name)); + sequence = n; + --sequence; // Decrement so ++sequence will return n. +} -- cgit v1.2.1