Skip to content

Commit 262fb89

Browse files
committed
Perf[MQB]: use concrete dispatcher event types
Signed-off-by: Evgeny Malygin <emalygin@bloomberg.net>
1 parent a99ce85 commit 262fb89

78 files changed

Lines changed: 5157 additions & 2029 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

src/groups/mqb/CMakeLists.txt

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ set_property(SOURCE "mqba/mqba_application.cpp"
2020
APPEND
2121
PROPERTY COMPILE_DEFINITIONS "BMQ_BUILD_TYPE=${CMAKE_BUILD_TYPE}")
2222

23-
set(MQB_PRIVATE_PACKAGES mqba mqbblp mqbc mqbcmd mqbconfm mqbi mqbmock mqbnet mqbs mqbsi mqbsl mqbu)
23+
set(MQB_PRIVATE_PACKAGES mqba mqbblp mqbc mqbcmd mqbconfm mqbevt mqbi mqbmock mqbnet mqbs mqbsi mqbsl mqbu)
2424
target_bmq_style_uor( mqb PRIVATE_PACKAGES ${MQB_PRIVATE_PACKAGES})
2525

2626
# Turn off warnings in generated code.

src/groups/mqb/group/mqb.mem

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ mqbc
55
mqbcfg
66
mqbcmd
77
mqbconfm
8+
mqbevt
89
mqbi
910
mqbmock
1011
mqbnet

src/groups/mqb/mqba/mqba_adminsession.cpp

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,9 @@
7272
#include <bsls_performancehint.h>
7373
#include <bsls_timeinterval.h>
7474

75+
// MQB
76+
#include <mqbevt_callbackevent.h>
77+
7578
namespace BloombergLP {
7679
namespace mqba {
7780

@@ -443,8 +446,8 @@ void AdminSession::onDispatcherEvent(const mqbi::DispatcherEvent& event)
443446

444447
switch (event.type()) {
445448
case mqbi::DispatcherEventType::e_CALLBACK: {
446-
const mqbi::DispatcherCallbackEvent* realEvent =
447-
event.asCallbackEvent();
449+
const mqbevt::CallbackEvent* const realEvent =
450+
event.the<mqbevt::CallbackEvent>();
448451

449452
BSLS_ASSERT_SAFE(!realEvent->callback().empty());
450453
realEvent->callback()();

src/groups/mqb/mqba/mqba_clientsession.cpp

Lines changed: 45 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -151,6 +151,12 @@
151151
#include <mqbblp_clustercatalog.h>
152152
#include <mqbblp_queueengineutil.h>
153153
#include <mqbcfg_brokerconfig.h>
154+
#include <mqbevt_ackevent.h>
155+
#include <mqbevt_callbackevent.h>
156+
#include <mqbevt_confirmevent.h>
157+
#include <mqbevt_pushevent.h>
158+
#include <mqbevt_putevent.h>
159+
#include <mqbevt_rejectevent.h>
154160
#include <mqbi_cluster.h>
155161
#include <mqbi_queue.h>
156162
#include <mqbnet_tcpsessionfactory.h>
@@ -1376,7 +1382,7 @@ void ClientSession::processConfigureStream(
13761382
opLogger));
13771383
}
13781384

1379-
void ClientSession::onAckEvent(const mqbi::DispatcherAckEvent& event)
1385+
void ClientSession::onAckEvent(const mqbevt::AckEvent& event)
13801386
{
13811387
// executed by the *CLIENT* dispatcher thread
13821388

@@ -1472,7 +1478,7 @@ void ClientSession::onAckEvent(const mqbi::DispatcherAckEvent& event)
14721478
"onAckEvent");
14731479
}
14741480

1475-
void ClientSession::onConfirmEvent(const mqbi::DispatcherConfirmEvent& event)
1481+
void ClientSession::onConfirmEvent(const mqbevt::ConfirmEvent& event)
14761482
{
14771483
// executed by the *CLIENT* dispatcher thread
14781484

@@ -1556,7 +1562,7 @@ void ClientSession::onConfirmEvent(const mqbi::DispatcherConfirmEvent& event)
15561562
}
15571563
}
15581564

1559-
void ClientSession::onRejectEvent(const mqbi::DispatcherRejectEvent& event)
1565+
void ClientSession::onRejectEvent(const mqbevt::RejectEvent& event)
15601566
{
15611567
// executed by the *CLIENT* dispatcher thread
15621568

@@ -1720,7 +1726,7 @@ bool ClientSession::validateMessage(mqbi::QueueHandle** queueHandle,
17201726
return true;
17211727
}
17221728

1723-
void ClientSession::onPushEvent(const mqbi::DispatcherPushEvent& event)
1729+
void ClientSession::onPushEvent(const mqbevt::PushEvent& event)
17241730
{
17251731
// executed by the *CLIENT* dispatcher thread
17261732

@@ -1917,7 +1923,7 @@ void ClientSession::onPushEvent(const mqbi::DispatcherPushEvent& event)
19171923
}
19181924
}
19191925

1920-
void ClientSession::onPutEvent(const mqbi::DispatcherPutEvent& event)
1926+
void ClientSession::onPutEvent(const mqbevt::PutEvent& event)
19211927
{
19221928
// executed by the *CLIENT* dispatcher thread
19231929

@@ -2641,33 +2647,45 @@ void ClientSession::processEvent(const bmqp::Event& event,
26412647
}
26422648

26432649
// Not a control or leader message, it's either a put or a confirm ..
2644-
mqbi::DispatcherEventType::Enum eventType;
26452650

2651+
bsl::shared_ptr<bdlbb::Blob> blobSp =
2652+
d_state.d_blobSpPool_p->getObject();
2653+
*blobSp = *(event.blob());
2654+
2655+
// Dispatch the event
2656+
// TODO(678098): revisit, use per-IO thread event source
26462657
if (event.isPutEvent()) {
2647-
eventType = mqbi::DispatcherEventType::e_PUT;
2658+
bsl::shared_ptr<mqbevt::PutEvent> event_sp =
2659+
dispatcher()
2660+
->getDefaultEventSource()
2661+
->getEvent<mqbevt::PutEvent>();
2662+
event_sp->setBlob(blobSp).setSource(this);
2663+
dispatcher()->dispatchEvent(bslmf::MovableRefUtil::move(event_sp),
2664+
this);
26482665
}
26492666
else if (event.isConfirmEvent()) {
2650-
eventType = mqbi::DispatcherEventType::e_CONFIRM;
2667+
bsl::shared_ptr<mqbevt::ConfirmEvent> event_sp =
2668+
dispatcher()
2669+
->getDefaultEventSource()
2670+
->getEvent<mqbevt::ConfirmEvent>();
2671+
event_sp->setBlob(blobSp).setSource(this);
2672+
dispatcher()->dispatchEvent(bslmf::MovableRefUtil::move(event_sp),
2673+
this);
26512674
}
26522675
else if (event.isRejectEvent()) {
2653-
eventType = mqbi::DispatcherEventType::e_REJECT;
2676+
bsl::shared_ptr<mqbevt::RejectEvent> event_sp =
2677+
dispatcher()
2678+
->getDefaultEventSource()
2679+
->getEvent<mqbevt::RejectEvent>();
2680+
event_sp->setBlob(blobSp).setSource(this);
2681+
dispatcher()->dispatchEvent(bslmf::MovableRefUtil::move(event_sp),
2682+
this);
26542683
}
26552684
else {
26562685
BALL_LOG_ERROR << "#CLIENT_UNEXPECTED_EVENT " << description()
26572686
<< ": Unexpected event type: " << event;
26582687
return; // RETURN
26592688
}
2660-
2661-
// Dispatch the event
2662-
// TODO(678098): revisit, use per-IO thread event source
2663-
mqbi::Dispatcher::DispatcherEventSp dispEvent =
2664-
dispatcher()->getDefaultEventSource()->getEvent();
2665-
bsl::shared_ptr<bdlbb::Blob> blobSp =
2666-
d_state.d_blobSpPool_p->getObject();
2667-
*blobSp = *(event.blob());
2668-
(*dispEvent).setType(eventType).setSource(this).setBlob(blobSp);
2669-
dispatcher()->dispatchEvent(bslmf::MovableRefUtil::move(dispEvent),
2670-
this);
26712689
}
26722690
}
26732691

@@ -2824,23 +2842,23 @@ void ClientSession::onDispatcherEvent(const mqbi::DispatcherEvent& event)
28242842

28252843
switch (event.type()) {
28262844
case mqbi::DispatcherEventType::e_CONFIRM: {
2827-
onConfirmEvent(*(event.asConfirmEvent()));
2845+
onConfirmEvent(*(event.the<mqbevt::ConfirmEvent>()));
28282846
} break;
28292847
case mqbi::DispatcherEventType::e_REJECT: {
2830-
onRejectEvent(*(event.asRejectEvent()));
2848+
onRejectEvent(*(event.the<mqbevt::RejectEvent>()));
28312849
} break;
28322850
case mqbi::DispatcherEventType::e_PUSH: {
2833-
onPushEvent(*(event.asPushEvent()));
2851+
onPushEvent(*(event.the<mqbevt::PushEvent>()));
28342852
} break;
28352853
case mqbi::DispatcherEventType::e_PUT: {
2836-
onPutEvent(*(event.asPutEvent()));
2854+
onPutEvent(*(event.the<mqbevt::PutEvent>()));
28372855
} break;
28382856
case mqbi::DispatcherEventType::e_ACK: {
2839-
onAckEvent(*(event.asAckEvent()));
2857+
onAckEvent(*(event.the<mqbevt::AckEvent>()));
28402858
} break;
28412859
case mqbi::DispatcherEventType::e_CALLBACK: {
2842-
const mqbi::DispatcherCallbackEvent* realEvent =
2843-
event.asCallbackEvent();
2860+
const mqbevt::CallbackEvent* const realEvent =
2861+
event.the<mqbevt::CallbackEvent>();
28442862

28452863
BSLS_ASSERT_SAFE(!realEvent->callback().empty());
28462864
flush(); // Flush any pending messages to guarantee ordering of events

src/groups/mqb/mqba/mqba_clientsession.h

Lines changed: 12 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,13 @@ class ClusterCatalog;
9292
namespace bmqst {
9393
class StatContext;
9494
}
95+
namespace mqbevt {
96+
class AckEvent;
97+
class ConfirmEvent;
98+
class PushEvent;
99+
class PutEvent;
100+
class RejectEvent;
101+
}
95102

96103
namespace mqba {
97104

@@ -486,19 +493,19 @@ class ClientSession : public mqbnet::Session,
486493
const bsl::shared_ptr<bmqsys::OperationLogger>& opLogger);
487494

488495
/// Process the specified ack `event`.
489-
void onAckEvent(const mqbi::DispatcherAckEvent& event);
496+
void onAckEvent(const mqbevt::AckEvent& event);
490497

491498
/// Process the specified confirm `event`.
492-
void onConfirmEvent(const mqbi::DispatcherConfirmEvent& event);
499+
void onConfirmEvent(const mqbevt::ConfirmEvent& event);
493500

494501
/// Process the specified reject `event`.
495-
void onRejectEvent(const mqbi::DispatcherRejectEvent& event);
502+
void onRejectEvent(const mqbevt::RejectEvent& event);
496503

497504
/// Process the specified push `event` received from the dispatcher.
498-
void onPushEvent(const mqbi::DispatcherPushEvent& event);
505+
void onPushEvent(const mqbevt::PushEvent& event);
499506

500507
/// Process the specified put `event`.
501-
void onPutEvent(const mqbi::DispatcherPutEvent& event);
508+
void onPutEvent(const mqbevt::PutEvent& event);
502509

503510
/// Validate a message of the specified `eventType` using the specified
504511
/// `queueId`. Return true if the message is valid and false otherwise.

src/groups/mqb/mqba/mqba_clientsession.t.cpp

Lines changed: 35 additions & 47 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,9 @@
1919
// MQB
2020
#include <mqbcfg_brokerconfig.h>
2121
#include <mqbcfg_messages.h>
22+
#include <mqbevt_ackevent.h>
23+
#include <mqbevt_pushevent.h>
24+
#include <mqbevt_putevent.h>
2225
#include <mqbi_queue.h>
2326
#include <mqbmock_cluster.h>
2427
#include <mqbmock_dispatcher.h>
@@ -641,10 +644,6 @@ T assertFail()
641644

642645
/// The `TestBench` holds system components together.
643646
class TestBench {
644-
private:
645-
// PRIVATE TYPES
646-
typedef mqbmock::Dispatcher::EventGuard EventGuard;
647-
648647
public:
649648
// DATA
650649
bdlbb::PooledBlobBufferFactory d_bufferFactory;
@@ -824,13 +823,11 @@ class TestBench {
824823
guid,
825824
queueId);
826825

827-
mqbi::Dispatcher::DispatcherEventSp event =
828-
bsl::allocate_shared<mqbi::DispatcherEvent>(d_allocator_p);
829-
(*event)
830-
.setType(mqbi::DispatcherEventType::e_ACK)
831-
.setAckMessage(ackMessage);
826+
bsl::shared_ptr<mqbevt::AckEvent> event_sp =
827+
bsl::allocate_shared<mqbevt::AckEvent>(d_allocator_p);
828+
(*event_sp).setAckMessage(ackMessage);
832829

833-
dispatch(event);
830+
dispatch(event_sp);
834831
}
835832

836833
/// Sends a `Put` event for the specified `queueId`, `msgGUID` and
@@ -870,16 +867,6 @@ class TestBench {
870867
&d_bufferFactory,
871868
cat);
872869

873-
mqbi::Dispatcher::DispatcherEventSp event =
874-
bsl::allocate_shared<mqbi::DispatcherEvent>(d_allocator_p);
875-
(*event)
876-
.setType(mqbi::DispatcherEventType::e_PUT)
877-
.setIsRelay(true) // Relay message
878-
.setSource(&d_cs) // DispatcherClient *value
879-
.setPutHeader(putHeader)
880-
.setBlob(eventBlob) // const bsl::shared_ptr<bdlbb::Blob>& value
881-
.setCompressionAlgorithmType(cat);
882-
883870
// Internal-ticket D167598037.
884871
// Verify that PutMessageIterator does not change the input.
885872
bmqp::Event rawEvent(eventBlob.get(),
@@ -898,7 +885,12 @@ class TestBench {
898885
BSLS_ASSERT(pIt2.next());
899886
BSLS_ASSERT(pIt3.next());
900887

901-
dispatch(event);
888+
bsl::shared_ptr<mqbevt::PutEvent> event_sp =
889+
bsl::allocate_shared<mqbevt::PutEvent>(d_allocator_p);
890+
(*event_sp).setPutHeader(putHeader).setBlob(eventBlob).setSource(
891+
&d_cs);
892+
893+
dispatch(event_sp);
902894
}
903895

904896
/// Sends a `Push` event for the specified `queueId`, `msgGUID` and
@@ -912,20 +904,20 @@ class TestBench {
912904
{
913905
PVV("Sending PUSH with queueId=" << queueId << ", guid=" << msgGUID);
914906

915-
mqbi::Dispatcher::DispatcherEventSp event =
916-
bsl::allocate_shared<mqbi::DispatcherEvent>(d_allocator_p);
907+
bsl::shared_ptr<mqbevt::PushEvent> event_sp =
908+
bsl::allocate_shared<mqbevt::PushEvent>(d_allocator_p);
917909

918-
(*event)
919-
.setType(mqbi::DispatcherEventType::e_PUSH)
920-
.setSource(&d_cs) // DispatcherClient *value
910+
(*event_sp)
911+
.setSource(&d_cs)
921912
.setQueueId(queueId)
922913
.setBlob(blob)
923914
.setGuid(msgGUID)
924915
.setMessagePropertiesInfo(logic)
925916
.setCompressionAlgorithmType(cat);
926917

927-
dispatch(event);
918+
dispatch(event_sp);
928919
}
920+
929921
bool validateData(const bdlbb::Blob& blob,
930922
int offset,
931923
int length = k_PAYLOAD_LENGTH)
@@ -1013,7 +1005,8 @@ class TestBench {
10131005
/// using the specified `event`.
10141006
void dispatch(mqbi::Dispatcher::DispatcherEventSp event)
10151007
{
1016-
EventGuard guard(d_mockDispatcher._withEvent(&d_cs, event));
1008+
mqbmock::Dispatcher::EventGuard guard(
1009+
d_mockDispatcher._withEvent(&d_cs, event));
10171010

10181011
d_cs.onDispatcherEvent(*event);
10191012
}
@@ -1943,9 +1936,6 @@ static void test9_newStylePush()
19431936

19441937
BMQTST_ASSERT_EQ(bmqt::EventBuilderResult::e_SUCCESS, rc);
19451938

1946-
mqbi::Dispatcher::DispatcherEventSp putEvent =
1947-
bsl::allocate_shared<mqbi::DispatcherEvent>(
1948-
bmqtst::TestHelperUtil::allocator());
19491939
bmqp::Event rawEvent(peb.blob().get(),
19501940
bmqtst::TestHelperUtil::allocator());
19511941

@@ -1957,14 +1947,15 @@ static void test9_newStylePush()
19571947
rawEvent.loadPutMessageIterator(&putIt, false);
19581948
BSLS_ASSERT(putIt.next());
19591949

1960-
(*putEvent)
1961-
.setType(mqbi::DispatcherEventType::e_PUT)
1962-
.setIsRelay(true) // Relay message
1963-
.setSource(&tb.d_cs) // DispatcherClient *value
1950+
bsl::shared_ptr<mqbevt::PutEvent> event_sp =
1951+
bsl::allocate_shared<mqbevt::PutEvent>(
1952+
bmqtst::TestHelperUtil::allocator());
1953+
(*event_sp)
1954+
.setSource(&tb.d_cs)
19641955
.setPutHeader(putIt.header())
1965-
.setBlob(peb.blob()); // const bsl::shared_ptr<bdlbb::Blob>& value
1956+
.setBlob(peb.blob());
19661957

1967-
tb.dispatch(putEvent);
1958+
tb.dispatch(event_sp);
19681959

19691960
// Check if a message was sent
19701961
const bsl::vector<MyMockQueueHandle::Post>& postMessages =
@@ -2060,9 +2051,6 @@ static void test10_newStyleCompressedPush()
20602051

20612052
BMQTST_ASSERT_EQ(bmqt::EventBuilderResult::e_SUCCESS, rc);
20622053

2063-
mqbi::Dispatcher::DispatcherEventSp putEvent =
2064-
bsl::allocate_shared<mqbi::DispatcherEvent>(
2065-
bmqtst::TestHelperUtil::allocator());
20662054
bmqp::Event rawEvent(peb.blob().get(),
20672055
bmqtst::TestHelperUtil::allocator());
20682056

@@ -2074,15 +2062,15 @@ static void test10_newStyleCompressedPush()
20742062
rawEvent.loadPutMessageIterator(&putIt, false);
20752063
BSLS_ASSERT(putIt.next());
20762064

2077-
(*putEvent)
2078-
.setType(mqbi::DispatcherEventType::e_PUT)
2079-
.setIsRelay(true) // Relay message
2080-
.setSource(&tb.d_cs) // DispatcherClient *value
2065+
bsl::shared_ptr<mqbevt::PutEvent> event_sp =
2066+
bsl::allocate_shared<mqbevt::PutEvent>(
2067+
bmqtst::TestHelperUtil::allocator());
2068+
(*event_sp)
2069+
.setSource(&tb.d_cs)
20812070
.setPutHeader(putIt.header())
2082-
.setBlob(peb.blob()) // const bsl::shared_ptr<bdlbb::Blob>& value
2083-
.setCompressionAlgorithmType(bmqt::CompressionAlgorithmType::e_ZLIB);
2071+
.setBlob(peb.blob());
20842072

2085-
tb.dispatch(putEvent);
2073+
tb.dispatch(event_sp);
20862074

20872075
// Check if a message was sent
20882076
const bsl::vector<MyMockQueueHandle::Post>& postMessages =

0 commit comments

Comments
 (0)