Crow  1.1
A C++ microframework for the web
 
Loading...
Searching...
No Matches
websocket.h
1#pragma once
2#include <array>
3#include <memory>
4#include <optional>
5#include <string>
6#include <thread>
7#include "crow/http_response.h"
8#include "crow/logging.h"
9#include "crow/socket_adaptors.h"
10#include "crow/http_request.h"
11#include "crow/tcp_socket_options.h"
12#include "crow/TinySHA1.hpp"
13#include "crow/utility.h"
14
15namespace crow // NOTE: Already documented in "crow/app.h"
16{
17#ifdef CROW_USE_BOOST
18 namespace asio = boost::asio;
19 using error_code = boost::system::error_code;
20#else
21 using error_code = asio::error_code;
22#endif
23
24 /**
25 * \namespace crow::websocket
26 * \brief Namespace that includes the \ref Connection class
27 * and \ref connection struct. Useful for WebSockets connection.
28 *
29 * Used specially in crow/websocket.h, crow/app.h and crow/routing.h
30 */
31 namespace websocket
32 {
33 enum class WebSocketReadState
34 {
35 MiniHeader,
36 Len16,
37 Len64,
38 Mask,
39 Payload,
40 };
41
42 // Codes taken from https://www.rfc-editor.org/rfc/rfc6455#section-7.4.1
43 enum CloseStatusCode : uint16_t {
44 NormalClosure = 1000,
45 EndpointGoingAway = 1001,
46 ProtocolError = 1002,
47 UnacceptableData = 1003,
48 InconsistentData = 1007,
49 PolicyViolated = 1008,
50 MessageTooBig = 1009,
51 ExtensionsNotNegotiated = 1010,
52 UnexpectedCondition = 1011,
53
54 // Reserved for applications only, should not send/receive these to/from clients
55 NoStatusCodePresent = 1005,
56 ClosedAbnormally = 1006,
57 TLSHandshakeFailure = 1015,
58
59 StartStatusCodesForLibraries = 3000,
60 StartStatusCodesForPrivateUse = 4000,
61 // Status code should be between 1000 and 4999 inclusive
62 StartStatusCodes = NormalClosure,
63 EndStatusCodes = 4999,
64 };
65
66 /// A base class for websocket connection.
68 {
69 virtual void send_binary(std::string msg) = 0;
70 virtual void send_text(std::string msg) = 0;
71 virtual void send_ping(std::string msg) = 0;
72 virtual void send_pong(std::string msg) = 0;
73 virtual void close(std::string const& msg = "quit", uint16_t status_code = CloseStatusCode::NormalClosure) = 0;
74 virtual std::string get_remote_ip() = 0;
75 virtual uint16_t get_remote_port() = 0;
76 virtual std::string get_subprotocol() const = 0;
77 virtual ~connection() = default;
78
79 void userdata(void* u) { userdata_ = u; }
80 void* userdata() { return userdata_; }
81
82 private:
83 void* userdata_;
84 };
85
86 // Modified version of the illustration in RFC6455 Section-5.2
87 //
88 //
89 // 0 1 2 3 -byte
90 // 0 1 2 3 4 5 6 7 0 1 2 3 4 5 6 7 0 1 2 3 4 5 6 7 0 1 2 3 4 5 6 7 -bit
91 // +-+-+-+-+-------+-+-------------+-------------------------------+
92 // |F|R|R|R| opcode|M| Payload len | Extended payload length |
93 // |I|S|S|S| (4) |A| (7) | (16/64) |
94 // |N|V|V|V| |S| | (if payload len==126/127) |
95 // | |1|2|3| |K| | |
96 // +-+-+-+-+-------+-+-------------+ - - - - - - - - - - - - - - - +
97 // | Extended payload length continued, if payload len == 127 |
98 // + - - - - - - - - - - - - - - - +-------------------------------+
99 // | |Masking-key, if MASK set to 1 |
100 // +-------------------------------+-------------------------------+
101 // | Masking-key (continued) | Payload Data |
102 // +-------------------------------- - - - - - - - - - - - - - - - +
103 // : Payload Data continued ... :
104 // + - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - - +
105 // | Payload Data continued ... |
106 // +---------------------------------------------------------------+
107 //
108
109 /// A websocket connection.
110
111 template<typename Adaptor, typename Handler>
112 class Connection : public connection, public std::enable_shared_from_this<Connection<Adaptor, Handler>>
113 {
114 public:
115 /// Factory for a connection.
116 ///
117 /// Requires a request with an "Upgrade: websocket" header.<br>
118 /// Automatically handles the handshake.
119 static void create(const crow::request& req, Adaptor adaptor, Handler* handler,
120 uint64_t max_payload, const std::vector<std::string>& subprotocols,
121 std::function<void(crow::websocket::connection&)> open_handler,
122 std::function<void(crow::websocket::connection&, const std::string&, bool)> message_handler,
123 std::function<void(crow::websocket::connection&, const std::string&, uint16_t)> close_handler,
124 std::function<void(crow::websocket::connection&, const std::string&)> error_handler,
125 std::function<void(const crow::request&, std::optional<crow::response>&, void**)> accept_handler,
126 bool mirror_protocols,
127 const detail::socket::tcp_socket_options& tcp_options = {})
128 {
129 auto conn = std::shared_ptr<Connection>(new Connection(std::move(adaptor),
130 handler, max_payload,
131 std::move(open_handler),
132 std::move(message_handler),
133 std::move(close_handler),
134 std::move(error_handler),
135 std::move(accept_handler)));
136
137 // Apply TCP socket options to WebSocket connection
138 detail::socket::apply_tcp_socket_options(conn->adaptor_.socket(), tcp_options);
139
140 // Perform handshake validation
141 if (!utility::string_equals(req.get_header_value("upgrade"), "websocket"))
142 {
143 conn->adaptor_.close();
144 return;
145 }
146
147 std::string requested_subprotocols_header = req.get_header_value("Sec-WebSocket-Protocol");
148 if (!subprotocols.empty() || !requested_subprotocols_header.empty())
149 {
150 auto requested_subprotocols = utility::split(requested_subprotocols_header, ", ");
151 auto subprotocol = utility::find_first_of(subprotocols.begin(), subprotocols.end(), requested_subprotocols.begin(), requested_subprotocols.end());
152 if (subprotocol != subprotocols.end())
153 {
154 conn->subprotocol_ = *subprotocol;
155 }
156 }
157
158 if (mirror_protocols & !requested_subprotocols_header.empty())
159 {
160 conn->subprotocol_ = requested_subprotocols_header;
161 }
162
163 if (conn->accept_handler_)
164 {
165 void* ud = nullptr;
166 std::optional<crow::response> res;
167 conn->accept_handler_(req, res, &ud);
168 if (res)
169 {
170 std::vector<asio::const_buffer> buffers;
171 auto server_name = "";
172 std::string content_length_buffer;
173 res->write_header_into_buffer(buffers, content_length_buffer, req.keep_alive, server_name);
174 buffers.emplace_back(res->body.data(), res->body.size());
175 error_code ec;
176 asio::write(conn->adaptor_.socket(), buffers, ec);
177 conn->adaptor_.close();
178 return;
179 }
180 conn->userdata(ud);
181 }
182
183 // Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==
184 // Sec-WebSocket-Version: 13
185 std::string magic = req.get_header_value("Sec-WebSocket-Key") + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11";
186 sha1::SHA1 s;
187 s.processBytes(magic.data(), magic.size());
188 uint8_t digest[20];
189 s.getDigestBytes(digest);
190
191 conn->handler_->add_websocket(conn);
192 conn->start(crow::utility::base64encode((unsigned char*)digest, 20));
193 }
194
195 ~Connection() noexcept override = default;
196
197 template<typename Callable>
199 {
200 Callable callable;
201 std::weak_ptr<void> watch;
202
203 void operator()()
204 {
205 if (auto anchor = watch.lock())
206 {
207 std::move(callable)();
208 }
209 }
210 };
211
212 /// Send data through the socket.
213 template<typename CompletionHandler>
214 void dispatch(CompletionHandler&& handler)
215 {
216 asio::dispatch(adaptor_.get_io_context(),
217 WeakWrappedMessage<typename std::decay<CompletionHandler>::type>{
218 std::forward<CompletionHandler>(handler), anchor_});
219 }
220
221 /// Send data through the socket and return immediately.
222 template<typename CompletionHandler>
223 void post(CompletionHandler&& handler)
224 {
225 asio::post(adaptor_.get_io_context(),
226 WeakWrappedMessage<typename std::decay<CompletionHandler>::type>{
227 std::forward<CompletionHandler>(handler), anchor_});
228 }
229
230 /// Send a "Ping" message.
231
232 ///
233 /// Usually invoked to check if the other point is still online.
234 void send_ping(std::string msg) override
235 {
236 send_data(0x9, std::move(msg));
237 }
238
239 /// Send a "Pong" message.
240
241 ///
242 /// Usually automatically invoked as a response to a "Ping" message.
243 void send_pong(std::string msg) override
244 {
245 send_data(0xA, std::move(msg));
246 }
247
248 /// Send a binary encoded message.
249 void send_binary(std::string msg) override
250 {
251 send_data(0x2, std::move(msg));
252 }
253
254 /// Send a plaintext message.
255 void send_text(std::string msg) override
256 {
257 send_data(0x1, std::move(msg));
258 }
259
260 /// Send a close signal.
261
262 ///
263 /// Sets a flag to destroy the object once the message is sent.
264 void close(std::string const& msg, uint16_t status_code) override
265 {
266 dispatch([shared_this = this->shared_from_this(), msg, status_code]() mutable {
267 shared_this->has_sent_close_ = true;
268 if (shared_this->has_recv_close_ && !shared_this->is_close_handler_called_)
269 {
270 shared_this->is_close_handler_called_ = true;
271 if (shared_this->close_handler_)
272 shared_this->close_handler_(*shared_this, msg, status_code);
273 }
274 auto header = shared_this->build_header(0x8, msg.size() + 2);
275 char status_buf[2];
276 *(uint16_t*)(status_buf) = htons(status_code);
277
278 shared_this->write_buffers_.emplace_back(std::move(header));
279 shared_this->write_buffers_.emplace_back(std::string(status_buf, 2));
280 shared_this->write_buffers_.emplace_back(msg);
281 shared_this->do_write();
282 });
283 }
284
285 std::string get_remote_ip() override
286 {
287 return adaptor_.address();
288 }
289
290 /// Returns the TCP port of the remote endpoint, 0 when unavailable (Unix domain socket or socket already closed).
291 uint16_t get_remote_port() override
292 {
293 return adaptor_.remote_port();
294 }
295
296 void set_max_payload_size(uint64_t payload)
297 {
298 max_payload_bytes_ = payload;
299 }
300
301 /// Returns the matching client/server subprotocol, empty string if none matched.
302 std::string get_subprotocol() const override
303 {
304 return subprotocol_;
305 }
306
307 protected:
308 /// Generate the websocket headers using an opcode and the message size (in bytes).
309 std::string build_header(int opcode, size_t size)
310 {
311 char buf[2 + 8] = "\x80\x00";
312 buf[0] += opcode;
313 if (size < 126)
314 {
315 buf[1] += static_cast<char>(size);
316 return {buf, buf + 2};
317 }
318 else if (size < 0x10000)
319 {
320 buf[1] += 126;
321 *(uint16_t*)(buf + 2) = htons(static_cast<uint16_t>(size));
322 return {buf, buf + 4};
323 }
324 else
325 {
326 buf[1] += 127;
327 *reinterpret_cast<uint64_t*>(buf + 2) = ((1 == htonl(1)) ? static_cast<uint64_t>(size) : (static_cast<uint64_t>(htonl((size)&0xFFFFFFFF)) << 32) | htonl(static_cast<uint64_t>(size) >> 32));
328 return {buf, buf + 10};
329 }
330 }
331
332 /// Send the HTTP upgrade response.
333
334 ///
335 /// Finishes the handshake process, then starts reading messages from the socket.
336 void start(std::string&& hello)
337 {
338 static const std::string header =
339 "HTTP/1.1 101 Switching Protocols\r\n"
340 "Upgrade: websocket\r\n"
341 "Connection: Upgrade\r\n"
342 "Sec-WebSocket-Accept: ";
343 write_buffers_.emplace_back(header);
344 write_buffers_.emplace_back(std::move(hello));
345 write_buffers_.emplace_back(crlf);
346 if (!subprotocol_.empty())
347 {
348 write_buffers_.emplace_back("Sec-WebSocket-Protocol: ");
349 write_buffers_.emplace_back(subprotocol_);
350 write_buffers_.emplace_back(crlf);
351 }
352 write_buffers_.emplace_back(crlf);
353 do_write();
354 if (open_handler_)
355 open_handler_(*this);
356 do_read();
357 }
358
359 /// Read a websocket message.
360
361 ///
362 /// Involves:<br>
363 /// Handling headers (opcodes, size).<br>
364 /// Unmasking the payload.<br>
365 /// Reading the actual payload.<br>
366 void do_read()
367 {
368 if (has_sent_close_ && has_recv_close_)
369 {
370 close_connection_ = true;
371 adaptor_.shutdown_readwrite();
372 adaptor_.close();
374 return;
375 }
376
377 is_reading = true;
378 switch (state_)
379 {
380 case WebSocketReadState::MiniHeader:
381 {
382 mini_header_ = 0;
383 //asio::async_read(adaptor_.socket(), asio::buffer(&mini_header_, 1),
384 adaptor_.socket().async_read_some(
385 asio::buffer(&mini_header_, 2),
386 [shared_this = this->shared_from_this()](const error_code& ec, std::size_t
387#ifdef CROW_ENABLE_DEBUG
388 bytes_transferred
389#endif
390 )
391
392 {
393 shared_this->is_reading = false;
394 shared_this->mini_header_ = ntohs(shared_this->mini_header_);
395#ifdef CROW_ENABLE_DEBUG
396
397 if (!ec && bytes_transferred != 2)
398 {
399 throw std::runtime_error("WebSocket:MiniHeader:async_read fail:asio bug?");
400 }
401#endif
402
403 if (!ec)
404 {
405 if ((shared_this->mini_header_ & 0x80) == 0x80)
406 shared_this->has_mask_ = true;
407 else //if the websocket specification is enforced and the message isn't masked, terminate the connection
408 {
409#ifndef CROW_ENFORCE_WS_SPEC
410 shared_this->has_mask_ = false;
411#else
412 shared_this->close_connection_ = true;
413 shared_this->adaptor_.shutdown_readwrite();
414 shared_this->adaptor_.close();
415 if (shared_this->error_handler_)
416 shared_this->error_handler_(*shared_this, "Client connection not masked.");
417 shared_this->check_destroy(CloseStatusCode::UnacceptableData);
418#endif
419 }
420
421 if ((shared_this->mini_header_ & 0x7f) == 127)
422 {
423 shared_this->state_ = WebSocketReadState::Len64;
424 }
425 else if ((shared_this->mini_header_ & 0x7f) == 126)
426 {
427 shared_this->state_ = WebSocketReadState::Len16;
428 }
429 else
430 {
431 shared_this->remaining_length_ = shared_this->mini_header_ & 0x7f;
432 shared_this->state_ = WebSocketReadState::Mask;
433 }
434 shared_this->do_read();
435 }
436 else
437 {
438 shared_this->close_connection_ = true;
439 shared_this->adaptor_.shutdown_readwrite();
440 shared_this->adaptor_.close();
441 if (shared_this->error_handler_)
442 shared_this->error_handler_(*shared_this, ec.message());
443 shared_this->check_destroy();
444 }
445 });
446 }
447 break;
448 case WebSocketReadState::Len16:
449 {
450 remaining_length_ = 0;
451 remaining_length16_ = 0;
452 asio::async_read(
453 adaptor_.socket(), asio::buffer(&remaining_length16_, 2),
454 [shared_this = this->shared_from_this()](const error_code& ec, std::size_t
455#ifdef CROW_ENABLE_DEBUG
456 bytes_transferred
457#endif
458 ) {
459 shared_this->is_reading = false;
460 shared_this->remaining_length16_ = ntohs(shared_this->remaining_length16_);
461 shared_this->remaining_length_ = shared_this->remaining_length16_;
462#ifdef CROW_ENABLE_DEBUG
463 if (!ec && bytes_transferred != 2)
464 {
465 throw std::runtime_error("WebSocket:Len16:async_read fail:asio bug?");
466 }
467#endif
468
469 if (!ec)
470 {
471 shared_this->state_ = WebSocketReadState::Mask;
472 shared_this->do_read();
473 }
474 else
475 {
476 shared_this->close_connection_ = true;
477 shared_this->adaptor_.shutdown_readwrite();
478 shared_this->adaptor_.close();
479 if (shared_this->error_handler_)
480 shared_this->error_handler_(*shared_this, ec.message());
481 shared_this->check_destroy();
482 }
483 });
484 }
485 break;
486 case WebSocketReadState::Len64:
487 {
488 asio::async_read(
489 adaptor_.socket(), asio::buffer(&remaining_length_, 8),
490 [shared_this = this->shared_from_this()](const error_code& ec, std::size_t
491#ifdef CROW_ENABLE_DEBUG
492 bytes_transferred
493#endif
494 ) {
495 shared_this->is_reading = false;
496 shared_this->remaining_length_ = ((1 == ntohl(1)) ? (shared_this->remaining_length_) : (static_cast<uint64_t>(ntohl((shared_this->remaining_length_)&0xFFFFFFFF)) << 32) | ntohl((shared_this->remaining_length_) >> 32));
497#ifdef CROW_ENABLE_DEBUG
498 if (!ec && bytes_transferred != 8)
499 {
500 throw std::runtime_error("WebSocket:Len16:async_read fail:asio bug?");
501 }
502#endif
503
504 if (!ec)
505 {
506 shared_this->state_ = WebSocketReadState::Mask;
507 shared_this->do_read();
508 }
509 else
510 {
511 shared_this->close_connection_ = true;
512 shared_this->adaptor_.shutdown_readwrite();
513 shared_this->adaptor_.close();
514 if (shared_this->error_handler_)
515 shared_this->error_handler_(*shared_this, ec.message());
516 shared_this->check_destroy();
517 }
518 });
519 }
520 break;
521 case WebSocketReadState::Mask:
522 if (message_.size() > max_payload_bytes_ ||
523 remaining_length_ > max_payload_bytes_ - message_.size())
524 {
525 close_connection_ = true;
526 adaptor_.close();
527 if (error_handler_)
528 error_handler_(*this, "Message length exceeds maximum payload.");
529 check_destroy(MessageTooBig);
530 }
531 else if (has_mask_)
532 {
533 asio::async_read(
534 adaptor_.socket(), asio::buffer((char*)&mask_, 4),
535 [shared_this = this->shared_from_this()](const error_code& ec, std::size_t
536#ifdef CROW_ENABLE_DEBUG
537 bytes_transferred
538#endif
539 ) {
540 shared_this->is_reading = false;
541#ifdef CROW_ENABLE_DEBUG
542 if (!ec && bytes_transferred != 4)
543 {
544 throw std::runtime_error("WebSocket:Mask:async_read fail:asio bug?");
545 }
546#endif
547
548 if (!ec)
549 {
550 shared_this->state_ = WebSocketReadState::Payload;
551 shared_this->do_read();
552 }
553 else
554 {
555 shared_this->close_connection_ = true;
556 if (shared_this->error_handler_)
557 shared_this->error_handler_(*shared_this, ec.message());
558 shared_this->adaptor_.shutdown_readwrite();
559 shared_this->adaptor_.close();
560 shared_this->check_destroy();
561 }
562 });
563 }
564 else
565 {
566 state_ = WebSocketReadState::Payload;
567 do_read();
568 }
569 break;
570 case WebSocketReadState::Payload:
571 {
572 auto to_read = static_cast<std::uint64_t>(buffer_.size());
573 if (remaining_length_ < to_read)
574 to_read = remaining_length_;
575 adaptor_.socket().async_read_some(
576 asio::buffer(buffer_, static_cast<std::size_t>(to_read)),
577 [shared_this = this->shared_from_this()](const error_code& ec, std::size_t bytes_transferred) {
578 shared_this->is_reading = false;
579
580 if (!ec)
581 {
582 shared_this->fragment_.insert(shared_this->fragment_.end(), shared_this->buffer_.begin(), shared_this->buffer_.begin() + bytes_transferred);
583 shared_this->remaining_length_ -= bytes_transferred;
584 if (shared_this->remaining_length_ == 0)
585 {
586 if (shared_this->handle_fragment())
587 {
588 shared_this->state_ = WebSocketReadState::MiniHeader;
589 shared_this->do_read();
590 }
591 }
592 else
593 shared_this->do_read();
594 }
595 else
596 {
597 shared_this->close_connection_ = true;
598 if (shared_this->error_handler_)
599 shared_this->error_handler_(*shared_this, ec.message());
600 shared_this->adaptor_.shutdown_readwrite();
601 shared_this->adaptor_.close();
602 shared_this->check_destroy();
603 }
604 });
605 }
606 break;
607 }
608 }
609
610 /// Check if the FIN bit is set.
611 bool is_FIN()
612 {
613 return mini_header_ & 0x8000;
614 }
615
616 /// Extract the opcode from the header.
617 int opcode()
618 {
619 return (mini_header_ & 0x0f00) >> 8;
620 }
621
622 /// Process the payload fragment.
623
624 ///
625 /// Unmasks the fragment, checks the opcode, merges fragments into 1 message body, and calls the appropriate handler.
627 {
628 if (has_mask_)
629 {
630 for (decltype(fragment_.length()) i = 0; i < fragment_.length(); i++)
631 {
632 fragment_[i] ^= ((char*)&mask_)[i % 4];
633 }
634 }
635 switch (opcode())
636 {
637 case 0: // Continuation
638 {
639 message_ += fragment_;
640 if (is_FIN())
641 {
642 if (message_handler_)
643 message_handler_(*this, message_, is_binary_);
644 message_.clear();
645 }
646 }
647 break;
648 case 1: // Text
649 {
650 is_binary_ = false;
651 message_ += fragment_;
652 if (is_FIN())
653 {
654 if (message_handler_)
655 message_handler_(*this, message_, is_binary_);
656 message_.clear();
657 }
658 }
659 break;
660 case 2: // Binary
661 {
662 is_binary_ = true;
663 message_ += fragment_;
664 if (is_FIN())
665 {
666 if (message_handler_)
667 message_handler_(*this, message_, is_binary_);
668 message_.clear();
669 }
670 }
671 break;
672 case 0x8: // Close
673 {
674 has_recv_close_ = true;
675
676
677 uint16_t status_code = NoStatusCodePresent;
678 std::string::size_type message_start = 2;
679 if (fragment_.size() >= 2)
680 {
681 status_code = ntohs(((uint16_t*)fragment_.data())[0]);
682 } else {
683 // no message will crash substr
684 message_start = 0;
685 }
686
687 if (!has_sent_close_)
688 {
689 close(fragment_.substr(message_start), status_code);
690 }
691 else
692 {
693
694 close_connection_ = true;
695 if (!is_close_handler_called_)
696 {
697 if (close_handler_)
698 close_handler_(*this, fragment_.substr(message_start), status_code);
699 is_close_handler_called_ = true;
700 }
701 adaptor_.shutdown_readwrite();
702 adaptor_.close();
703
704 // Close handler must have been called at this point so code does not matter
705 check_destroy();
706 return false;
707 }
708 }
709 break;
710 case 0x9: // Ping
711 {
712 send_pong(fragment_);
713 }
714 break;
715 case 0xA: // Pong
716 {
717 pong_received_ = true;
718 }
719 break;
720 }
721
722 fragment_.clear();
723 return true;
724 }
725
726 /// Send the buffers' data through the socket.
727
728 ///
729 /// Also destroys the object if the Close flag is set.
730 void do_write()
731 {
732 if (sending_buffers_.empty()) {
733 if (write_buffers_.empty()) return;
734
735 sending_buffers_.swap(write_buffers_);
736 std::vector<asio::const_buffer> buffers;
737 buffers.reserve(sending_buffers_.size());
738 for (auto &s: sending_buffers_)
739 {
740 buffers.emplace_back(asio::buffer(s));
741 }
742 auto watch = std::weak_ptr<void>{anchor_};
743 asio::async_write(
744 adaptor_.socket(), buffers,
745 [shared_this = this->shared_from_this(), watch](const error_code &ec, std::size_t /*bytes_transferred*/) {
746 auto anchor = watch.lock();
747 if (anchor == nullptr)
748 return;
749
750 if (!ec && !shared_this->close_connection_)
751 {
752 shared_this->sending_buffers_.clear();
753 if (!shared_this->write_buffers_.empty())
754 shared_this->do_write();
755 if (shared_this->has_sent_close_)
756 shared_this->close_connection_ = true;
757 }
758 else
759 {
760 shared_this->sending_buffers_.clear();
761 shared_this->close_connection_ = true;
762 shared_this->check_destroy();
763 }
764 });
765 }
766 }
767
768 /// Destroy the Connection.
769 void check_destroy(websocket::CloseStatusCode code = CloseStatusCode::ClosedAbnormally)
770 {
771 // Note that if the close handler was not yet called at this point we did not receive a close packet (or send one)
772 // and thus we use ClosedAbnormally unless instructed otherwise
773 if (!is_close_handler_called_)
774 {
775 if (close_handler_)
776 {
777 close_handler_(*this, "uncleanly", code);
778 }
779 }
780
781 handler_->remove_websocket(this->shared_from_this());
782 }
783
784
786 {
787 std::string payload;
788 Connection* self;
789 int opcode;
790
791 void operator()()
792 {
793 self->send_data_impl(this);
794 }
795 };
796
797 void send_data_impl(SendMessageType* s)
798 {
799 auto header = build_header(s->opcode, s->payload.size());
800 write_buffers_.emplace_back(std::move(header));
801 write_buffers_.emplace_back(std::move(s->payload));
802 do_write();
803 }
804
805 void send_data(int opcode, std::string&& msg)
806 {
807 SendMessageType event_arg{
808 std::move(msg),
809 this,
810 opcode};
811
812 post(std::move(event_arg));
813 }
814
815 private:
816 Connection(Adaptor&& adaptor, Handler* handler, uint64_t max_payload,
817 std::function<void(crow::websocket::connection&)> open_handler,
818 std::function<void(crow::websocket::connection&, const std::string&, bool)> message_handler,
819 std::function<void(crow::websocket::connection&, const std::string&, uint16_t)> close_handler,
820 std::function<void(crow::websocket::connection&, const std::string&)> error_handler,
821 std::function<void(const crow::request&, std::optional<crow::response>&, void**)> accept_handler):
822 adaptor_(std::move(adaptor)),
823 handler_(handler),
824 max_payload_bytes_(max_payload),
825 open_handler_(std::move(open_handler)),
826 message_handler_(std::move(message_handler)),
827 close_handler_(std::move(close_handler)),
828 error_handler_(std::move(error_handler)),
829 accept_handler_(std::move(accept_handler))
830 {}
831
832 Adaptor adaptor_;
833 Handler* handler_;
834
835 std::vector<std::string> sending_buffers_;
836 std::vector<std::string> write_buffers_;
837
838 std::array<char, 4096> buffer_;
839 bool is_binary_;
840 std::string message_;
841 std::string fragment_;
842 WebSocketReadState state_{WebSocketReadState::MiniHeader};
843 uint16_t remaining_length16_{0};
844 uint64_t remaining_length_{0};
845 uint64_t max_payload_bytes_{UINT64_MAX};
846 std::string subprotocol_;
847 bool close_connection_{false};
848 bool is_reading{false};
849 bool has_mask_{false};
850 uint32_t mask_;
851 uint16_t mini_header_;
852 bool has_sent_close_{false};
853 bool has_recv_close_{false};
854 bool error_occurred_{false};
855 bool pong_received_{false};
856 bool is_close_handler_called_{false};
857
858 std::shared_ptr<void> anchor_ = std::make_shared<int>(); // Value is just for placeholding
859
860 std::function<void(crow::websocket::connection&)> open_handler_;
861 std::function<void(crow::websocket::connection&, const std::string&, bool)> message_handler_;
862 std::function<void(crow::websocket::connection&, const std::string&, uint16_t status_code)> close_handler_;
863 std::function<void(crow::websocket::connection&, const std::string&)> error_handler_;
864 std::function<void(const crow::request&, std::optional<crow::response>&, void**)> accept_handler_;
865 };
866 } // namespace websocket
867} // namespace crow
TinySHA1 - a header only implementation of the SHA1 algorithm in C++. Based on the implementation in ...
A websocket connection.
Definition websocket.h:113
void dispatch(CompletionHandler &&handler)
Send data through the socket.
Definition websocket.h:214
void do_read()
Read a websocket message.
Definition websocket.h:366
void send_pong(std::string msg) override
Send a "Pong" message.
Definition websocket.h:243
bool handle_fragment()
Process the payload fragment.
Definition websocket.h:626
std::string build_header(int opcode, size_t size)
Generate the websocket headers using an opcode and the message size (in bytes).
Definition websocket.h:309
void close(std::string const &msg, uint16_t status_code) override
Send a close signal.
Definition websocket.h:264
void send_text(std::string msg) override
Send a plaintext message.
Definition websocket.h:255
std::string get_subprotocol() const override
Returns the matching client/server subprotocol, empty string if none matched.
Definition websocket.h:302
int opcode()
Extract the opcode from the header.
Definition websocket.h:617
void do_write()
Send the buffers' data through the socket.
Definition websocket.h:730
void start(std::string &&hello)
Send the HTTP upgrade response.
Definition websocket.h:336
bool is_FIN()
Check if the FIN bit is set.
Definition websocket.h:611
uint16_t get_remote_port() override
Returns the TCP port of the remote endpoint, 0 when unavailable (Unix domain socket or socket already...
Definition websocket.h:291
void send_ping(std::string msg) override
Send a "Ping" message.
Definition websocket.h:234
void send_binary(std::string msg) override
Send a binary encoded message.
Definition websocket.h:249
void check_destroy(websocket::CloseStatusCode code=CloseStatusCode::ClosedAbnormally)
Destroy the Connection.
Definition websocket.h:769
static void create(const crow::request &req, Adaptor adaptor, Handler *handler, uint64_t max_payload, const std::vector< std::string > &subprotocols, std::function< void(crow::websocket::connection &)> open_handler, std::function< void(crow::websocket::connection &, const std::string &, bool)> message_handler, std::function< void(crow::websocket::connection &, const std::string &, uint16_t)> close_handler, std::function< void(crow::websocket::connection &, const std::string &)> error_handler, std::function< void(const crow::request &, std::optional< crow::response > &, void **)> accept_handler, bool mirror_protocols, const detail::socket::tcp_socket_options &tcp_options={})
Definition websocket.h:119
void post(CompletionHandler &&handler)
Send data through the socket and return immediately.
Definition websocket.h:223
A tiny SHA1 algorithm implementation used internally in the Crow server (specifically in crow/websock...
Definition TinySHA1.hpp:48
The main namespace of the library. In this namespace is defined the most important classes and functi...
Definition tcp_socket_options.h:29
An HTTP request.
Definition http_request.h:47
bool keep_alive
Whether or not the server should send a connection: Keep-Alive header to the client.
Definition http_request.h:57
A base class for websocket connection.
Definition websocket.h:68