| 355 | } |
| 356 | |
| 357 | void HttpSession::start_client_disconnect_monitor(std::shared_ptr<CancellationToken> token) { |
| 358 | if (!token || token->cancelled() || !socket_.is_open()) { |
| 359 | return; |
| 360 | } |
| 361 | |
| 362 | auto self = shared_from_this(); |
| 363 | socket_.async_wait(tcp::socket::wait_read, |
| 364 | [self, token](beast::error_code ec) { |
| 365 | if (!token || token->cancelled() || token->completed()) { |
| 366 | return; |
| 367 | } |
| 368 | |
| 369 | if (ec) { |
| 370 | if (ec != net::error::operation_aborted && !token->completed()) { |
| 371 | token->cancel(); |
| 372 | header_print("🔴 ", "Client socket wait failed; cancelling active request: " + ec.message()); |
| 373 | } |
| 374 | return; |
| 375 | } |
| 376 | |
| 377 | char peek_buffer = 0; |
| 378 | boost::system::error_code peek_ec; |
| 379 | std::size_t bytes_read = self->socket_.receive( |
| 380 | net::buffer(&peek_buffer, 1), |
| 381 | boost::asio::socket_base::message_peek, |
| 382 | peek_ec); |
| 383 | |
| 384 | if ((peek_ec || bytes_read == 0) && !token->completed()) { |
| 385 | token->cancel(); |
| 386 | header_print("🔴 ", "Client disconnected; cancelling active request"); |
| 387 | return; |
| 388 | } |
| 389 | }); |
| 390 | } |
| 391 | |
| 392 | ///@brief write streaming response |
| 393 | ///@param data the data |