This sample demonstrates how to measure the handler latency.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19#include <iostream>
20
26
27#include "RdiListener.h"
28#include "EmdiListener.h"
29#include "EobiListener.h"
30
31#include "../Common/Signal.h"
32#include "../Common/Options.h"
33#include "../Common/Settings.h"
34#include "../Common/Utils.h"
35
38
41
42
44int main(int argc, char* argv[])
45{
46
47 const Samples::AppConfiguration<
48 Samples::NetworkInterfaceConfiguration
49 , Samples::EnvironmentConfiguration
50 , Samples::LogDirectoryConfiguration
51 , Samples::AffinityConfiguration
52 , Samples::MarketSegmentConfiguration
53 , Samples::ProductConfiguration
54 , Samples::PacketCountConfiguration
55 > cfg{"Benchmark", argc, argv};
56
57 try
58 {
59 Samples::SignalHelper::manageSignals();
60
63
64 auto rdiSettings = Samples::makeSettings<RdiHandlerSettings>(cfg);
65 fillEnvironment(rdiSettings, cfg);
66
67 RdiListener rdiListener;
69 rdiHandler.bindFeedEngine(feedEngine);
70
71 rdiHandler.registerErrorListener(&rdiListener);
72 rdiHandler.registerWarningListener(&rdiListener);
73 rdiHandler.registerReferenceDataListener(&rdiListener);
74
75 std::clog << "Will start the RDI Handler ..." << std::endl;
76 std::clog << "Please press Ctrl+C key to stop..." << std::endl;
77
78 rdiHandler.start();
79
80 while(!rdiListener.referenceDataReceived() && !Samples::SignalHelper::interruptDetected())
82
83 std::clog << "Will stop the RDI Handler ..." << std::endl;
84 rdiHandler.stop();
85
87
88 if(Samples::SignalHelper::interruptDetected())
89 return 0;
90
92
93 switch (cfg.product())
94 {
95 case Samples::Product::Emdi:
96 process(rdiHandler.findEmdiDescriptors(mktSeg), feedEngine, cfg);
97 break;
98
99 case Samples::Product::Eobi:
100 process(rdiHandler.findEobiDescriptors(mktSeg), feedEngine, cfg);
101 break;
102 }
103 }
104 catch(const std::exception& ex)
105 {
106 std::cerr << "EXCEPTION: " << ex.what() << std::endl;
107 }
108
109 return 0;
110}
111
112template <typename Settings, typename Configuration, typename Descriptor>
113Settings makeBmSettings(const Configuration& cfg, const Descriptor& descriptor)
114{
115 auto settings = Samples::makeSettings<Settings>(cfg);
116
117 settings.buildInternalOrderBooks = false;
120 settings.interfaceDescriptor.incrementalFeed = descriptor.incrementalFeed;
121 settings.interfaceDescriptor.snapshotFeed = descriptor.snapshotFeed;
122
123 return settings;
124}
125
126template <typename Configuration>
128{
129 if(descriptors.empty())
130 throw std::runtime_error("No products found.");
131
132 const auto& descriptor = descriptors[0];
133
134 EmdiListener listener{cfg.packetsCount()};
135
136 const auto settings = makeBmSettings<EmdiHandlerSettings>(cfg, descriptor);
137 auto handler = Samples::makeUnique<EmdiHandler>(settings);
138
139 handler->bindFeedEngine(feedEngine);
140
141 handler->setPartitionIdFilters(descriptor.partitionIdFilters);
142 handler->setMarketSegmentIdFilters(descriptor.marketSegmentIdFilters);
143
144 listener.subscribe(*handler);
145
146 handler->start();
147
148 while(!Samples::SignalHelper::interruptDetected() || listener.done())
150
151 handler->stop();
152
154
155 listener.saveLatencies("emdi-results.csv");
156 listener.processLatencies();
157}
158
159template <typename Configuration>
161{
162 if(descriptors.empty())
163 throw std::runtime_error("No products found.");
164
165 const auto& descriptor = descriptors[0];
166
167 EobiListener listener{cfg.packetsCount()};
168
169 const auto settings = makeBmSettings<EobiHandlerSettings>(cfg, descriptor);
170 auto handler = Samples::makeUnique<EobiHandler>(settings);
171
172 handler->bindFeedEngine(feedEngine);
173
174 handler->setPartitionIdFilters(descriptor.partitionIdFilters);
175 handler->setMarketSegmentIdFilters(descriptor.marketSegmentIdFilters);
176
177 listener.subscribe(*handler);
178
179 handler->start();
180
181 while(!Samples::SignalHelper::interruptDetected() || listener.done())
183
184 handler->stop();
185
187
188 listener.saveLatencies("eobi-results.csv");
189 listener.processLatencies();
190}
The Feed Engine machinery.
Eurex Reference Data Interface Handler.
The given class implements feed engine concept using pool of working threads and standard socket API.
static void affinity(const ThreadAffinity &)
Sets the processor affinity mask for the current thread.
EobiDescriptor::Collection EobiDescriptors
EmdiDescriptor::Collection EmdiDescriptors
IInterfaceDescriptorProvider::MarketSegments MarketSegments
bool process(FeedEngine &engine)
@ Fatal
Fatal error, cannot continue.
@ TraceToFile
Trace to the log file.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19#include <iostream>
20#include <sstream>
21
22#include "RdiListener.h"
23
24
26
27
28
29
30
31RdiListener::RdiListener()
32 : referenceDataReady_(false)
33{
34}
35
36RdiListener::~RdiListener() = default;
37
38void RdiListener::onError(
ErrorCode::Enum code,
const std::string& description)
39{
40 std::clog <<
"Error occurred, errorCode = " <<
enumToString(code) <<
". Description: '" << description <<
"'" << std::endl;
41}
42
43void RdiListener::onWarning(const std::string& description)
44{
45 std::stringstream ss;
46 std::clog << "Warning occurred. Description: '" << description << "'" << std::endl;
47}
48
49void RdiListener::onSnapshotCycleStart()
50{
51}
52
54{
55}
56
58{
59}
60
62{
63}
64
66{
67}
68
70{
71}
72
74{
75}
76
77void RdiListener::onSnapshotCycleEnd()
78{
79 referenceDataReady_ = true;
80}
81
82bool RdiListener::referenceDataReceived() const
83{
84 return referenceDataReady_;
85}
std::string enumToString(HandlerState::Enum)
Returns string representation of HandlerState value.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19#include <map>
20#include <iostream>
21#include <sstream>
22
24
25#include "EobiListener.h"
26
29
30
31
32
33
34EobiListener::EobiListener(size_t count)
35 : ListenerBase(count)
36{
37}
38
39EobiListener::~EobiListener() = default;
40
42{
49}
50
51void EobiListener::onError(
ErrorCode::Enum code,
const std::string& description)
52{
53 std::clog <<
"Error occurred, errorCode = " <<
enumToString(code) <<
". Description: '" << description <<
"'" << std::endl;
54}
55
56void EobiListener::onWarning(const std::string& description)
57{
58 std::clog << "Warning occurred. Description: '" << description << "'" << std::endl;
59}
60
61void EobiListener::onOrderAdd(
const OrderAdd& orderAdd)
62{
66}
67
68void EobiListener::onOrderModify(
const OrderModify& orderModify)
69{
73}
74
76{
80}
81
82void EobiListener::onOrderDelete(
const OrderDelete& orderDelete)
83{
87}
88
89void EobiListener::onOrderMassDelete(
const OrderMassDelete& orderMassDelete)
90{
94}
95
97{
101}
102
104{
108}
109
111{
115}
116
118{
122}
123
124void EobiListener::onTopOfBook(
const TopOfBook& topOfBook)
125{
129}
130
131void EobiListener::onExecutionSummary(
const ExecutionSummary& executionSummary)
132{
136}
137
139{
143}
144
146{
150}
151
152void EobiListener::onTradeReport(
const TradeReport& tradeReport)
153{
157}
158
159void EobiListener::onTradeReversal(
const TradeReversal& tradeReversal)
160{
164}
165
167{
171}
172
174{
178}
179
180void EobiListener::onMassInstrumentStateChange(
182{
186}
187
189{
193}
194
196{
200}
201
203{
207}
208
210{
214}
EobiHandler & registerWarningListener(WarningListener *listener)
EobiHandler & registerStateChangeListener(StateChangeListener *listener)
EobiHandler & registerErrorListener(ErrorListener *listener)
EobiHandler & registerReferenceDataListener(ReferenceDataListener *listener)
EobiHandler & registerTradeDataListener(TradeDataListener *listener)
EobiHandler & registerOrderDataListener(OrderDataListener *listener)
const DataSource & dataSource() const
Returns data source.
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19#include <map>
20#include <climits>
21#include <iostream>
22#include <fstream>
23#include <cmath>
24#include <algorithm>
25
26#include "ListenerBase.h"
27
28
30
31
32
33
34
35ListenerBase::ListenerBase(size_t count)
36 : maxCount_(count)
37{
38 latencies_.reserve(maxCount_);
39}
40
41ListenerBase::~ListenerBase() = default;
42
44{
45 if ONIXS_EUREX_EMDI_UNLIKELY(latencies_.size() >= maxCount_)
46 return;
47
49 latencies_.push_back(latency.
ticks());
50}
51
52ListenerBase::Latency ListenerBase::calculateAdjustment()
53{
54 const int iterations = 10000;
55 Latencies latencies;
56
57 latencies.reserve(iterations);
58
59 for(int i = 0; i < iterations; ++i)
60 {
62 latencies.push_back(latency.
ticks());
63 }
64
65 std::sort(latencies.begin(), latencies.end());
66
67 return(latencies[latencies.size() / 2]);
68}
69
70void ListenerBase::processLatencies()
71{
72 size_t count = latencies_.size();
73
74 if(count == 0)
75 {
76 std::clog << "Nothing to process(latencies list is empty)." << std::endl;
77 return;
78 }
79
80 std::sort(latencies_.begin(), latencies_.end());
81
82 const Latency adjustment = calculateAdjustment();
83
84 const long double medianLatency =
85 static_cast<double>((latencies_[latencies_.size() / 2]) - adjustment);
86
87 std::clog << "Results:" << std::endl;
88 std::clog << "Latency(nanoseconds):" << std::endl;
89 std::clog << "Count: " << count << std::endl;
90 std::clog << "Median: " << medianLatency << std::endl;
91}
92
93void ListenerBase::saveLatencies(const std::string& filename)
94{
95 std::ofstream csvFile;
96
97 csvFile.open(filename.c_str());
98
99 for(auto && latency : latencies_)
100 csvFile << latency << std::endl;
101
102 csvFile.close();
103}
static Timestamp utcNow()
Returns current UTC time.