summaryrefslogtreecommitdiff
path: root/src/multiplex.c
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-09-14 08:50:45 -0400
committerAndrew Opalach <andrew@akon.city> 2026-09-14 08:50:45 -0400
commita2dad19f83ac00101489d7e00751ff9e17719344 (patch)
treebdaf339e76983ffe85d508e4ff3f3e4065624971 /src/multiplex.c
parent36f09a0c35657f5c7379624c8a2463aafb953bd4 (diff)
downloadlibnaunet-a2dad19f83ac00101489d7e00751ff9e17719344.tar.gz
libnaunet-a2dad19f83ac00101489d7e00751ff9e17719344.tar.bz2
libnaunet-a2dad19f83ac00101489d7e00751ff9e17719344.zip
Direct mode fixesHEADmaster
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/multiplex.c')
-rw-r--r--src/multiplex.c63
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);
}