summaryrefslogtreecommitdiff
path: root/cpp/src/qpid/ha/ReplicatingSubscription.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'cpp/src/qpid/ha/ReplicatingSubscription.cpp')
-rw-r--r--cpp/src/qpid/ha/ReplicatingSubscription.cpp64
1 files changed, 36 insertions, 28 deletions
diff --git a/cpp/src/qpid/ha/ReplicatingSubscription.cpp b/cpp/src/qpid/ha/ReplicatingSubscription.cpp
index 8c88382b94..9d7df51e3b 100644
--- a/cpp/src/qpid/ha/ReplicatingSubscription.cpp
+++ b/cpp/src/qpid/ha/ReplicatingSubscription.cpp
@@ -20,6 +20,7 @@
*/
#include "ReplicatingSubscription.h"
+#include "HaBroker.h"
#include "Primary.h"
#include "qpid/broker/Queue.h"
#include "qpid/broker/SessionContext.h"
@@ -40,6 +41,7 @@ using namespace std;
const string ReplicatingSubscription::QPID_REPLICATING_SUBSCRIPTION("qpid.replicating-subscription");
const string ReplicatingSubscription::QPID_HIGH_SEQUENCE_NUMBER("qpid.high-sequence-number");
const string ReplicatingSubscription::QPID_LOW_SEQUENCE_NUMBER("qpid.low-sequence-number");
+const string ReplicatingSubscription::QPID_BROKER_INFO("qpid.broker-info");
namespace {
const string DOLLAR("$");
@@ -59,7 +61,7 @@ class DequeueRemover
}
void operator()(const QueuedMessage& message) {
- if (message.position >= start && message.position <= end) {
+ if (message.position >= start && message.position <= end) {
//i.e. message is within the intial range and has not been dequeued,
//so remove it from the dequeues
dequeues.remove(message.position);
@@ -127,7 +129,7 @@ ReplicatingSubscription::Factory::create(
boost::shared_ptr<ReplicatingSubscription> rs;
if (arguments.isSet(QPID_REPLICATING_SUBSCRIPTION)) {
rs.reset(new ReplicatingSubscription(
- LogPrefix(haBroker),
+ haBroker,
parent, name, queue, ack, acquire, exclusive, tag,
resumeId, resumeTtl, arguments));
queue->addObserver(rs);
@@ -173,7 +175,7 @@ ostream& operator<<(ostream& o, const QueueRange& qr) {
}
ReplicatingSubscription::ReplicatingSubscription(
- LogPrefix lp,
+ HaBroker& hb,
SemanticState* parent,
const string& name,
Queue::shared_ptr queue,
@@ -186,17 +188,24 @@ ReplicatingSubscription::ReplicatingSubscription(
const framing::FieldTable& arguments
) : ConsumerImpl(parent, name, queue, ack, acquire, exclusive, tag,
resumeId, resumeTtl, arguments),
- logPrefix(lp, queue->getName()),
+ haBroker(hb),
+ logPrefix(hb),
dummy(new Queue(mask(name))),
ready(false)
{
try {
- // FIXME aconway 2012-05-22: use hostname from brokerinfo
- // Separate the remote part from a "local-remote" address for logging.
- string address = parent->getSession().getConnection().getUrl();
- size_t i = address.find('-');
- if (i != string::npos) address = address.substr(i+1);
- logSuffix = " (" + address + ")";
+ // Set a log prefix message that identifies the remote broker.
+ // FIXME aconway 2012-05-24: use URL instead of host:port, include transport?
+ ostringstream os;
+ os << queue->getName() << "@";
+ FieldTable ft;
+ if (arguments.getTable(ReplicatingSubscription::QPID_BROKER_INFO, ft)) {
+ BrokerInfo info(ft);
+ os << info.getHostName() << ":" << info.getPort();
+ }
+ else
+ os << parent->getSession().getConnection().getUrl();
+ logPrefix.setMessage(os.str());
QueueRange primary(*queue);
QueueRange backup(arguments);
@@ -228,16 +237,16 @@ ReplicatingSubscription::ReplicatingSubscription(
<< " backup range " << backup
<< " primary range " << primary
<< " position " << position
- << " dequeues " << dequeues << logSuffix);
+ << " dequeues " << dequeues);
}
catch (const std::exception& e) {
throw Exception(QPID_MSG(logPrefix << "Error setting up replication: "
- << e.what() << logSuffix));
+ << e.what()));
}
}
ReplicatingSubscription::~ReplicatingSubscription() {
- QPID_LOG(debug, logPrefix << "Detroyed replicating subscription" << logSuffix);
+ QPID_LOG(debug, logPrefix << "Detroyed replicating subscription");
}
// Called in subscription's connection thread when the subscription is created.
@@ -258,7 +267,7 @@ void ReplicatingSubscription::initialize() {
}
else {
QPID_LOG(debug, logPrefix << "Backup subscription catching up from "
- << position << " to " << readyPosition << logSuffix);
+ << position << " to " << readyPosition);
}
}
@@ -267,7 +276,7 @@ bool ReplicatingSubscription::deliver(QueuedMessage& qm) {
try {
// Add position events for the subscribed queue, not the internal event queue.
if (qm.queue == getQueue().get()) {
- QPID_LOG(trace, logPrefix << "replicating " << qm << logSuffix);
+ QPID_LOG(trace, logPrefix << "Replicating " << qm);
{
sys::Mutex::ScopedLock l(lock);
assert(position == qm.position);
@@ -296,8 +305,8 @@ bool ReplicatingSubscription::deliver(QueuedMessage& qm) {
else
return ConsumerImpl::deliver(qm); // Message is for internal event queue.
} catch (const std::exception& e) {
- QPID_LOG(critical, logPrefix << "error replicating " << qm
- << logSuffix << ": " << e.what());
+ QPID_LOG(critical, logPrefix << "Error replicating " << qm
+ << ": " << e.what());
throw;
}
}
@@ -305,7 +314,7 @@ bool ReplicatingSubscription::deliver(QueuedMessage& qm) {
void ReplicatingSubscription::setReady(const sys::Mutex::ScopedLock&) {
if (ready) return;
ready = true;
- QPID_LOG(info, logPrefix << "Caught up at " << getPosition() << logSuffix);
+ QPID_LOG(info, logPrefix << "Caught up at " << getPosition());
// Notify Primary that a subscription is ready.
if (Primary::get()) Primary::get()->readyReplica(getQueue()->getName());
}
@@ -319,7 +328,7 @@ void ReplicatingSubscription::complete(
{
// Handle completions for the subscribed queue, not the internal event queue.
if (qm.queue == getQueue().get()) {
- QPID_LOG(trace, logPrefix << "completed " << qm << logSuffix);
+ QPID_LOG(trace, logPrefix << "Completed " << qm);
Delayed::iterator i= delayed.find(qm.position);
// The same message can be completed twice, by acknowledged and
// dequeued, remove it from the set so it only gets completed
@@ -337,7 +346,7 @@ void ReplicatingSubscription::complete(
// Called in arbitrary connection thread *with the queue lock held*
void ReplicatingSubscription::enqueued(const QueuedMessage& qm) {
// Delay completion
- QPID_LOG(trace, logPrefix << "delaying completion of " << qm << logSuffix);
+ QPID_LOG(trace, logPrefix << "Delaying completion of " << qm);
qm.payload->getIngressCompletion().startCompleter();
{
sys::Mutex::ScopedLock l(lock);
@@ -350,7 +359,7 @@ void ReplicatingSubscription::enqueued(const QueuedMessage& qm) {
void ReplicatingSubscription::cancelComplete(
const Delayed::value_type& v, const sys::Mutex::ScopedLock&)
{
- QPID_LOG(trace, logPrefix << "cancel completed " << v.second << logSuffix);
+ QPID_LOG(trace, logPrefix << "Cancel completed " << v.second);
v.second.payload->getIngressCompletion().finishCompleter();
}
@@ -361,8 +370,8 @@ void ReplicatingSubscription::cancel()
boost::dynamic_pointer_cast<QueueObserver>(shared_from_this()));
{
sys::Mutex::ScopedLock l(lock);
- QPID_LOG(debug, logPrefix << "cancel backup subscription to "
- << getQueue()->getName() << logSuffix);
+ QPID_LOG(debug, logPrefix << "Cancel backup subscription to "
+ << getQueue()->getName());
for_each(delayed.begin(), delayed.end(),
boost::bind(&ReplicatingSubscription::cancelComplete, this, _1, boost::ref(l)));
delayed.clear();
@@ -385,8 +394,7 @@ bool ReplicatingSubscription::hideDeletedError() { return true; }
void ReplicatingSubscription::sendDequeueEvent(const sys::Mutex::ScopedLock&)
{
if (dequeues.empty()) return;
- QPID_LOG(trace, logPrefix << "sending dequeues " << dequeues
- << " from " << getQueue()->getName() << logSuffix);
+ QPID_LOG(trace, logPrefix << "Sending dequeues " << dequeues);
string buf(dequeues.encodedSize(),'\0');
framing::Buffer buffer(&buf[0], buf.size());
dequeues.encode(buffer);
@@ -401,7 +409,7 @@ void ReplicatingSubscription::sendDequeueEvent(const sys::Mutex::ScopedLock&)
void ReplicatingSubscription::dequeued(const QueuedMessage& qm)
{
{
- QPID_LOG(trace, logPrefix << "dequeued " << qm << logSuffix);
+ QPID_LOG(trace, logPrefix << "Dequeued " << qm);
sys::Mutex::ScopedLock l(lock);
dequeues.add(qm.position);
// If we have not yet sent this message to the backup, then
@@ -415,8 +423,8 @@ void ReplicatingSubscription::dequeued(const QueuedMessage& qm)
void ReplicatingSubscription::sendPositionEvent(SequenceNumber pos, const sys::Mutex::ScopedLock&)
{
if (pos == backupPosition) return; // No need to send.
- QPID_LOG(trace, logPrefix << "sending position " << pos
- << ", was " << backupPosition << logSuffix);
+ QPID_LOG(trace, logPrefix << "Sending position " << pos
+ << ", was " << backupPosition);
string buf(pos.encodedSize(),'\0');
framing::Buffer buffer(&buf[0], buf.size());
pos.encode(buffer);