NexusForce 1.0.0
A rigorously engineered full-stack C++ backend library.
载入中...
搜索中...
未找到
websocket.hpp
浏览该文件的文档.
1#ifndef NEFORCE_NETWORK_HTTP_WEBSOCKET_HPP__
2#define NEFORCE_NETWORK_HTTP_WEBSOCKET_HPP__
3
10
22NEFORCE_BEGIN_NAMESPACE__
23NEFORCE_BEGIN_HTTP__
24
29
47
72
81 TEXT = 0x1,
82 BINARY = 0x2,
83 CLOSE = 0x8,
84 PING = 0x9,
85 PONG = 0xA
86};
87
88#pragma pack(push, 1)
106#pragma pack(pop)
107
109
110
135class NEFORCE_API websocket_server {
136public:
139
140private:
142 vector<session_ptr> sessions_;
143 mutable mutex sessions_mutex_;
144 io_context* ctx_{nullptr};
145
146public:
147 websocket_server() = default;
149
150 websocket_server(const websocket_server&) = delete;
151 websocket_server& operator=(const websocket_server&) = delete;
152
153 websocket_server(websocket_server&&) noexcept = delete;
154 websocket_server& operator=(websocket_server&&) noexcept = delete;
155
163 void set_io_context(io_context& ctx) noexcept { ctx_ = &ctx; }
164
170 void route(const string& path, session_handler handler) { route_handlers_[path] = _NEFORCE move(handler); }
171
179
184 void remove_session(const session_ptr& session);
185 void stop();
186
193
198 size_t session_count() const noexcept {
199 lock<mutex> lk(sessions_mutex_);
200 return sessions_.size();
201 }
202
203 // TODO: STOMP sub-protocol support — STOMP 1.2 over WebSocket with destination-based routing, subscription management, ACK modes
204 // TODO: SockJS fallback — HTTP long-polling / XHR streaming fallback for browsers without WebSocket support
205};
206
229class NEFORCE_API websocket_session : public enable_shared_from_this<websocket_session> {
230public:
231 using message_handler = function<void(const string&, websocket_opcode)>;
232 using close_handler = function<void(websocket_status, const string&)>;
233 using error_handler = function<void(const exception&)>;
234
235private:
236 unique_ptr<tcp_socket> socket_;
237 websocket_server* server_;
238
239 atomic<bool> running_{false};
240 atomic_flag closed_once_;
241
242 thread read_thread_;
243 thread write_thread_;
244 thread heartbeat_thread_;
245
246 io_context* ctx_{nullptr};
247 bool event_driven_ = false;
248 size_t heartbeat_timer_id_ = 0;
249 byte_vector read_buffer_;
250 bool write_registered_ = false;
251
252 mutex write_mutex_;
253 condition_variable write_cv_;
254 queue<byte_vector> write_queue_;
255 queue<byte_vector> ctrl_queue_;
256
257 string fragment_buffer_;
258 websocket_opcode fragment_opcode_ = websocket_opcode::TEXT;
259 bool in_fragment_ = false;
260
261 atomic<bool> ping_pending_{false};
262 atomic<int64_t> last_pong_ms_{0};
263
264 websocket_deflate_config deflate_config_;
265 unique_ptr<websocket_deflate> deflate_compressor_;
266 unique_ptr<websocket_deflate> deflate_decompressor_;
267 byte_vector deflate_fragment_buffer_;
268
269 message_handler on_message_;
270 close_handler on_close_;
271 error_handler on_error_;
272
273 bool queue_frame(byte_vector frame, bool is_control = false);
274 void write_loop();
275
276 void read_loop();
277 bool read_frame();
278
279 bool dispatch(const websocket_frame_header& hdr, websocket_opcode opcode, string payload);
280 void deliver_message(const string& data, websocket_opcode opcode);
281
282 void send_close_frame(websocket_status status, const string& reason);
283 void handle_close_frame(string payload);
284
285 void heartbeat_loop();
286 void do_stop(websocket_status status, const string& reason, bool notify_server = true);
287
288 void start_event_driven();
289
290 void on_readable(int fd, uint32_t events, error_code ec);
291 void on_writable(int fd, uint32_t events, error_code ec);
292 void on_heartbeat_timer();
293
294 void flush_event_writes();
295 void try_parse_frames();
296
297public:
304
309
310 websocket_session(const websocket_session&) = delete;
311 websocket_session& operator=(const websocket_session&) = delete;
312
318 void start();
319
325 void close(websocket_status status = websocket_status::NORMAL_CLOSURE, const string& reason = "");
326
330 void stop();
331
338 bool send(const string& data, websocket_opcode opcode = websocket_opcode::TEXT);
339
345 bool send_binary(const string& data) { return send(data, websocket_opcode::BINARY); }
346
351 bool is_open() const noexcept { return running_ && socket_->is_open(); }
352
357 void set_message_handler(message_handler handler) { on_message_ = _NEFORCE move(handler); }
358
363 void set_close_handler(close_handler handler) { on_close_ = _NEFORCE move(handler); }
364
369 void set_error_handler(error_handler handler) { on_error_ = _NEFORCE move(handler); }
370
379
384 const websocket_deflate_config& deflate_config() const noexcept { return deflate_config_; }
385
390 bool has_deflate_config() const noexcept { return deflate_config_.active; }
391
399 void set_io_context(io_context& ctx) noexcept {
400 ctx_ = &ctx;
401 event_driven_ = true;
402 }
403
408 tcp_socket& socket() noexcept { return *socket_; }
409
414 const tcp_socket& socket() const noexcept { return *socket_; }
415
420 ssl_socket* ssl_socket_ptr() noexcept { return dynamic_cast<ssl_socket*>(socket_.get()); }
421
426 const ssl_socket* ssl_socket_ptr() const noexcept { return dynamic_cast<const ssl_socket*>(socket_.get()); }
427};
428 // WebSocket
430 // HTTP
432
433NEFORCE_END_HTTP__
434NEFORCE_END_NAMESPACE__
435#endif // NEFORCE_NETWORK_HTTP_WEBSOCKET_HPP__
原子类型完整实现
函数包装器主模板声明
function< void(session_ptr)> session_handler
会话处理器类型
void set_io_context(io_context &ctx) noexcept
设置异步 I/O 执行上下文
bool handle_upgrade(const http_request &request, unique_ptr< tcp_socket > sock)
处理WebSocket升级请求
size_t session_count() const noexcept
获取活动会话数量
shared_ptr< websocket_session > session_ptr
会话智能指针类型
void broadcast(const string &data, websocket_opcode opcode=websocket_opcode::TEXT)
向所有会话广播消息
void remove_session(const session_ptr &session)
移除会话
void route(const string &path, session_handler handler)
注册WebSocket路由
bool send_binary(const string &data)
发送二进制消息
const ssl_socket * ssl_socket_ptr() const noexcept
获取SSL socket常量指针
function< void(const string &, websocket_opcode)> message_handler
消息处理器类型
tcp_socket & socket() noexcept
获取底层socket引用
void set_deflate_config(const websocket_deflate_config &cfg)
设置permessage-deflate配置
void set_error_handler(error_handler handler)
设置错误处理器
function< void(websocket_status, const string &)> close_handler
关闭处理器类型
ssl_socket * ssl_socket_ptr() noexcept
获取SSL socket指针
const tcp_socket & socket() const noexcept
获取底层socket常量引用
bool is_open() const noexcept
检查连接是否开启
bool send(const string &data, websocket_opcode opcode=websocket_opcode::TEXT)
发送文本/二进制消息
bool has_deflate_config() const noexcept
检查是否已启用deflate压缩
function< void(const exception &)> error_handler
错误处理器类型
websocket_session(unique_ptr< tcp_socket > sock, websocket_server *server=nullptr)
构造函数
const websocket_deflate_config & deflate_config() const noexcept
获取permessage-deflate协商配置
void set_io_context(io_context &ctx) noexcept
设置 io_context(启用事件驱动模式)
void set_message_handler(message_handler handler)
设置消息处理器
void close(websocket_status status=websocket_status::NORMAL_CLOSURE, const string &reason="")
关闭连接
void set_close_handler(close_handler handler)
设置关闭处理器
统一异步操作核心
锁管理器模板
非递归互斥锁
文件路径类
队列容器适配器
共享智能指针类模板
SSL/TLS安全Socket类
独占智能指针
动态大小数组容器
条件变量行为
通用函数包装器
vector< byte_t > byte_vector
字节向量类型别名
unsigned char byte_t
字节类型,定义为无符号字符
unsigned char uint8_t
8位无符号整数类型
unsigned short uint16_t
16位无符号整数类型
@ BINARY
二进制 'b'
http_server_request http_request
HTTP请求类型别名
constexpr Iterator2 move(Iterator1 first, Iterator1 last, Iterator2 result) noexcept(noexcept(inner::__move_aux(first, last, result)))
移动范围元素
constexpr decltype(auto) data(Container &cont) noexcept(noexcept(cont.data()))
获取容器的底层数据指针
websocket_opcode
WebSocket帧操作码
websocket_status
WebSocket关闭状态码
@ INVALID_FRAME_PAYLOAD_DATA
无效帧负载
@ UNSUPPORTED_DATA
不支持的数据类型
HTTP服务器消息结构
统一异步操作核心
队列容器适配器
SSL/TLS安全Socket实现
通用原子类型模板
byte_t fin
是否最后一帧 (bit 7)
byte_t payload_len
负载长度 (bits 0-6)
线程管理类
无序映射容器
弱智能指针实现
WebSocket permessage-deflate 扩展