summaryrefslogtreecommitdiff
path: root/src/liana/server.c
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2026-09-14 08:57:42 -0400
committerAndrew Opalach <andrew@akon.city> 2026-09-14 08:57:42 -0400
commit8f208c26b6fa1a9f3372679c047cab559c06e26b (patch)
tree323d894d6ff8e1ed1445c40cb1e2f5d3cee5e8e8 /src/liana/server.c
parentc66c7c64ebd16287b892f8a780cffcabafba3799 (diff)
downloadcamu-8f208c26b6fa1a9f3372679c047cab559c06e26b.tar.gz
camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.tar.bz2
camu-8f208c26b6fa1a9f3372679c047cab559c06e26b.zip
Server-side fixes from DIRECT_MODE testing
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana/server.c')
-rw-r--r--src/liana/server.c160
1 files changed, 118 insertions, 42 deletions
diff --git a/src/liana/server.c b/src/liana/server.c
index 791fe0e..8305c61 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -1,6 +1,9 @@
#define AL_LOG_SECTION "liana"
+//#define AL_LOG_ENABLE_TRACE
#include <al/log.h>
+#include "../server/common.h"
+
#include "server.h"
#include "handlers.h"
#include "list.h"
@@ -29,6 +32,7 @@ static void data_packet_sent_callback(void *userdata, struct nn_packet *packet)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
nn_packet_pool_lock(&conn->pool);
+ disown_packet(packet);
nn_packet_pool_return(&conn->pool, packet);
nn_packet_pool_unlock(&conn->pool);
}
@@ -36,11 +40,12 @@ static void data_packet_sent_callback(void *userdata, struct nn_packet *packet)
static void data_packets_sent_callback(void *userdata, struct nn_packet **packets, u32 count)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
+ // If in DIRECT_MODE and the client disconnects first, free_connection_stream() has already run.
nn_packet_pool_lock(&conn->pool);
for (u32 i = 0; i < count; i++) {
- if (packets[i]) {
- nn_packet_pool_return(&conn->pool, packets[i]);
- }
+ al_assert(packets[i]);
+ disown_packet(packets[i]);
+ nn_packet_pool_return(&conn->pool, packets[i]);
}
nn_packet_pool_unlock(&conn->pool);
}
@@ -50,11 +55,17 @@ static void packet_pool_callback(void *userdata, struct nn_packet *packet)
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
if (!nn_packet_stream_send_packet(conn->stream, packet)) {
nn_packet_pool_lock(&conn->pool);
+ disown_packet(packet);
nn_packet_pool_return(&conn->pool, packet);
nn_packet_pool_unlock(&conn->pool);
}
}
+//#define SPORADIC_ERROR_PACKET
+#ifdef SPORADIC_ERROR_PACKET
+#include <al/random.h>
+#endif
+
static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
@@ -65,25 +76,59 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
conn->seek_pos = LIANA_TIMESTAMP_INVALID;
}
+ // Set TCP_CORK with the intention to notify socket that we are likely
+ // about to write a lot of data and sending partial frames/small chunks
+ // won't be helpful.
+ struct nn_socket *sock = &conn->stream->sock;
+ if (sock->type == NNWT_SOCKET_TCP) {
+ nn_socket_set_cork(&conn->stream->sock, 1);
+#ifdef AL_LOG_ENABLE_TRACE
+ u32 rcvbuf = nn_socket_get_recv_buf(sock);
+ u32 sndbuf = nn_socket_get_send_buf(sock);
+#endif
+ nn_socket_set_recv_buf(sock, KB(8));
+ nn_socket_set_send_buf(sock, MB(2));
+#ifdef AL_LOG_ENABLE_TRACE
+ log_trace("rcvbuf: %u -> %u", rcvbuf, nn_socket_get_recv_buf(sock));
+ log_trace("sndbuf: %u -> %u", sndbuf, nn_socket_get_send_buf(sock));
+#endif
+ }
+
+ // A possible throughput optimization here would be to bundle multiple
+ // AVPackets into a single packet up to a certain size.
for (;;) {
struct nn_packet *packet = nn_packet_pool_get(&conn->pool);
if (!packet) {
- return 0;
+ goto out;
}
- conn->handler->step(conn->handler);
- conn->handler->write_packet(conn->handler, packet);
-
+#ifdef SPORADIC_ERROR_PACKET
+ bool error = al_random_int(0, 1000) == 17;
+ if (error) {
+ nn_packet_write_u8(packet, LIANA_PACKET_ERROR);
+ conn->handler->status = CAMU_ERR_ERROR;
+ } else {
+#endif
+ conn->handler->step(conn->handler);
+ conn->handler->write_packet(conn->handler, packet);
#ifdef LIANA_SERVER_LOOP
- if (conn->node->duration > 0 && conn->handler->status == CAMU_ERR_EOF) {
- nn_packet_pool_lock(&conn->pool);
- nn_packet_pool_return(&conn->pool, packet);
- nn_packet_pool_unlock(&conn->pool);
- continue;
+ if (conn->node->duration > 0 && conn->handler->status == CAMU_ERR_EOF) {
+ nn_packet_pool_lock(&conn->pool);
+ disown_packet(packet);
+ nn_packet_pool_return(&conn->pool, packet);
+ nn_packet_pool_unlock(&conn->pool);
+ continue;
+ }
+#endif
+#ifdef SPORADIC_ERROR_PACKET
}
#endif
- nn_packet_pool_submit(&conn->pool, packet);
+ if (!nn_packet_pool_submit(&conn->pool, packet)) {
+ disown_packet(packet);
+ nn_packet_pool_return(&conn->pool, packet);
+ nn_packet_pool_unlock(&conn->pool);
+ }
// Check status after submitting so the EOF packet gets sent.
if (conn->handler->status != CAMU_OK) break;
@@ -91,13 +136,17 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
nn_packet_pool_flush(&conn->pool);
+out:
+ if (sock->type == NNWT_SOCKET_TCP) {
+ nn_socket_set_cork(sock, 0);
+ }
+
return 0;
}
static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
(void)userdata;
- // @TODO: This should invalidate the connection instead of asserting.
nn_packet_stream_return_packet(stream, packet);
al_assert(false);
}
@@ -127,7 +176,7 @@ static void free_connection(struct lia_node_connection *conn)
cch_entry_return_handle(node->entry, &conn->handle);
nn_packet_pool_free(&conn->pool);
bool removed = al_array_remove(node->connections, conn);
- al_assert(removed);
+ al_assert(removed && !al_array_contains(node->connections, conn));
al_free(conn);
if (should_free_node(node)) {
free_node(node);
@@ -145,6 +194,13 @@ static void free_connection_stream(struct lia_node_connection *conn)
static void disable_connection_and_wait(struct lia_node_connection *conn)
{
nn_packet_pool_disable(&conn->pool);
+ nn_packet_pool_lock(&conn->pool);
+ struct nn_packet *packet;
+ while ((packet = nn_packet_pool_pop(&conn->pool))) {
+ disown_packet(packet);
+ nn_packet_pool_return(&conn->pool, packet);
+ }
+ nn_packet_pool_unlock(&conn->pool);
cch_handle_disable(&conn->handle);
nn_thread_join(&conn->thread);
}
@@ -188,12 +244,6 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s
start_connection_handler(conn, mask);
}
-static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet)
-{
- (void)userdata;
- nn_packet_free(packet);
-}
-
static void subscribe_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
@@ -209,6 +259,8 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
u64 mask = nn_packet_read_u64(packet);
u64 seek_pos = nn_packet_read_u64(packet);
+ nn_packet_stream_return_packet(stream, packet);
+
al_assert(!conn->ref);
conn->ref = true;
@@ -225,15 +277,12 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
nn_packet_write_str(rpacket, cch_entry_get_handler(conn->node->entry));
conn->handler->write_info(conn->handler, rpacket);
stream->packet_callback = subscribe_packet_callback;
- stream->packet_sent_callback = subscribe_packet_sent_callback;
stream->connection_closed_callback = subscribe_connection_closed_callback;
nn_packet_stream_send_packet(stream, rpacket);
} else {
start_connection_handler(conn, mask);
conn->handler->subscribe(conn->handler, mask);
}
-
- nn_packet_stream_return_packet(stream, packet);
}
static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
@@ -244,6 +293,12 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
al_free(stream);
}
+static void packet_dequeued_callback(void *userdata, struct nn_packet *packet)
+{
+ (void)userdata;
+ nn_packet_write_size(packet);
+}
+
static void packet_sent_callback(void *userdata, struct nn_packet *packet)
{
(void)userdata;
@@ -272,21 +327,18 @@ static void signal_callback(void *userdata)
nn_thread_join(&conn->thread);
nn_signal_stop(&conn->signal);
+ al_array_remove(node->requests, conn);
struct nn_packet_stream *stream = conn->stream;
struct nn_packet *packet = conn->packet;
conn->packet = NULL;
- if (packet) {
- al_array_remove(node->requests, conn);
- }
-
if (!packet || conn->errored) {
conn->handler->free(&conn->handler);
cch_entry_return_handle(node->entry, &conn->handle);
}
- if (!packet) { // Connection was closed before init was done.
+ if (!packet) { // Connection was closed before init_thread() finished.
al_free(conn);
if (should_free_node(node)) {
free_node(node);
@@ -298,10 +350,14 @@ static void signal_callback(void *userdata)
nn_packet_stream_return_packet(stream, packet);
demote_and_disconnect_stream(server, stream);
al_free(conn);
+ if (should_free_node(node)) {
+ free_node(node);
+ }
} else {
+ al_assert(!should_free_node(node));
conn->id = get_incremental_id(server);
al_array_push(node->connections, conn);
- nn_packet_pool_init(&conn->pool, 1024, server->loop, packet_pool_callback, conn);
+ nn_packet_pool_init(&conn->pool, 3072, 2, server->loop, packet_pool_callback, conn);
handle_connection(conn, packet);
}
}
@@ -328,22 +384,22 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id)
static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u32 id)
{
- struct lia_node_connection *conn;
+ struct lia_node_connection *conn, *ret = NULL;
al_array_foreach(node->connections, i, conn) {
- if (conn->id == id) return conn;
+ if (conn->id == id) {
+ al_assert(!ret);
+ ret = conn;
+ }
}
- return NULL;
+ return ret;
}
static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
- struct lia_node *node = conn->node;
struct nn_packet *packet = conn->packet;
- // Checked in signal_callback and will signal to cleanup the connection.
- conn->packet = NULL;
+ conn->packet = NULL; // Request cleanup in signal_callback().
nn_packet_stream_return_packet(stream, packet);
- al_array_remove(node->requests, conn);
}
static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
@@ -357,7 +413,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
u32 connection_id = nn_packet_read_u32(packet);
struct lia_node *node = get_node_from_id(server, node_id);
- if (!node) {
+ if (!node || node->closed) {
goto err;
}
@@ -403,7 +459,6 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str
return;
err:
- // Return packet before disconnecting.
nn_packet_stream_return_packet(stream, packet);
nn_packet_stream_disconnect(stream);
}
@@ -412,6 +467,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_server *server = (struct lia_server *)userdata;
stream->packet_callback = packet_callback;
+ stream->packet_dequeued_callback = packet_dequeued_callback;
stream->packet_sent_callback = packet_sent_callback;
al_array_push(server->dormant_connections, stream);
return true;
@@ -474,7 +530,7 @@ static void duration_signal_callback(void *userdata)
}
}
-void lia_node_get_duration(struct lia_node *node)
+void lia_node_probe_duration(struct lia_node *node)
{
// It is vital we don't block the loop during handler->init().
nn_signal_init(&node->signal, node->server->loop, duration_signal_callback, node);
@@ -492,12 +548,14 @@ void lia_node_close(struct lia_node *node)
} else {
struct lia_node_connection *conn;
al_array_foreach_rev(node->requests, i, conn) {
- nn_packet_stream_disconnect(conn->stream);
+ // We cannot disconnect the stream here because pre_init_connection_closed_callback()
+ // doesn't remove it from node->requests.
+ conn->errored = true;
}
al_array_foreach_rev(node->connections, i, conn) {
if (conn->stream) {
conn->disconnected = true;
- //nn_packet_stream_disconnect(conn->stream);
+ nn_packet_stream_disconnect(conn->stream);
} else {
free_connection(conn);
}
@@ -513,6 +571,24 @@ void lia_server_close(struct lia_server *server)
}
}
+void lia_server_abort(struct lia_server *server)
+{
+ struct lia_node *node;
+ al_array_foreach(server->nodes, i, node) {
+ struct lia_node_connection *conn;
+ al_array_foreach_rev(node->requests, j, conn) {
+ // We can't assert(conn->errored) here as we haven't fully waited on the loop.
+#ifdef NAUNET_HAS_THREAD_CANCEL
+ // This could result in a mutex_destroy() on a locked mutex.
+ nn_thread_cancel(&conn->thread);
+ nn_signal_send(&conn->signal);
+#else
+ (void)conn;
+#endif
+ }
+ }
+}
+
void lia_server_force_disconnect_nodes(struct lia_server *server)
{
struct lia_node *node;