From a2dad19f83ac00101489d7e00751ff9e17719344 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 14 Sep 2026 08:50:45 -0400 Subject: Direct mode fixes Signed-off-by: Andrew Opalach --- src/multiplex.c | 63 ++++++++++++++++++++++++++++++++++++++++----------------- 1 file changed, 45 insertions(+), 18 deletions(-) (limited to 'src/multiplex.c') diff --git a/src/multiplex.c b/src/multiplex.c index 0911540..b1db007 100644 --- a/src/multiplex.c +++ b/src/multiplex.c @@ -1,6 +1,7 @@ #define AL_LOG_SECTION "multiplex" //#define AL_LOG_ENABLE_TRACE #include +#include #include "multiplex.h" @@ -92,25 +93,53 @@ void nn_multiplex_socket_close(struct nn_multiplex_socket *multi) nn_socket_cleanup(&multi->sock); } -static struct nn_multiplex_direct multiplex_direct_global = { 0 }; +static struct nn_multiplex_direct direct_global = { 0 }; void nn_multiplex_direct_init(bool (*connection_callback)(void *, u8, struct nn_packet_stream *), void *userdata) { - multiplex_direct_global.closing_bridge = al_alloc_object(struct nn_packet_stream); - multiplex_direct_global.connection_callback = connection_callback; - multiplex_direct_global.userdata = userdata; + direct_global.connection_callback = connection_callback; + direct_global.userdata = userdata; +} + +static void queue_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) +{ + struct nn_multiplex_bridge *bridge = (struct nn_multiplex_bridge *)userdata; + (void)stream; + al_array_push(bridge->queue, packet); } static void direct_connect(struct nn_packet_stream *client, u8 id) { - struct nn_packet_stream *server = al_alloc_object(struct nn_packet_stream); - nn_packet_stream_init(server, NULL, NULL, NULL); - multiplex_direct_global.connection_callback(multiplex_direct_global.userdata, id, server); - server->direct = client; - client->direct = server; - // Explicitly signal connected on the server first. - nn_packet_stream_set_connected(server); + struct nn_multiplex_bridge *bridge = al_alloc_object(struct nn_multiplex_bridge); + bridge->broken = false; + al_array_init(bridge->queue); + bridge->bridge = al_alloc_object(struct nn_packet_stream); + nn_packet_stream_init(bridge->bridge, NULL, NULL, NULL); + // Server-side client, emulates the result of socket_accept(). + struct nn_packet_stream *cl = al_alloc_object(struct nn_packet_stream); + nn_packet_stream_init(cl, NULL, NULL, NULL); + cl->is_direct = true; + client->is_direct = true; + atomic_store(void)(&cl->direct, client, AL_ATOMIC_RELAXED); + atomic_store(void)(&client->direct, cl, AL_ATOMIC_RELAXED); + cl->bridge = bridge; + cl->packet_callback = queue_packet_callback; + cl->userdata = bridge; + // The order of connect client -> global callback -> connect server-side client cannot change. nn_packet_stream_set_connected(client); + direct_global.connection_callback(direct_global.userdata, id, cl); + nn_packet_stream_set_connected(cl); + // Resend any packets that might have been sent in client->connection_callback(). + struct nn_packet *packet; + al_array_foreach(bridge->queue, i, packet) { + al_assert(cl->packet_callback != queue_packet_callback); + cl->packet_callback(cl->userdata, cl, packet); + if (bridge->broken) break; + } + bridge->queue.count = 0; + if (bridge->broken) { + nn_multiplex_bridge_free(bridge); + } } void nn_multiplex_direct_connect(struct nn_packet_stream *client, u8 id) @@ -124,12 +153,10 @@ void nn_multiplex_direct_reconnect(struct nn_packet_stream *client) direct_connect(client, client->id); } -struct nn_packet_stream *nn_multiplex_direct_get_bridge(void) -{ - return multiplex_direct_global.closing_bridge; -} - -void nn_multiplex_direct_close(void) +void nn_multiplex_bridge_free(struct nn_multiplex_bridge *bridge) { - al_free(multiplex_direct_global.closing_bridge); + nn_packet_stream_free(bridge->bridge); + al_free(bridge->bridge); + al_array_free(bridge->queue); + al_free(bridge); } -- cgit v1.2.3-101-g0448