websocket_server.cpp 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708
  1. #include "websocket_server.hpp"
  2. #include "common/uuid.hpp"
  3. #include "logging/logger.hpp"
  4. #include <cstring>
  5. namespace smartbotic::webserver {
  6. // Per-session data
  7. struct PerSessionData {
  8. WebSocketClient* client;
  9. };
  10. // Forward declarations for LWS callbacks
  11. static int ws_callback(struct lws* wsi, enum lws_callback_reasons reason,
  12. void* user, void* in, size_t len);
  13. // Global pointer for callback access (set per-context in user data)
  14. static WebSocketServer* g_ws_server = nullptr;
  15. // Protocol definition
  16. static struct lws_protocols protocols[] = {
  17. {
  18. "smartbotic",
  19. ws_callback,
  20. sizeof(PerSessionData),
  21. 1024 * 64, // rx buffer size
  22. 0, nullptr, 0
  23. },
  24. { nullptr, nullptr, 0, 0, 0, nullptr, 0 }
  25. };
  26. static int ws_callback(struct lws* wsi, enum lws_callback_reasons reason,
  27. void* user, void* in, size_t len) {
  28. auto* pss = static_cast<PerSessionData*>(user);
  29. switch (reason) {
  30. case LWS_CALLBACK_ESTABLISHED:
  31. if (g_ws_server) {
  32. g_ws_server->onConnect(wsi);
  33. }
  34. break;
  35. case LWS_CALLBACK_CLOSED:
  36. if (g_ws_server) {
  37. g_ws_server->onDisconnect(wsi);
  38. }
  39. break;
  40. case LWS_CALLBACK_RECEIVE:
  41. if (g_ws_server && in && len > 0) {
  42. return g_ws_server->onReceive(wsi, static_cast<const char*>(in), len);
  43. }
  44. break;
  45. case LWS_CALLBACK_SERVER_WRITEABLE:
  46. if (g_ws_server) {
  47. return g_ws_server->onWritable(wsi);
  48. }
  49. break;
  50. default:
  51. break;
  52. }
  53. return 0;
  54. }
  55. WebSocketServer::WebSocketServer(const WebSocketServerConfig& config, auth::JwtUtils& jwt)
  56. : config_(config), jwt_(jwt) {
  57. g_ws_server = this;
  58. }
  59. WebSocketServer::~WebSocketServer() {
  60. stop();
  61. g_ws_server = nullptr;
  62. }
  63. void WebSocketServer::start() {
  64. if (running_) {
  65. return;
  66. }
  67. struct lws_context_creation_info info{};
  68. memset(&info, 0, sizeof(info));
  69. info.port = config_.port;
  70. info.protocols = protocols;
  71. info.gid = -1;
  72. info.uid = -1;
  73. info.options = LWS_SERVER_OPTION_HTTP_HEADERS_SECURITY_BEST_PRACTICES_ENFORCE;
  74. context_ = lws_create_context(&info);
  75. if (!context_) {
  76. LOG_ERROR("Failed to create WebSocket context");
  77. return;
  78. }
  79. running_ = true;
  80. service_thread_ = std::thread(&WebSocketServer::serviceLoop, this);
  81. LOG_INFO("WebSocket server started on port {}", config_.port);
  82. }
  83. void WebSocketServer::stop() {
  84. if (!running_) {
  85. return;
  86. }
  87. running_ = false;
  88. // Cancel the service loop immediately
  89. if (context_) {
  90. lws_cancel_service(context_);
  91. }
  92. if (service_thread_.joinable()) {
  93. service_thread_.join();
  94. }
  95. if (context_) {
  96. lws_context_destroy(context_);
  97. context_ = nullptr;
  98. }
  99. LOG_INFO("WebSocket server stopped");
  100. }
  101. void WebSocketServer::serviceLoop() {
  102. while (running_) {
  103. lws_service(context_, 50); // 50ms timeout
  104. // Process any pending writable requests (thread-safe wakeup mechanism)
  105. std::vector<struct lws*> to_signal;
  106. {
  107. std::lock_guard<std::mutex> lock(pending_mutex_);
  108. to_signal.assign(pending_writable_.begin(), pending_writable_.end());
  109. pending_writable_.clear();
  110. }
  111. for (auto* wsi : to_signal) {
  112. lws_callback_on_writable(wsi);
  113. }
  114. }
  115. }
  116. int WebSocketServer::onConnect(struct lws* wsi) {
  117. std::unique_lock lock(clients_mutex_);
  118. auto client = std::make_unique<WebSocketClient>();
  119. client->wsi = wsi;
  120. client->id = common::UUID::generate();
  121. client_id_map_[client->id] = wsi;
  122. clients_[wsi] = std::move(client);
  123. LOG_DEBUG("WebSocket client connected: {}", clients_[wsi]->id);
  124. return 0;
  125. }
  126. void WebSocketServer::onDisconnect(struct lws* wsi) {
  127. // Which workflows this client was watching has to be read before it is
  128. // erased, but the roster can only be republished after - and publishing
  129. // takes the same lock. So: collect, drop the lock, then publish.
  130. std::vector<std::string> was_watching;
  131. std::string gone_client_id;
  132. {
  133. std::unique_lock lock(clients_mutex_);
  134. auto it = clients_.find(wsi);
  135. if (it != clients_.end()) {
  136. LOG_DEBUG("WebSocket client disconnected: {}", it->second->id);
  137. gone_client_id = it->second->id;
  138. for (const auto& sub : it->second->subscriptions) {
  139. std::string workflow_id = presenceWorkflowId(sub);
  140. if (!workflow_id.empty()) was_watching.push_back(workflow_id);
  141. }
  142. client_id_map_.erase(it->second->id);
  143. clients_.erase(it);
  144. }
  145. }
  146. // A lock outlives its holder only if nothing notices they have gone, which
  147. // is exactly how an editor ends up with a node nobody can touch again.
  148. std::vector<std::string> unlocked;
  149. releaseLocksOf(gone_client_id, &unlocked);
  150. for (const auto& workflow_id : was_watching) {
  151. publishPresence(workflow_id);
  152. }
  153. for (const auto& workflow_id : unlocked) {
  154. publishLocks(workflow_id);
  155. }
  156. }
  157. int WebSocketServer::onReceive(struct lws* wsi, const char* data, size_t len) {
  158. WebSocketClient* client = nullptr;
  159. // Get client pointer under lock, then release before processing
  160. {
  161. std::shared_lock lock(clients_mutex_);
  162. auto it = clients_.find(wsi);
  163. if (it == clients_.end()) {
  164. return 0;
  165. }
  166. client = it->second.get();
  167. }
  168. // Process message without holding the lock (processMessage may call sendToClient
  169. // which needs unique_lock, so we must not hold shared_lock here)
  170. try {
  171. std::string str(data, len);
  172. auto message = nlohmann::json::parse(str);
  173. processMessage(*client, message);
  174. } catch (const std::exception& e) {
  175. LOG_WARN("Failed to parse WebSocket message: {}", e.what());
  176. }
  177. return 0;
  178. }
  179. int WebSocketServer::onWritable(struct lws* wsi) {
  180. std::unique_lock lock(clients_mutex_);
  181. auto it = clients_.find(wsi);
  182. if (it == clients_.end() || it->second->send_queue.empty()) {
  183. return 0;
  184. }
  185. auto& queue = it->second->send_queue;
  186. std::string& msg = queue.front();
  187. std::vector<unsigned char> buf(LWS_PRE + msg.size());
  188. memcpy(buf.data() + LWS_PRE, msg.data(), msg.size());
  189. int written = lws_write(wsi, buf.data() + LWS_PRE, msg.size(), LWS_WRITE_TEXT);
  190. if (written < 0) {
  191. LOG_ERROR("WebSocket write failed");
  192. return -1;
  193. }
  194. queue.erase(queue.begin());
  195. if (!queue.empty()) {
  196. lws_callback_on_writable(wsi);
  197. }
  198. return 0;
  199. }
  200. void WebSocketServer::processMessage(WebSocketClient& client, const nlohmann::json& message) {
  201. std::string type = message.value("type", "");
  202. if (type == "auth") {
  203. handleAuth(client, message);
  204. } else if (type == "subscribe") {
  205. handleSubscribe(client, message);
  206. } else if (type == "unsubscribe") {
  207. handleUnsubscribe(client, message);
  208. } else if (type == "lock") {
  209. handleLock(client, message);
  210. } else if (type == "unlock") {
  211. handleUnlock(client, message);
  212. } else if (type == "node_moved") {
  213. handleNodeMoved(client, message);
  214. } else if (type == "graph_edit") {
  215. handleGraphEdit(client, message);
  216. } else if (message_handler_) {
  217. message_handler_(client.id, message);
  218. }
  219. }
  220. void WebSocketServer::handleAuth(WebSocketClient& client, const nlohmann::json& message) {
  221. std::string token = message.value("token", "");
  222. // Remove "Bearer " prefix if present
  223. if (token.starts_with("Bearer ")) {
  224. token = token.substr(7);
  225. }
  226. auto result = jwt_.verifyToken(token);
  227. if (result.ok()) {
  228. client.authenticated = true;
  229. client.user_id = result.value().user_id;
  230. client.username = result.value().username;
  231. nlohmann::json response;
  232. response["type"] = "auth_success";
  233. response["userId"] = client.user_id;
  234. response["username"] = client.username;
  235. // This connection's own id. A lock is per connection, so an editor has
  236. // to be able to tell its own lock from one held by the same person in
  237. // another tab.
  238. response["clientId"] = client.id;
  239. sendToClient(client.id, response);
  240. LOG_DEBUG("WebSocket client authenticated: {}", client.id);
  241. } else {
  242. nlohmann::json response;
  243. response["type"] = "auth_error";
  244. response["error"] = result.error().message();
  245. sendToClient(client.id, response);
  246. }
  247. }
  248. void WebSocketServer::handleSubscribe(WebSocketClient& client, const nlohmann::json& message) {
  249. if (!client.authenticated) {
  250. nlohmann::json response;
  251. response["type"] = "error";
  252. response["error"] = "Not authenticated";
  253. sendToClient(client.id, response);
  254. return;
  255. }
  256. auto channels = message.value("channels", std::vector<std::string>{});
  257. for (const auto& channel : channels) {
  258. subscribe(client.id, channel);
  259. // Publishing after subscribing means the arriving client is itself on
  260. // the channel, so it gets the roster as its first message rather than
  261. // waiting for somebody else to come or go.
  262. std::string workflow_id = presenceWorkflowId(channel);
  263. if (!workflow_id.empty()) publishPresence(workflow_id);
  264. // The same for locks, which are otherwise only published when they
  265. // change: an editor arriving after somebody started working would see
  266. // an unlocked canvas and walk straight into a node already taken.
  267. std::string locks_workflow = channelWorkflowId(channel, ".locks");
  268. if (!locks_workflow.empty()) publishLocks(locks_workflow);
  269. }
  270. }
  271. void WebSocketServer::handleUnsubscribe(WebSocketClient& client, const nlohmann::json& message) {
  272. auto channels = message.value("channels", std::vector<std::string>{});
  273. for (const auto& channel : channels) {
  274. unsubscribe(client.id, channel);
  275. std::string workflow_id = presenceWorkflowId(channel);
  276. if (!workflow_id.empty()) publishPresence(workflow_id);
  277. }
  278. }
  279. std::string WebSocketServer::presenceChannel(const std::string& workflow_id) {
  280. return "workflows." + workflow_id + ".presence";
  281. }
  282. std::string WebSocketServer::channelWorkflowId(const std::string& channel,
  283. const std::string& suffix) {
  284. static const std::string prefix = "workflows.";
  285. if (!channel.starts_with(prefix) || !channel.ends_with(suffix)) return "";
  286. // A wildcard subscription is a listener, not somebody with the workflow
  287. // open, and counting it would put phantom names in the roster.
  288. std::string id = channel.substr(prefix.size(),
  289. channel.size() - prefix.size() - suffix.size());
  290. if (id.empty() || id.find('*') != std::string::npos) return "";
  291. return id;
  292. }
  293. std::string WebSocketServer::presenceWorkflowId(const std::string& channel) {
  294. return channelWorkflowId(channel, ".presence");
  295. }
  296. void WebSocketServer::publishPresence(const std::string& workflow_id) {
  297. const std::string channel = presenceChannel(workflow_id);
  298. nlohmann::json viewers = nlohmann::json::array();
  299. {
  300. std::shared_lock lock(clients_mutex_);
  301. // One person with the editor open in two tabs is one person, so the
  302. // roster is by user rather than by connection.
  303. std::unordered_map<std::string, size_t> seen;
  304. for (const auto& [wsi, client] : clients_) {
  305. (void)wsi;
  306. if (!client->authenticated) continue;
  307. if (!client->subscriptions.contains(channel)) continue;
  308. auto it = seen.find(client->user_id);
  309. if (it != seen.end()) {
  310. viewers[it->second]["connections"] =
  311. viewers[it->second]["connections"].get<int>() + 1;
  312. continue;
  313. }
  314. seen.emplace(client->user_id, viewers.size());
  315. viewers.push_back({
  316. {"userId", client->user_id},
  317. {"username", client->username},
  318. {"connections", 1},
  319. });
  320. }
  321. }
  322. broadcast(channel, {{"workflowId", workflow_id}, {"viewers", viewers}});
  323. }
  324. // ---------------------------------------------------------------------------
  325. // Node locks and live movement
  326. //
  327. // While somebody is dragging or configuring a node, everyone else is kept off
  328. // that one node - not off the workflow. Two people working on different parts
  329. // of the same flow is the normal case and should stay possible; two people
  330. // dragging the same node is the one that produces nonsense.
  331. // ---------------------------------------------------------------------------
  332. void WebSocketServer::handleLock(WebSocketClient& client, const nlohmann::json& message) {
  333. if (!client.authenticated) return;
  334. const std::string workflow_id = message.value("workflowId", "");
  335. const std::string node_id = message.value("nodeId", "");
  336. const std::string kind = message.value("kind", "editing");
  337. if (workflow_id.empty() || node_id.empty()) return;
  338. bool granted = false;
  339. std::string holder_username;
  340. {
  341. std::lock_guard<std::mutex> lock(locks_mutex_);
  342. auto& nodes = locks_[workflow_id];
  343. auto it = nodes.find(node_id);
  344. if (it == nodes.end() || it->second.client_id == client.id) {
  345. // Re-claiming your own lock is how the kind changes from dragging
  346. // to editing without a release in between.
  347. nodes[node_id] = NodeLock{client.id, client.user_id, client.username, kind};
  348. granted = true;
  349. } else {
  350. holder_username = it->second.username;
  351. }
  352. }
  353. if (!granted) {
  354. // Said out loud rather than ignored: a click that silently does nothing
  355. // reads as the editor being broken.
  356. nlohmann::json response;
  357. response["type"] = "lock_denied";
  358. response["workflowId"] = workflow_id;
  359. response["nodeId"] = node_id;
  360. response["heldBy"] = holder_username;
  361. sendToClient(client.id, response);
  362. return;
  363. }
  364. publishLocks(workflow_id);
  365. }
  366. void WebSocketServer::handleUnlock(WebSocketClient& client, const nlohmann::json& message) {
  367. const std::string workflow_id = message.value("workflowId", "");
  368. const std::string node_id = message.value("nodeId", "");
  369. if (workflow_id.empty() || node_id.empty()) return;
  370. bool changed = false;
  371. {
  372. std::lock_guard<std::mutex> lock(locks_mutex_);
  373. auto wf = locks_.find(workflow_id);
  374. if (wf != locks_.end()) {
  375. auto it = wf->second.find(node_id);
  376. // Only the holder can let go - otherwise a stale release from an
  377. // editor that has already moved on would free somebody else's node.
  378. if (it != wf->second.end() && it->second.client_id == client.id) {
  379. wf->second.erase(it);
  380. changed = true;
  381. }
  382. }
  383. }
  384. if (changed) publishLocks(workflow_id);
  385. }
  386. void WebSocketServer::handleNodeMoved(WebSocketClient& client, const nlohmann::json& message) {
  387. if (!client.authenticated) return;
  388. const std::string workflow_id = message.value("workflowId", "");
  389. const std::string node_id = message.value("nodeId", "");
  390. if (workflow_id.empty() || node_id.empty()) return;
  391. // Only the editor holding the node may say where it is. Without this,
  392. // anyone could move anybody's node by sending the message directly.
  393. {
  394. std::lock_guard<std::mutex> lock(locks_mutex_);
  395. auto wf = locks_.find(workflow_id);
  396. if (wf == locks_.end()) return;
  397. auto it = wf->second.find(node_id);
  398. if (it == wf->second.end() || it->second.client_id != client.id) return;
  399. }
  400. // Not stored: a position in flight is worth nothing once it has been
  401. // delivered, and the authoritative one is whatever gets saved.
  402. broadcast("workflows." + workflow_id + ".moves", {
  403. {"workflowId", workflow_id},
  404. {"nodeId", node_id},
  405. {"position", message.value("position", nlohmann::json::object())},
  406. {"userId", client.user_id},
  407. {"clientId", client.id},
  408. });
  409. }
  410. void WebSocketServer::handleGraphEdit(WebSocketClient& client, const nlohmann::json& message) {
  411. if (!client.authenticated) return;
  412. const std::string workflow_id = message.value("workflowId", "");
  413. if (workflow_id.empty()) return;
  414. // Having the workflow open is the check that matters here. Anyone who can
  415. // see it can already save whatever they like to it, so demanding a lock per
  416. // edit would buy nothing and would make an auto-layout - which moves every
  417. // node at once and holds no lock at all - impossible to share.
  418. const std::string channel = "workflows." + workflow_id + ".edits";
  419. {
  420. std::shared_lock lock(clients_mutex_);
  421. if (!client.subscriptions.contains(channel)) return;
  422. }
  423. nlohmann::json out = message.value("delta", nlohmann::json::object());
  424. out["workflowId"] = workflow_id;
  425. // Stamped here rather than trusted from the sender, so an edit cannot be
  426. // attributed to somebody else.
  427. out["userId"] = client.user_id;
  428. out["username"] = client.username;
  429. out["clientId"] = client.id;
  430. broadcast(channel, out);
  431. }
  432. void WebSocketServer::releaseLocksOf(const std::string& client_id,
  433. std::vector<std::string>* touched_workflows) {
  434. std::lock_guard<std::mutex> lock(locks_mutex_);
  435. for (auto& [workflow_id, nodes] : locks_) {
  436. bool changed = false;
  437. for (auto it = nodes.begin(); it != nodes.end();) {
  438. if (it->second.client_id == client_id) {
  439. it = nodes.erase(it);
  440. changed = true;
  441. } else {
  442. ++it;
  443. }
  444. }
  445. if (changed && touched_workflows) touched_workflows->push_back(workflow_id);
  446. }
  447. }
  448. void WebSocketServer::publishLocks(const std::string& workflow_id) {
  449. nlohmann::json held = nlohmann::json::array();
  450. {
  451. std::lock_guard<std::mutex> lock(locks_mutex_);
  452. auto wf = locks_.find(workflow_id);
  453. if (wf != locks_.end()) {
  454. for (const auto& [node_id, holder] : wf->second) {
  455. held.push_back({
  456. {"nodeId", node_id},
  457. {"userId", holder.user_id},
  458. {"username", holder.username},
  459. {"clientId", holder.client_id},
  460. {"kind", holder.kind},
  461. });
  462. }
  463. }
  464. }
  465. broadcast("workflows." + workflow_id + ".locks",
  466. {{"workflowId", workflow_id}, {"locks", held}});
  467. }
  468. void WebSocketServer::broadcast(const std::string& channel, const nlohmann::json& data) {
  469. nlohmann::json message;
  470. message["channel"] = channel;
  471. message["data"] = data;
  472. std::string msg_str = message.dump();
  473. std::vector<struct lws*> clients_to_signal;
  474. {
  475. std::unique_lock lock(clients_mutex_); // Must be unique_lock - we modify send_queue
  476. for (auto& [wsi, client] : clients_) {
  477. if (!client->authenticated) {
  478. continue;
  479. }
  480. for (const auto& sub : client->subscriptions) {
  481. if (matchesChannel(sub, channel)) {
  482. client->send_queue.push_back(msg_str);
  483. clients_to_signal.push_back(wsi);
  484. break;
  485. }
  486. }
  487. }
  488. }
  489. // Signal pending writes outside of clients_mutex_ to avoid deadlock
  490. if (!clients_to_signal.empty()) {
  491. std::lock_guard<std::mutex> lock(pending_mutex_);
  492. for (auto* wsi : clients_to_signal) {
  493. pending_writable_.insert(wsi);
  494. }
  495. // Wake up the service loop
  496. lws_cancel_service(context_);
  497. }
  498. LOG_DEBUG("Broadcast to channel {}: {} clients matched", channel, clients_to_signal.size());
  499. }
  500. void WebSocketServer::sendToClient(const std::string& client_id, const nlohmann::json& message) {
  501. struct lws* wsi_to_signal = nullptr;
  502. {
  503. std::unique_lock lock(clients_mutex_);
  504. auto it = client_id_map_.find(client_id);
  505. if (it == client_id_map_.end()) {
  506. return;
  507. }
  508. auto client_it = clients_.find(it->second);
  509. if (client_it == clients_.end()) {
  510. return;
  511. }
  512. client_it->second->send_queue.push_back(message.dump());
  513. wsi_to_signal = it->second;
  514. }
  515. // Signal pending write outside of clients_mutex_
  516. if (wsi_to_signal) {
  517. std::lock_guard<std::mutex> lock(pending_mutex_);
  518. pending_writable_.insert(wsi_to_signal);
  519. lws_cancel_service(context_);
  520. }
  521. }
  522. void WebSocketServer::sendToUser(const std::string& user_id, const nlohmann::json& message) {
  523. std::vector<struct lws*> clients_to_signal;
  524. std::string msg_str = message.dump();
  525. {
  526. std::unique_lock lock(clients_mutex_);
  527. for (auto& [wsi, client] : clients_) {
  528. if (client->user_id == user_id) {
  529. client->send_queue.push_back(msg_str);
  530. clients_to_signal.push_back(wsi);
  531. }
  532. }
  533. }
  534. // Signal pending writes outside of clients_mutex_
  535. if (!clients_to_signal.empty()) {
  536. std::lock_guard<std::mutex> lock(pending_mutex_);
  537. for (auto* wsi : clients_to_signal) {
  538. pending_writable_.insert(wsi);
  539. }
  540. lws_cancel_service(context_);
  541. }
  542. }
  543. void WebSocketServer::subscribe(const std::string& client_id, const std::string& channel) {
  544. std::unique_lock lock(clients_mutex_);
  545. auto it = client_id_map_.find(client_id);
  546. if (it == client_id_map_.end()) {
  547. return;
  548. }
  549. auto client_it = clients_.find(it->second);
  550. if (client_it == clients_.end()) {
  551. return;
  552. }
  553. client_it->second->subscriptions.insert(channel);
  554. LOG_DEBUG("Client {} subscribed to {}", client_id, channel);
  555. }
  556. void WebSocketServer::unsubscribe(const std::string& client_id, const std::string& channel) {
  557. std::unique_lock lock(clients_mutex_);
  558. auto it = client_id_map_.find(client_id);
  559. if (it == client_id_map_.end()) {
  560. return;
  561. }
  562. auto client_it = clients_.find(it->second);
  563. if (client_it == clients_.end()) {
  564. return;
  565. }
  566. client_it->second->subscriptions.erase(channel);
  567. }
  568. bool WebSocketServer::matchesChannel(const std::string& subscription, const std::string& channel) {
  569. // Support wildcard matching: "executions.*" matches "executions.123.started"
  570. if (subscription == channel) {
  571. return true;
  572. }
  573. if (subscription.ends_with(".*")) {
  574. std::string prefix = subscription.substr(0, subscription.length() - 1);
  575. return channel.starts_with(prefix);
  576. }
  577. return false;
  578. }
  579. void WebSocketServer::setMessageHandler(WebSocketMessageHandler handler) {
  580. message_handler_ = std::move(handler);
  581. }
  582. size_t WebSocketServer::getConnectionCount() const {
  583. std::shared_lock lock(clients_mutex_);
  584. return clients_.size();
  585. }
  586. } // namespace smartbotic::webserver