| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103 |
- #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 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<bool> running_{false};
- 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
|