#pragma once #include #include #include #include #include #include #include #include #include #include #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 subscriptions; std::vector send_queue; }; // WebSocket message handler using WebSocketMessageHandler = std::function; 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* 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 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> locks_; mutable std::mutex locks_mutex_; std::unordered_map> clients_; std::unordered_map client_id_map_; mutable std::shared_mutex clients_mutex_; // Track clients needing writable callback (for cross-thread signaling) std::unordered_set pending_writable_; std::mutex pending_mutex_; WebSocketMessageHandler message_handler_; }; } // namespace smartbotic::webserver