7#define NETWORK_USE_EXPERIMENTAL
23 : server_id_(server_id)
47 lifecycle_.wait_for_stop();
60 if (lifecycle_.is_running())
64 "QUIC server is already running",
65 "messaging_quic_server::start_server");
69 lifecycle_.set_running();
71 auto result = do_start_impl(port);
74 lifecycle_.mark_stopped();
82 if (!lifecycle_.is_running())
86 "QUIC server is not running",
87 "messaging_quic_server::stop_server");
90 if (!lifecycle_.prepare_stop())
94 "QUIC server is already stopping",
95 "messaging_quic_server::stop_server");
98 auto result = do_stop_impl();
99 lifecycle_.mark_stopped();
106 return start_server(port);
111 return stop_server();
127 auto cleanup = do_stop_impl();
128 if (cleanup.is_err())
return cleanup;
131 io_context_ = std::make_unique<asio::io_context>();
132 work_guard_ = std::make_unique<
133 asio::executor_work_guard<asio::io_context::executor_type>>(
134 asio::make_work_guard(*io_context_));
137 udp_socket_ = std::make_unique<asio::ip::udp::socket>(
139 asio::ip::udp::endpoint(asio::ip::udp::v4(), port));
142 cleanup_timer_ = std::make_unique<asio::steady_timer>(*io_context_);
148 start_cleanup_timer();
154 thread_pool_ = std::make_shared<integration::basic_thread_pool>(2);
158 io_context_future_ = thread_pool_->
submit(
164 "[messaging_quic_server] Starting io_context on thread pool");
167 "[messaging_quic_server] io_context stopped");
169 catch (
const std::exception& e)
172 "[messaging_quic_server] Exception in io_context: "
173 + std::string(e.what()));
178 + std::to_string(port));
181 catch (
const std::system_error& e)
183 if (e.code() == asio::error::address_in_use ||
184 e.code() == std::errc::address_in_use)
187 "Failed to bind to port: address already in use",
188 "messaging_quic_server::do_start_impl",
189 "Port: " + std::to_string(port));
191 else if (e.code() == asio::error::access_denied ||
192 e.code() == std::errc::permission_denied)
195 "Failed to bind to port: permission denied",
196 "messaging_quic_server::do_start_impl",
197 "Port: " + std::to_string(port));
201 "Failed to start server: " + std::string(e.what()),
202 "messaging_quic_server::do_start_impl",
203 "Port: " + std::to_string(port));
205 catch (
const std::exception& e)
208 "Failed to start server: " + std::string(e.what()),
209 "messaging_quic_server::do_start_impl",
210 "Port: " + std::to_string(port));
231 if (io_context_future_.valid())
233 io_context_future_.wait();
240 cleanup_timer_->cancel();
247 udp_socket_->cancel(ec);
248 if (udp_socket_->is_open())
250 udp_socket_->close(ec);
260 cleanup_timer_.reset();
261 thread_pool_.reset();
267 catch (
const std::exception& e)
270 "Failed to stop server: " + std::string(e.what()),
271 "messaging_quic_server::do_stop_impl",
272 "Server ID: " + server_id_);
277 -> std::vector<std::shared_ptr<session::quic_session>>
280 std::vector<std::shared_ptr<session::quic_session>> result;
282 for (
const auto& [
id, session] :
sessions_)
284 result.push_back(session);
290 -> std::shared_ptr<session::quic_session>
292 std::shared_lock<std::shared_mutex> lock(sessions_mutex_);
293 auto it = sessions_.find(session_id);
294 if (it != sessions_.end())
311 std::shared_ptr<session::quic_session> session;
313 std::unique_lock<std::shared_mutex> lock(sessions_mutex_);
314 auto it = sessions_.find(session_id);
315 if (it == sessions_.end())
319 "messaging_quic_server::disconnect_session",
320 "Session ID: " + session_id);
322 session = it->second;
328 return session->close(error_code);
335 std::vector<std::shared_ptr<session::quic_session>> sessions_to_close;
337 std::unique_lock<std::shared_mutex> lock(sessions_mutex_);
338 sessions_to_close.reserve(sessions_.size());
339 for (
auto& [
id, session] : sessions_)
341 sessions_to_close.push_back(session);
346 for (
auto& session : sessions_to_close)
350 auto result = session->close(error_code);
359 auto sessions_list = sessions();
360 for (
auto& session : sessions_list)
362 if (session && session->is_active())
364 std::vector<uint8_t> data_copy(data);
365 auto result = session->send(std::move(data_copy));
369 + session->session_id() +
": "
370 + result.error().message);
378 const std::vector<std::string>& session_ids,
381 for (
const auto& session_id : session_ids)
383 auto session = get_session(session_id);
384 if (session && session->is_active())
386 std::vector<uint8_t> data_copy(data);
387 auto result = session->send(std::move(data_copy));
391 + session_id +
": " + result.error().message);
398#if KCENON_WITH_COMMON_SYSTEM
399auto messaging_quic_server::set_monitor(
400 kcenon::common::interfaces::IMonitor* monitor) ->
void
405auto messaging_quic_server::get_monitor() const
406 ->
kcenon::common::interfaces::IMonitor*
414 if (!is_running() || !udp_socket_)
419 auto self = shared_from_this();
420 udp_socket_->async_receive_from(
421 asio::buffer(recv_buffer_),
423 [
this, self](std::error_code ec, std::size_t bytes_received)
432 if (ec != asio::error::operation_aborted)
435 "[messaging_quic_server] Receive error: " + ec.message());
437 invoke_error_callback(ec);
443 std::span<const uint8_t> packet_data(recv_buffer_.data(),
445 handle_packet(packet_data, recv_endpoint_);
453 const asio::ip::udp::endpoint& from)
463 if (header_result.is_err())
466 + from.address().to_string());
475 dcid = hdr.dest_conn_id;
477 header_result.value().first);
480 auto session = find_or_create_session(dcid, from);
484 "[messaging_quic_server] Could not find or create session for packet");
489 session->handle_packet(data);
494 const asio::ip::udp::endpoint& endpoint)
495 -> std::shared_ptr<session::quic_session>
499 std::shared_lock<std::shared_mutex> lock(sessions_mutex_);
500 for (
const auto& [
id, session] : sessions_)
502 if (session->matches_connection_id(dcid))
510 if (session_count() >= config_.max_connections)
513 "[messaging_quic_server] Connection limit reached, rejecting new connection");
518 auto session_id = generate_session_id();
521 asio::ip::udp::socket session_socket(
522 *io_context_, asio::ip::udp::endpoint(asio::ip::udp::v4(), 0));
523 session_socket.connect(endpoint);
525 auto quic_socket = std::make_shared<internal::quic_socket>(
529 if (!config_.cert_file.empty() && !config_.key_file.empty())
532 quic_socket->accept(config_.cert_file, config_.key_file);
533 if (accept_result.is_err())
536 "[messaging_quic_server] Failed to accept connection: "
537 + accept_result.error().message);
543 std::make_shared<session::quic_session>(quic_socket, session_id);
546 auto self = weak_from_this();
548 session->set_receive_callback(
549 [self, session_id](
const std::vector<uint8_t>& data)
551 if (
auto server = self.lock())
553 auto sess = server->get_session(session_id);
556 server->invoke_receive_callback(sess, data);
561 session->set_stream_receive_callback(
562 [self, session_id](uint64_t stream_id,
563 const std::vector<uint8_t>& data,
566 if (
auto server = self.lock())
568 auto sess = server->get_session(session_id);
571 server->invoke_stream_receive_callback(sess, stream_id, data, fin);
576 session->set_close_callback(
579 if (
auto server = self.lock())
581 server->on_session_close(session_id);
587 std::unique_lock<std::shared_mutex> lock(sessions_mutex_);
588 sessions_[session_id] = session;
592 session->start_session();
599 invoke_connection_callback(session);
601#if KCENON_WITH_COMMON_SYSTEM
604 monitor_->record_metric(
"active_connections",
605 static_cast<double>(session_count()));
610 + session_id +
" from " + endpoint.address().to_string());
617 auto counter = session_counter_.fetch_add(1);
618 std::ostringstream oss;
619 oss << server_id_ <<
"-" << counter;
626 std::shared_ptr<session::quic_session> session;
628 std::unique_lock<std::shared_mutex> lock(sessions_mutex_);
629 auto it = sessions_.find(session_id);
630 if (it != sessions_.end())
632 session = it->second;
643 invoke_disconnection_callback(session);
645#if KCENON_WITH_COMMON_SYSTEM
648 monitor_->record_metric(
"active_connections",
649 static_cast<double>(session_count()));
660 if (!cleanup_timer_ || !is_running())
666 cleanup_timer_->expires_after(std::chrono::seconds(30));
668 auto self = shared_from_this();
669 cleanup_timer_->async_wait(
670 [
this, self](
const std::error_code& ec)
672 if (!ec && is_running())
674 cleanup_dead_sessions();
675 start_cleanup_timer();
682 std::vector<std::string> dead_session_ids;
685 std::shared_lock<std::shared_mutex> lock(sessions_mutex_);
686 for (
const auto& [
id, session] : sessions_)
688 if (!session || !session->is_active())
690 dead_session_ids.push_back(
id);
695 for (
const auto&
id : dead_session_ids)
697 on_session_close(
id);
700 if (!dead_session_ids.empty())
703 "[messaging_quic_server] Cleaned up "
704 + std::to_string(dead_session_ids.size())
705 +
" dead sessions. Active: " + std::to_string(session_count()));
714 std::shared_ptr<session::quic_session> session) ->
void
716 callbacks_.invoke<
to_index(callback_index::connection)>(session);
720 std::shared_ptr<session::quic_session> session) ->
void
722 callbacks_.invoke<
to_index(callback_index::disconnection)>(session);
726 std::shared_ptr<session::quic_session> session,
727 const std::vector<uint8_t>& data) ->
void
729 callbacks_.invoke<
to_index(callback_index::receive)>(session, data);
733 std::shared_ptr<session::quic_session> session,
735 const std::vector<uint8_t>& data,
738 callbacks_.invoke<
to_index(callback_index::stream_receive)>(session, stream_id, data, fin);
743 callbacks_.invoke<
to_index(callback_index::error)>(ec);
752 callbacks_.set<
to_index(callback_index::connection)>(std::move(callback));
757 callbacks_.set<
to_index(callback_index::disconnection)>(std::move(callback));
762 callbacks_.set<
to_index(callback_index::receive)>(std::move(callback));
767 callbacks_.set<
to_index(callback_index::stream_receive)>(std::move(callback));
772 callbacks_.set<
to_index(callback_index::error)>(std::move(callback));
783 set_connection_callback(
784 [cb = std::move(callback)](std::shared_ptr<session::quic_session> session) {
796 set_disconnection_callback(
797 [cb = std::move(callback)](std::shared_ptr<session::quic_session> session) {
800 cb(session->session_id());
808 set_receive_callback(
809 [cb = std::move(callback)](std::shared_ptr<session::quic_session> session,
810 const std::vector<uint8_t>& data) {
813 cb(session->session_id(), data);
821 set_stream_receive_callback(
822 [cb = std::move(callback)](std::shared_ptr<session::quic_session> session,
824 const std::vector<uint8_t>& data,
828 cb(session->session_id(), stream_id, data, fin);
840 interface_error_cb_ = std::move(callback);
auto cleanup_dead_sessions() -> void
auto invoke_error_callback(std::error_code ec) -> void
Invokes the error callback.
auto is_running() const -> bool override
Checks if the server is currently running.
std::function< void(std::error_code)> error_callback_t
Callback type for errors.
auto on_session_close(const std::string &session_id) -> void
utils::lifecycle_manager lifecycle_
auto stop_server() -> VoidResult
Stops the server and releases all resources.
auto disconnect_all(uint64_t error_code=0) -> void
Disconnect all active sessions.
messaging_quic_server(std::string_view server_id)
Constructs a QUIC server with a given identifier.
auto sessions() const -> std::vector< std::shared_ptr< session::quic_session > >
Get all active sessions.
auto stop() -> VoidResult override
Stops the QUIC server.
auto disconnect_session(const std::string &session_id, uint64_t error_code=0) -> VoidResult
Disconnect a specific session.
auto broadcast(std::vector< uint8_t > &&data) -> VoidResult
Send data to all connected clients.
auto handle_packet(std::span< const uint8_t > data, const asio::ip::udp::endpoint &from) -> void
auto invoke_connection_callback(std::shared_ptr< session::quic_session > session) -> void
Invokes the connection callback.
auto start(uint16_t port) -> VoidResult override
Starts the QUIC server on the specified port.
auto set_stream_receive_callback(stream_receive_callback_t callback) -> void
Sets the callback for stream data reception (legacy version).
auto start_receive() -> void
std::shared_mutex sessions_mutex_
auto set_receive_callback(interfaces::i_quic_server::receive_callback_t callback) -> void override
Sets the callback for received data on default stream (interface version).
auto set_connection_callback(interfaces::i_quic_server::connection_callback_t callback) -> void override
Sets the callback for new connections (interface version).
auto session_count() const -> size_t
Get the number of active sessions.
auto start_server(unsigned short port) -> VoidResult
Start the server with default configuration.
std::function< void(std::shared_ptr< session::quic_session >, uint64_t, const std::vector< uint8_t > &, bool)> stream_receive_callback_t
Callback type for stream data (session, stream_id, data, fin)
auto do_start_impl(unsigned short port) -> VoidResult
QUIC-specific implementation of server start.
auto invoke_receive_callback(std::shared_ptr< session::quic_session > session, const std::vector< uint8_t > &data) -> void
Invokes the receive callback.
std::function< void(std::shared_ptr< session::quic_session >, const std::vector< uint8_t > &)> receive_callback_t
Callback type for received data (session, data)
auto generate_session_id() -> std::string
auto invoke_disconnection_callback(std::shared_ptr< session::quic_session > session) -> void
Invokes the disconnection callback.
std::function< void(std::shared_ptr< session::quic_session >)> disconnection_callback_t
Callback type for disconnections.
auto wait_for_stop() -> void override
Blocks until stop() is called.
std::map< std::string, std::shared_ptr< session::quic_session > > sessions_
auto get_session(const std::string &session_id) -> std::shared_ptr< session::quic_session >
Get a session by its ID.
~messaging_quic_server() noexcept override
Destructor; automatically calls stop_server() if running.
auto start_cleanup_timer() -> void
auto set_disconnection_callback(interfaces::i_quic_server::disconnection_callback_t callback) -> void override
Sets the callback for disconnections (interface version).
auto connection_count() const -> size_t override
Gets the number of active QUIC connections (interface version).
auto server_id() const -> const std::string &
Returns the server identifier.
auto do_stop_impl() -> VoidResult
QUIC-specific implementation of server stop.
auto find_or_create_session(const protocols::quic::connection_id &dcid, const asio::ip::udp::endpoint &endpoint) -> std::shared_ptr< session::quic_session >
auto invoke_stream_receive_callback(std::shared_ptr< session::quic_session > session, uint64_t stream_id, const std::vector< uint8_t > &data, bool fin) -> void
Invokes the stream receive callback.
auto set_stream_callback(interfaces::i_quic_server::stream_callback_t callback) -> void override
Sets the callback for stream data (interface version).
auto set_error_callback(interfaces::i_quic_server::error_callback_t callback) -> void override
Sets the callback for errors (interface version).
auto multicast(const std::vector< std::string > &session_ids, std::vector< uint8_t > &&data) -> VoidResult
Send data to specific sessions.
std::function< void(std::shared_ptr< session::quic_session >)> connection_callback_t
Callback type for new connections.
std::shared_ptr< kcenon::network::integration::thread_pool_interface > get_thread_pool()
Get current thread pool.
static network_context & instance()
Get the singleton instance.
virtual std::future< void > submit(std::function< void()> task)=0
Submit a task to the thread pool.
std::function< void( std::string_view, uint64_t, const std::vector< uint8_t > &, bool)> stream_callback_t
Callback type for stream data (session_id, stream_id, data, is_fin)
std::function< void(std::string_view, const std::vector< uint8_t > &)> receive_callback_t
Callback type for default stream data (session_id, data)
std::function< void(std::string_view, std::error_code)> error_callback_t
Callback type for errors (session_id, error)
std::function< void(std::shared_ptr< i_quic_session >)> connection_callback_t
Callback type for new connections.
std::function< void(std::string_view)> disconnection_callback_t
Callback type for disconnections (session_id)
static void report_connection_accepted()
Report a new connection accepted.
static void report_active_connections(size_t count)
Report active connections count.
QUIC Connection ID (RFC 9000 Section 5.1)
static auto parse_header(std::span< const uint8_t > data) -> Result< std::pair< packet_header, size_t > >
Parse a packet header (without header protection removal)
auto is_running() const -> bool
Checks if the component is currently running.
Feature flags for network_system.
Logger system integration interface for network_system.
#define NETWORK_LOG_WARN(msg)
#define NETWORK_LOG_INFO(msg)
#define NETWORK_LOG_ERROR(msg)
#define NETWORK_LOG_DEBUG(msg)
constexpr int internal_error
constexpr int bind_failed
constexpr int server_already_running
constexpr int server_not_started
constexpr auto to_index(E e) noexcept -> std::size_t
Helper to convert enum to std::size_t for callback_manager access.
VoidResult error_void(int code, const std::string &message, const std::string &source="network_system", const std::string &details="")
Network system metrics definitions and reporting utilities.
QUIC-specific session with stream multiplexing.
Configuration options for QUIC server.