| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128 |
- #pragma once
- #include <string>
- #include <memory>
- #include <thread>
- #include <atomic>
- #include <unordered_map>
- #include <unordered_set>
- #include <shared_mutex>
- #include <functional>
- #include <libwebsockets.h>
- #include <nlohmann/json.hpp>
- #include "auth/jwt_utils.hpp"
- namespace smartbotic::webserver {
- // WebSocket client connection
- struct WebSocketClient {
- struct lws* wsi = nullptr;
- std::string id;
- std::string user_id;
- std::string username;
- bool authenticated = false;
- std::unordered_set<std::string> subscriptions;
- std::vector<std::string> send_queue;
- };
- // WebSocket message handler
- using WebSocketMessageHandler = std::function<void(
- const std::string& client_id,
- const nlohmann::json& message
- )>;
- struct WebSocketServerConfig {
- int port = 8080; // Same port as HTTP (uses vhost)
- std::string path = "/ws";
- int max_payload_size = 1024 * 1024; // 1MB
- };
- class WebSocketServer {
- public:
- WebSocketServer(const WebSocketServerConfig& config, auth::JwtUtils& jwt);
- ~WebSocketServer();
- // Lifecycle
- void start();
- void stop();
- // Broadcasting
- void broadcast(const std::string& channel, const nlohmann::json& data);
- void sendToClient(const std::string& client_id, const nlohmann::json& message);
- void sendToUser(const std::string& user_id, const nlohmann::json& message);
- // Subscription management
- void subscribe(const std::string& client_id, const std::string& channel);
- void unsubscribe(const std::string& client_id, const std::string& channel);
- // Who currently has a workflow open. Derived from the subscriptions rather
- // than stored separately, so a client that drops off never leaves a ghost
- // behind - there is no second list that can disagree with the connections.
- void publishPresence(const std::string& workflow_id);
- static std::string presenceWorkflowId(const std::string& channel);
- static std::string channelWorkflowId(const std::string& channel, const std::string& suffix);
- static std::string presenceChannel(const std::string& workflow_id);
- // Which nodes somebody is currently working on. Held here rather than in
- // the browsers for the same reason as presence: a lock whose holder has
- // gone must go with them, and the connection closing is the only thing that
- // reliably knows.
- void publishLocks(const std::string& workflow_id);
- void releaseLocksOf(const std::string& client_id,
- std::vector<std::string>* touched_workflows);
- // Message handler
- void setMessageHandler(WebSocketMessageHandler handler);
- // Stats
- size_t getConnectionCount() const;
- // LWS callbacks (public for protocol access)
- int onConnect(struct lws* wsi);
- void onDisconnect(struct lws* wsi);
- int onReceive(struct lws* wsi, const char* data, size_t len);
- int onWritable(struct lws* wsi);
- private:
- void serviceLoop();
- void processMessage(WebSocketClient& client, const nlohmann::json& message);
- void handleAuth(WebSocketClient& client, const nlohmann::json& message);
- void handleSubscribe(WebSocketClient& client, const nlohmann::json& message);
- void handleUnsubscribe(WebSocketClient& client, const nlohmann::json& message);
- void handleLock(WebSocketClient& client, const nlohmann::json& message);
- void handleUnlock(WebSocketClient& client, const nlohmann::json& message);
- void handleNodeMoved(WebSocketClient& client, const nlohmann::json& message);
- void handleGraphEdit(WebSocketClient& client, const nlohmann::json& message);
- bool matchesChannel(const std::string& subscription, const std::string& channel);
- WebSocketServerConfig config_;
- auth::JwtUtils& jwt_;
- struct lws_context* context_ = nullptr;
- std::thread service_thread_;
- std::atomic<bool> running_{false};
- // workflow id -> node id -> who is holding it. A lock is per connection,
- // not per user: the same person in two tabs is two editors, and the tab
- // that did not claim the node must not be allowed to drag it either.
- struct NodeLock {
- std::string client_id;
- std::string user_id;
- std::string username;
- std::string kind; // what they are doing: "dragging", "editing"
- };
- std::unordered_map<std::string, std::unordered_map<std::string, NodeLock>> locks_;
- mutable std::mutex locks_mutex_;
- std::unordered_map<struct lws*, std::unique_ptr<WebSocketClient>> clients_;
- std::unordered_map<std::string, struct lws*> client_id_map_;
- mutable std::shared_mutex clients_mutex_;
- // Track clients needing writable callback (for cross-thread signaling)
- std::unordered_set<struct lws*> pending_writable_;
- std::mutex pending_mutex_;
- WebSocketMessageHandler message_handler_;
- };
- } // namespace smartbotic::webserver
|