#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 presenceChannel(const std::string& workflow_id); // 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); 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}; 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