MayaFlux 0.5.0
Digital-First Multimedia Processing Framework
Loading...
Searching...
No Matches
NetworkSubsystem.cpp
Go to the documentation of this file.
2
5
9
11
12namespace MayaFlux::Core {
13
15 : m_config(config)
16 , m_tokens {
17 .Buffer = Buffers::ProcessingToken::EVENT_RATE,
18 .Node = Nodes::ProcessingToken::EVENT_RATE,
19 .Task = Vruta::ProcessingToken::EVENT_DRIVEN
20 }
21 , m_io_context(std::make_unique<asio::io_context>())
22{
23}
24
29
30// ─────────────────────────────────────────────────────────────────────────────
31// ISubsystem lifecycle
32// ─────────────────────────────────────────────────────────────────────────────
33
37
39{
41 "Initializing Network Subsystem...");
42
43 m_handle = &handle;
44
45 if (m_config.udp.enabled) {
47 }
48 if (m_config.tcp.enabled) {
50 }
53 }
54
56
58
59 m_ready.store(true);
60
62 "Network Subsystem initialized with {} backend(s)", m_backends.size());
63}
64
66{
67 if (!m_ready.load()) {
69 "Cannot start NetworkSubsystem: not initialized");
70 return;
71 }
72
73 if (m_running.load()) {
74 return;
75 }
76
77 {
78 std::shared_lock lock(m_backends_mutex);
79 for (auto& [transport, backend] : m_backends) {
80 backend->start();
81 }
82 }
83
84 m_work_guard = std::make_unique<asio::executor_work_guard<asio::io_context::executor_type>>(
85 m_io_context->get_executor());
86
87 m_io_thread = std::jthread([this](const std::stop_token& token) {
88 while (!token.stop_requested()) {
89 try {
90 m_io_context->run();
91 break;
92 } catch (const std::exception& e) {
94 "IO context exception: {}", e.what());
95 }
96 }
97
99 "Network IO thread exiting");
100 });
101
102 m_running.store(true);
103
105 "Network Subsystem started");
106}
107
108void NetworkSubsystem::pause()
109{
110 stop();
111}
112
113void NetworkSubsystem::resume()
114{
115 start();
116}
117
118void NetworkSubsystem::stop()
119{
120 if (!m_running.load()) {
121 return;
122 }
123
124 m_running.store(false);
125
126 {
127 std::shared_lock lock(m_backends_mutex);
128 for (auto& [transport, backend] : m_backends) {
129 backend->stop();
130 }
131 }
132
135 }
136
137 m_work_guard.reset();
138 m_io_context->stop();
139
140 if (m_io_thread.joinable()) {
141 m_io_thread.request_stop();
142 m_io_thread.join();
143 }
144
145 m_io_context->restart();
146
148 "Network Subsystem stopped");
149}
150
151void NetworkSubsystem::shutdown()
152{
153 stop();
154
155 {
156 std::unique_lock lock(m_routing_mutex);
157 m_endpoint_routing.clear();
158 }
159
160 {
161 std::unique_lock lock(m_callbacks_mutex);
162 m_endpoint_callbacks.clear();
163 }
164
165 {
166 std::unique_lock lock(m_backends_mutex);
167 for (auto& [transport, backend] : m_backends) {
168 backend->shutdown();
169 }
170 m_backends.clear();
171 }
172
175 }
176
177 if (m_network_service) {
180 m_network_service.reset();
181 }
182
183 m_ready.store(false);
184
186 "Network Subsystem shutdown complete");
187}
188
189void NetworkSubsystem::wait_until_running()
190{
191 while (!m_running.load(std::memory_order_acquire))
192 std::this_thread::yield();
193}
194
195// ─────────────────────────────────────────────────────────────────────────────
196// Backend management
197// ─────────────────────────────────────────────────────────────────────────────
198
199bool NetworkSubsystem::add_backend(std::unique_ptr<INetworkBackend> backend)
200{
201 if (!backend) {
202 return false;
203 }
204
205 NetworkTransport transport = backend->get_transport();
206
207 std::unique_lock lock(m_backends_mutex);
208
209 if (m_backends.contains(transport)) {
211 "Network backend {} already registered", backend->get_name());
212 return false;
213 }
214
215 if (!backend->initialize()) {
217 "Failed to initialize network backend: {}", backend->get_name());
218 return false;
219 }
220
221 backend->set_receive_callback(
222 [this](uint64_t id, const uint8_t* data, size_t size, std::string_view addr) {
223 on_backend_receive(id, data, size, addr);
224 });
225
226 backend->set_state_callback(
227 [this](const EndpointInfo& info, EndpointState prev, EndpointState curr) {
228 on_backend_state_change(info, prev, curr);
229 });
230
231 if (auto* tcp = dynamic_cast<TCPBackend*>(backend.get())) {
232 tcp->set_endpoint_id_allocator([this]() -> uint64_t {
233 return m_next_endpoint_id.fetch_add(1);
234 });
235 }
236
238 "Added network backend: {}", backend->get_name());
239
240 m_backends[transport] = std::move(backend);
241 return true;
242}
243
244INetworkBackend* NetworkSubsystem::get_backend(NetworkTransport transport) const
245{
246 std::shared_lock lock(m_backends_mutex);
247 auto it = m_backends.find(transport);
248 return (it != m_backends.end()) ? it->second.get() : nullptr;
249}
250
251std::vector<INetworkBackend*> NetworkSubsystem::get_backends() const
252{
253 std::shared_lock lock(m_backends_mutex);
254 std::vector<INetworkBackend*> result;
255 result.reserve(m_backends.size());
256 for (const auto& [transport, backend] : m_backends) {
257 result.push_back(backend.get());
258 }
259 return result;
260}
261
262// ─────────────────────────────────────────────────────────────────────────────
263// Endpoint management
264// ─────────────────────────────────────────────────────────────────────────────
265
266uint64_t NetworkSubsystem::open_endpoint(const EndpointInfo& info)
267{
268 auto* backend = get_backend(info.transport);
269 if (!backend) {
271 "No backend for transport {}", static_cast<int>(info.transport));
272 return 0;
273 }
274
275 EndpointInfo ep = info;
276 ep.id = m_next_endpoint_id.fetch_add(1);
277
278 uint64_t backend_id = backend->open_endpoint(ep);
279 if (backend_id == 0) {
280 return 0;
281 }
282
283 {
284 std::unique_lock lock(m_routing_mutex);
285 m_endpoint_routing[ep.id] = info.transport;
286 }
287
288 return ep.id;
289}
290
291void NetworkSubsystem::close_endpoint(uint64_t endpoint_id)
292{
293 auto* backend = resolve_backend(endpoint_id);
294 if (!backend) {
295 return;
296 }
297
298 backend->close_endpoint(endpoint_id);
299
300 {
301 std::unique_lock lock(m_routing_mutex);
302 m_endpoint_routing.erase(endpoint_id);
303 }
304
305 {
306 std::unique_lock lock(m_callbacks_mutex);
307 m_endpoint_callbacks.erase(endpoint_id);
308 }
309}
310
311bool NetworkSubsystem::send(uint64_t endpoint_id, const uint8_t* data, size_t size)
312{
313 auto* backend = resolve_backend(endpoint_id);
314 if (!backend) {
315 return false;
316 }
317 return backend->send(endpoint_id, data, size);
318}
319
320bool NetworkSubsystem::send_to(uint64_t endpoint_id, const uint8_t* data, size_t size,
321 const std::string& address, uint16_t port)
322{
323 auto* backend = resolve_backend(endpoint_id);
324 if (!backend) {
325 return false;
326 }
327 return backend->send_to(endpoint_id, data, size, address, port);
328}
329
330EndpointState NetworkSubsystem::get_endpoint_state(uint64_t endpoint_id) const
331{
332 auto* backend = resolve_backend(endpoint_id);
333 if (!backend) {
334 return EndpointState::CLOSED;
335 }
336 return backend->get_endpoint_state(endpoint_id);
337}
338
339void NetworkSubsystem::set_endpoint_receive_callback(uint64_t endpoint_id,
340 NetworkReceiveCallback callback)
341{
342 std::unique_lock lock(m_callbacks_mutex);
343 m_endpoint_callbacks[endpoint_id] = std::move(callback);
344}
345
346std::vector<EndpointInfo> NetworkSubsystem::get_all_endpoints() const
347{
348 std::shared_lock lock(m_backends_mutex);
349 std::vector<EndpointInfo> result;
350 for (const auto& [transport, backend] : m_backends) {
351 auto eps = backend->get_endpoints();
352 result.insert(result.end(), eps.begin(), eps.end());
353 }
354 return result;
355}
356
357// ─────────────────────────────────────────────────────────────────────────────
358// Private: backend initialisation
359// ─────────────────────────────────────────────────────────────────────────────
360
361void NetworkSubsystem::initialize_udp_backend()
362{
363 auto udp = std::make_unique<UDPBackend>(m_config.udp, *m_io_context);
364 add_backend(std::move(udp));
365}
366
367void NetworkSubsystem::initialize_tcp_backend()
368{
369 auto tcp = std::make_unique<TCPBackend>(m_config.tcp, *m_io_context);
370 add_backend(std::move(tcp));
371}
372
373void NetworkSubsystem::initialize_shm_backend()
374{
376 "SharedMemory backend not yet implemented");
377}
378
379void NetworkSubsystem::register_backend_service()
380{
381 auto& registry = Registry::BackendRegistry::instance();
382
383 auto service = std::make_shared<Registry::Service::NetworkService>();
384
385 service->open_endpoint = [this](const EndpointInfo& info) {
386 return open_endpoint(info);
387 };
388
389 service->close_endpoint = [this](uint64_t id) {
390 close_endpoint(id);
391 };
392
393 service->send = [this](uint64_t id, const uint8_t* data, size_t size) {
394 return send(id, data, size);
395 };
396
397 service->send_to = [this](uint64_t id, const uint8_t* data, size_t size,
398 const std::string& addr, uint16_t port) {
399 return send_to(id, data, size, addr, port);
400 };
401
402 service->get_endpoint_state = [this](uint64_t id) {
403 return get_endpoint_state(id);
404 };
405
406 service->set_endpoint_receive_callback = [this](uint64_t id, NetworkReceiveCallback cb) {
407 set_endpoint_receive_callback(id, std::move(cb));
408 };
409
410 service->get_all_endpoints = [this]() {
411 return get_all_endpoints();
412 };
413
414 m_network_service = service;
415
416 registry.register_service<Registry::Service::NetworkService>(
417 [service]() -> void* {
418 return service.get();
419 });
420}
421
422// ─────────────────────────────────────────────────────────────────────────────
423// Private: callback routing
424// ─────────────────────────────────────────────────────────────────────────────
425
426void NetworkSubsystem::on_backend_receive(uint64_t endpoint_id, const uint8_t* data,
427 size_t size, std::string_view sender_addr)
428{
429 std::shared_lock lock(m_callbacks_mutex);
430 auto it = m_endpoint_callbacks.find(endpoint_id);
431 if (it != m_endpoint_callbacks.end() && it->second) {
432 it->second(endpoint_id, data, size, sender_addr);
433 }
434}
435
436void NetworkSubsystem::on_backend_state_change(const EndpointInfo& info,
438{
440 "Endpoint {} state: {} -> {}",
441 info.id, static_cast<int>(previous), static_cast<int>(current));
442}
443
444// ─────────────────────────────────────────────────────────────────────────────
445// Private: routing
446// ─────────────────────────────────────────────────────────────────────────────
447
448INetworkBackend* NetworkSubsystem::resolve_backend(uint64_t endpoint_id) const
449{
450 NetworkTransport transport {};
451 {
452 std::shared_lock lock(m_routing_mutex);
453 auto it = m_endpoint_routing.find(endpoint_id);
454 if (it == m_endpoint_routing.end()) {
455 return nullptr;
456 }
457 transport = it->second;
458 }
459
460 return get_backend(transport);
461}
462
463} // namespace MayaFlux::Core
#define MF_INFO(comp, ctx,...)
#define MF_ERROR(comp, ctx,...)
#define MF_WARN(comp, ctx,...)
#define MF_DEBUG(comp, ctx,...)
glm::vec2 current
Abstract interface for network transport backends.
std::unique_ptr< asio::io_context > m_io_context
void shutdown() override
Shutdown and cleanup subsystem resources.
std::unique_ptr< asio::executor_work_guard< asio::io_context::executor_type > > m_work_guard
void initialize(SubsystemProcessingHandle &handle) override
Initialize with a handle provided by SubsystemManager.
SubsystemProcessingHandle * m_handle
void start() override
Start the subsystem's processing/event loops.
NetworkSubsystem(const GlobalNetworkConfig &config)
void register_callbacks() override
Register callback hooks for this domain.
std::unordered_map< NetworkTransport, std::unique_ptr< INetworkBackend > > m_backends
std::shared_ptr< Registry::Service::NetworkService > m_network_service
Unified interface combining buffer and node processing for subsystems.
Connection-oriented reliable stream transport over TCP via standalone Asio.
static BackendRegistry & instance()
Get the global registry instance.
void unregister_service()
Unregister a service.
EndpointState
Observable connection state for an endpoint.
NetworkTransport
Identifies the transport protocol a backend implements.
std::function< void(uint64_t endpoint_id, const uint8_t *data, size_t size, std::string_view sender_addr)> NetworkReceiveCallback
Callback signature for inbound data on an endpoint.
@ Shutdown
Engine/subsystem shutdown and cleanup.
@ NetworkSubsystem
Network subsystem operations (endpoint management, data routing)
@ Init
Engine/subsystem initialization.
@ Core
Core engine, backend, subsystems.
void stop()
Stop active Portal::Network operations.
Definition Network.cpp:36
bool initialize(Registry::Service::NetworkService *service)
Initialize Portal::Network.
Definition Network.cpp:11
void shutdown()
Shutdown Portal::Network and release all resources.
Definition Network.cpp:50
bool is_initialized()
Return true if Portal::Network has been initialized.
Definition Network.cpp:66
Describes one logical send/receive endpoint managed by a backend.
Configuration for the NetworkSubsystem.
Backend network transport service interface.