diff options
Diffstat (limited to 'src/multiplex.c')
| -rw-r--r-- | src/multiplex.c | 63 |
1 files changed, 45 insertions, 18 deletions
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 <al/log.h> +#include <al/atomic.h> #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); } |