summaryrefslogtreecommitdiff
path: root/cpp/src/qpid/broker/Queue.cpp
diff options
context:
space:
mode:
authorAlan Conway <aconway@apache.org>2008-10-23 16:21:56 +0000
committerAlan Conway <aconway@apache.org>2008-10-23 16:21:56 +0000
commit9d6e5ee84b3c3a53dfe78d3b5b74986495e7abee (patch)
tree257bed543782bbb7ca80a411a6d26b04a9b0784d /cpp/src/qpid/broker/Queue.cpp
parent1b127dfaac12835181f61637fb751380aff78e7e (diff)
downloadqpid-python-9d6e5ee84b3c3a53dfe78d3b5b74986495e7abee.tar.gz
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
Diffstat (limited to 'cpp/src/qpid/broker/Queue.cpp')
-rw-r--r--cpp/src/qpid/broker/Queue.cpp27
1 files changed, 27 insertions, 0 deletions
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.
+}