2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
44#include <express/Response.h>
45#include <web/http/server/SocketContext.h>
51#ifndef DOXYGEN_SHOULD_SKIP_THIS
56#include <log/Logger.h>
57#include <nlohmann/json.hpp>
71 return sseDistributor;
80 response->getSocketContext()->onDisconnected([
this, eventReceiverId]() {
82 return eventReceiver
.getId() == eventReceiverId;
92 const std::string& data,
93 const std::string& event,
94 const std::string& id) {
95 if (response->isConnected()) {
97 response->sendFragment(
"event:" + event);
100 response->sendFragment(
"id:" + id);
102 response->sendFragment(
"data:" + data);
103 response->sendFragment();
109 const std::string& event,
110 const std::string& id) {
115 VLOG(0) <<
"Server sent event: " << event <<
"\n" << data;
150 {{
"at",
timePointToString(std::chrono::system_clock::now()
)}, {
"name", bridgeName}}
, "bridge_disabled", std::to_string(
id++)
);
155 {{
"at",
timePointToString(std::chrono::system_clock::now()
)}, {
"name", bridgeName}}
, "bridge_starting", std::to_string(
id++)
);
160 {{
"at",
timePointToString(std::chrono::system_clock::now()
)}, {
"name", bridgeName}}
, "bridge_started", std::to_string(
id++)
);
165 {{
"at",
timePointToString(std::chrono::system_clock::now()
)}, {
"name", bridgeName}}
, "bridge_stopping", std::to_string(
id++)
);
170 {{
"at",
timePointToString(std::chrono::system_clock::now()
)}, {
"name", bridgeName}}
, "bridge_stopped", std::to_string(
id++)
);
176 std::to_string(
id++)
);
182 std::to_string(
id++)
);
188 std::to_string(
id++)
);
193 "broker_disconnecting",
194 std::to_string(
id++)
);
199 "broker_disconnected",
200 std::to_string(
id++)
);
208 std::time_t time = std::chrono::system_clock::to_time_t(timePoint);
209 const std::tm* tm_ptr = std::gmtime(&time);
212 std::string onlineSince =
"Formatting error";
215 if (std::strftime(buffer,
sizeof(buffer),
"%Y-%m-%d %H:%M:%S", tm_ptr)) {
216 onlineSince = std::string(buffer) +
" UTC";
223 const std::chrono::time_point<std::chrono::system_clock>& later) {
224 using seconds_duration_type = std::chrono::duration<std::chrono::seconds::rep>::rep;
226 seconds_duration_type totalSeconds = std::chrono::duration_cast<std::chrono::seconds>(later - bevore).count();
229 seconds_duration_type days = totalSeconds / 86400;
230 seconds_duration_type remainder = totalSeconds % 86400;
231 seconds_duration_type hours = remainder / 3600;
232 remainder = remainder % 3600;
233 seconds_duration_type minutes = remainder / 60;
234 seconds_duration_type seconds = remainder % 60;
237 std::ostringstream oss;
239 oss << days <<
" day" << (days == 1 ?
"" :
"s") <<
", ";
241 oss << std::setw(2) << std::setfill(
'0') << hours <<
":" << std::setw(2) << std::setfill(
'0') << minutes <<
":" << std::setw(2)
242 << std::setfill(
'0') << seconds;
252 response->sendFragment(
":keep-alive");
253 response->sendFragment();
std::uint64_t getId() const
core::timer::Timer heartbeatTimer
std::weak_ptr< express::Response > response
std::shared_ptr< express::Response > getResponse() const
EventReceiver(std::uint64_t id, const std::shared_ptr< express::Response > &response)
const std::string & getData() const
const std::string & getEvent() const
const std::string & getId() const
Event(const std::string &data, const std::string &event, const std::string &id)
static void sendJsonEvent1(const std::shared_ptr< express::Response > &response, const nlohmann::json &json, const std::string &event="", const std::string &id="")
void sendEvent(const std::string &data, const std::string &event="", const std::string &id="")
std::uint64_t nextEventReceiverId
std::string bridgesStartedAt() const
static std::string durationToString(const std::chrono::time_point< std::chrono::system_clock > &bevore, const std::chrono::time_point< std::chrono::system_clock > &later=std::chrono::system_clock::now())
void bridgeStopped(const std::string &bridgeName)
void addEventReceiver(const std::shared_ptr< express::Response > &response, const std::string &lastEventId)
void brokerConnected(const std::string &bridgeName, const std::string &instanceName)
std::chrono::time_point< std::chrono::system_clock > bridgesStartTimePoint
void brokerConnecting(const std::string &bridgeName, const std::string &instanceName)
void brokerDisabled(const std::string &bridgeName, const std::string &instanceName)
void bridgeDisabled(const std::string &bridgeName)
std::list< Event > replayEvents
static void sendEvent(const std::shared_ptr< express::Response > &response, const std::string &data, const std::string &event, const std::string &id)
void bridgeStarting(const std::string &bridgeName)
void bridgeStarted(const std::string &bridgeName)
static SSEDistributor & instance()
std::list< EventReceiver > eventReceiverList
std::chrono::time_point< std::chrono::system_clock > onlineSinceTimePoint
void brokerDisconnected(const std::string &bridgeName, const std::string &instanceName)
void brokerDisconnecting(const std::string &bridgeName, const std::string &instanceName)
static std::string timePointToString(const std::chrono::time_point< std::chrono::system_clock > &timePoint)
void sendJsonEvent(const nlohmann::json &json, const std::string &event="", const std::string &id="")
void bridgeStopping(const std::string &bridgeName)