summaryrefslogtreecommitdiff
path: root/qpid/cpp/src/qpid/ha/ReplicatingSubscription.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'qpid/cpp/src/qpid/ha/ReplicatingSubscription.cpp')
-rw-r--r--qpid/cpp/src/qpid/ha/ReplicatingSubscription.cpp20
1 files changed, 14 insertions, 6 deletions
diff --git a/qpid/cpp/src/qpid/ha/ReplicatingSubscription.cpp b/qpid/cpp/src/qpid/ha/ReplicatingSubscription.cpp
index a0cfa393aa..908458fad3 100644
--- a/qpid/cpp/src/qpid/ha/ReplicatingSubscription.cpp
+++ b/qpid/cpp/src/qpid/ha/ReplicatingSubscription.cpp
@@ -161,7 +161,7 @@ void ReplicatingSubscription::initialize() {
{
sys::Mutex::ScopedLock l(lock); // Concurrent calls to dequeued()
dequeues += initDequeues; // Messages on backup that are not on primary.
- skip = backupIds - initDequeues; // Messages already on the backup.
+ skipEnqueue = backupIds - initDequeues; // Messages already on the backup.
// Queue front is moving but we know this subscriptions will start at a
// position >= front so if front is safe then position must be.
position = front;
@@ -169,7 +169,7 @@ void ReplicatingSubscription::initialize() {
QPID_LOG(debug, logPrefix << "Subscribed: front " << front
<< ", back " << back
<< ", guarded " << guard->getFirst()
- << ", on backup " << skip);
+ << ", on backup " << skipEnqueue);
checkReady(l);
}
@@ -215,9 +215,9 @@ bool ReplicatingSubscription::deliver(
position = m.getSequence();
try {
bool result = false;
- if (skip.contains(id)) {
+ if (skipEnqueue.contains(id)) {
QPID_LOG(trace, logPrefix << "Skip " << LogMessageId(*getQueue(), m));
- skip -= id;
+ skipEnqueue -= id;
guard->complete(id); // This will never be acknowledged.
notify();
result = true;
@@ -281,6 +281,9 @@ void ReplicatingSubscription::acknowledged(const broker::DeliveryRecord& r) {
// Called with lock held. Called in subscription's connection thread.
void ReplicatingSubscription::sendDequeueEvent(Mutex::ScopedLock& l)
{
+ ReplicationIdSet oldDequeues = dequeues;
+ dequeues -= skipDequeue; // Don't send skipped dequeues
+ skipDequeue -= oldDequeues; // Forget dequeues that would have been sent.
if (dequeues.empty()) return;
QPID_LOG(trace, logPrefix << "Sending dequeues " << dequeues);
sendEvent(DequeueEvent(dequeues), l);
@@ -332,9 +335,14 @@ bool ReplicatingSubscription::doDispatch()
}
}
-void ReplicatingSubscription::addSkip(const ReplicationIdSet& ids) {
+void ReplicatingSubscription::skipEnqueues(const ReplicationIdSet& ids) {
Mutex::ScopedLock l(lock);
- skip += ids;
+ skipEnqueue += ids;
+}
+
+void ReplicatingSubscription::skipDequeues(const ReplicationIdSet& ids) {
+ Mutex::ScopedLock l(lock);
+ skipDequeue += ids;
}
}} // namespace qpid::ha