diff options
Diffstat (limited to 'src/liana/server.c')
| -rw-r--r-- | src/liana/server.c | 160 |
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; |