From 7bb93561f1b60aa1ceef5a2ecc7635a02d4cf4c7 Mon Sep 17 00:00:00 2001 From: Rahmeen14 Date: Tue, 18 Aug 2026 14:39:21 +0100 Subject: [PATCH] overload the send function in TracingProducerImpl --- src/rmq/rmqa/rmqa_tracingproducerimpl.cpp | 27 +++++++++++++ src/rmq/rmqa/rmqa_tracingproducerimpl.h | 6 +++ src/tests/rmqa/rmqa_producerimpl.t.cpp | 47 +++++++++++++++++++++++ 3 files changed, 80 insertions(+) diff --git a/src/rmq/rmqa/rmqa_tracingproducerimpl.cpp b/src/rmq/rmqa/rmqa_tracingproducerimpl.cpp index 7d4b60e9..3f2682a1 100644 --- a/src/rmq/rmqa/rmqa_tracingproducerimpl.cpp +++ b/src/rmq/rmqa/rmqa_tracingproducerimpl.cpp @@ -105,6 +105,33 @@ rmqp::Producer::SendStatus TracingProducerImpl::send( timeout); } +rmqp::Producer::SendStatus TracingProducerImpl::send( + const rmqt::Message& message, + const bsl::string& routingKey, + rmqt::Mandatory::Value mandatoryFlag, + const rmqp::Producer::ConfirmationCallback& confirmCallback, + const bsls::TimeInterval& timeout) +{ + rmqt::Message newMessage(message); + // ideally we'd have move semantics on the message to avoid + // copying the metadata note that this is not a deep copy of + // the message payload + bsl::shared_ptr context = + d_tracing->createAndTag( + &(newMessage.properties()), routingKey, d_exchangeName, d_endpoint); + + return ProducerImpl::send(newMessage, + routingKey, + mandatoryFlag, + bdlf::BindUtil::bind(&callbackAndContext, + confirmCallback, + context, + bdlf::PlaceHolders::_1, + bdlf::PlaceHolders::_2, + bdlf::PlaceHolders::_3), + timeout); +} + rmqp::Producer::SendStatus TracingProducerImpl::trySend( const rmqt::Message& message, const bsl::string& routingKey, diff --git a/src/rmq/rmqa/rmqa_tracingproducerimpl.h b/src/rmq/rmqa/rmqa_tracingproducerimpl.h index 462e80fa..1df085d5 100644 --- a/src/rmq/rmqa/rmqa_tracingproducerimpl.h +++ b/src/rmq/rmqa/rmqa_tracingproducerimpl.h @@ -64,6 +64,12 @@ class TracingProducerImpl : public ProducerImpl { const rmqp::Producer::ConfirmationCallback& confirmCallback, const bsls::TimeInterval& timeout) BSLS_KEYWORD_OVERRIDE; + SendStatus send(const rmqt::Message& message, + const bsl::string& routingKey, + rmqt::Mandatory::Value mandatoryFlag, + const rmqp::Producer::ConfirmationCallback& confirmCallback, + const bsls::TimeInterval& timeout) BSLS_KEYWORD_OVERRIDE; + SendStatus trySend(const rmqt::Message& message, const bsl::string& routingKey, diff --git a/src/tests/rmqa/rmqa_producerimpl.t.cpp b/src/tests/rmqa/rmqa_producerimpl.t.cpp index db3b1331..be481c45 100644 --- a/src/tests/rmqa/rmqa_producerimpl.t.cpp +++ b/src/tests/rmqa/rmqa_producerimpl.t.cpp @@ -867,6 +867,53 @@ TEST_P(TracingProducerImplTests, SendConfirmCallsTracing) d_threadPool.drain(); } +TEST_P(TracingProducerImplTests, SendWithMandatoryFlagConfirmCallsTracing) +{ + // Regression: the send() overload accepting an explicit mandatory flag must + // also create a tracing context and tag the message, just like the + // four-argument send(). Previously TracingProducerImpl only overrode the + // four-argument send(), so publishing with an explicit mandatory flag + // silently bypassed tracing. + + // GIVEN + bsl::shared_ptr tracingContext( + bsl::make_shared()); + rmqt::Properties specialProperties = d_message.properties(); + specialProperties.headers = bsl::make_shared(); + specialProperties.headers->insert( + bsl::make_pair(bsl::string("special"), bsl::string("property"))); + + bsl::shared_ptr producer(d_factory->create( + 1, d_exchange, d_mockSendChannel, d_threadPool, d_eventLoop)); + + // The tagged message must be published, and the caller-supplied mandatory + // flag must be forwarded unchanged to the channel. + EXPECT_CALL(*d_mockSendChannel, + publishMessage(MessagePropertiesMatch(specialProperties), + bsl::string("routingKey"), + rmqt::Mandatory::DISCARD_UNROUTABLE)); + EXPECT_CALL( + *d_tracing, + createAndTag( + _, bsl::string("routingKey"), bsl::string("test-exchange"), _)) + .WillOnce( + DoAll(SetArgPointee<0>(specialProperties), Return(tracingContext))); + // WHEN + producer->send(d_message, + "routingKey", + rmqt::Mandatory::DISCARD_UNROUTABLE, + d_callback, + d_timeout); + + const rmqt::ConfirmResponse confirmResponse(rmqt::ConfirmResponse::ACK); + + EXPECT_CALL(*tracingContext, response(confirmResponse)).WillOnce(Return()); + + d_injectConfirm(d_message, d_exchange->name(), confirmResponse); + + d_threadPool.drain(); +} + RMQTESTUTIL_TESTSUITE_P(AllMembers, ProducerImplTests, Values(PRODUCER, TRACING_PRODUCER),