diff options
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 2 | ||||
| -rw-r--r-- | src/liana/list.c | 75 | ||||
| -rw-r--r-- | src/liana/list.h | 12 | ||||
| -rw-r--r-- | src/liana/server.c | 74 | ||||
| -rw-r--r-- | src/liana/server.h | 3 |
5 files changed, 116 insertions, 50 deletions
diff --git a/src/liana/client.c b/src/liana/client.c index e9dd536..bb39f3b 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -166,7 +166,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) stream->packet_sent_callback = packet_sent_callback; struct nn_packet *packet = nn_packet_create(); nn_packet_write_u32(packet, client->id); - nn_packet_write_u32(packet, 0); + nn_packet_write_u32(packet, client->connection_id); nn_packet_write_s32(packet, client->mask); nn_packet_write_u64(packet, client->pos); if (client->mask == 0) { diff --git a/src/liana/list.c b/src/liana/list.c index 2280eec..218c51e 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -45,6 +45,12 @@ enum { CLEAR }; +static inline u32 get_incremental_id(struct lia_list *list) +{ + list->increment = al_u32_inc_wrap(list->increment); + return list->increment; +} + void lia_list_init(struct lia_list *list, str *name) { al_str_clone(&list->name, name); @@ -242,11 +248,12 @@ static void adjust_current(struct lia_list *list, struct lia_list_entry *previou struct lia_list_entry *entry; al_array_foreach(list->entries, i, entry) { if (entry->opaque == previous->opaque) { - if (i == (u32)list->current) + if (i == (u32)list->current) { return; + } cmd->op = SKIPTO; cmd->sequence = i; - cmd->value.i = list->current; + cmd->arg0.i = list->current; list->current = i; break; } @@ -410,8 +417,8 @@ static bool handle_skip(struct lia_list *list, s32 sequence, s32 n) struct lia_list_cmd *cmd = list->cmd; cmd->op = SKIPTO; cmd->sequence = sequence; - cmd->value.i = sequence + n; - return handle_skipto(list, cmd->sequence, cmd->value.i); + cmd->arg0.i = sequence + n; + return handle_skipto(list, cmd->sequence, cmd->arg0.i); } static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) @@ -477,10 +484,14 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent entry = get_entry_from_id(list, id, &sequence); if (!entry) return; } + if (entry->duration == LIANA_TIMESTAMP_INVALID) { + al_log_warn("list", "Skipping seek on entry with no duration."); return; } + entry->ended = false; + entry->reset_id = get_incremental_id(list); u64 now = nn_get_timestamp(); u64 pos = (u64)(entry->duration * percent); @@ -512,20 +523,38 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent } } -static bool handle_end(struct lia_list *list, s32 id) +static bool handle_end(struct lia_list *list, u32 id, u32 reset_id) { s32 sequence; struct lia_list_entry *entry = get_entry_from_id(list, id, &sequence); - if (!entry) return true; + if (!entry) { + return true; + } + + if (reset_id != entry->reset_id) { + al_log_warn("list", "Got end() with out of order or incorrect reset id, ignoring."); + return true; + } if (entry->ended) { al_log_warn("list", "Got end() from an already ended resource, ignoring."); return true; } + al_log_debug("list", "end [#%u].", entry->id); entry->ended = true; entry->offset = entry->duration; +#ifdef LIANA_LIST_SCUFFED_LOOP + struct lia_list_cmd *cmd = list->cmd; + cmd->op = SEEK; + cmd->sequence = sequence; + cmd->arg0.u = id; + cmd->argf = 0.0; + pump_queue(list); + return false; +#endif + s32 size = (s32)list->entries.size; if (sequence == list->current) { s32 next = sequence + 1; @@ -544,7 +573,7 @@ static bool handle_end(struct lia_list *list, s32 id) struct lia_list_cmd *cmd = list->cmd; cmd->op = SKIPTO; cmd->sequence = sequence; - cmd->value.i = next; + cmd->arg0.i = next; pump_queue(list); return false; } else { @@ -626,25 +655,25 @@ static void run_queue(struct lia_list *list) handle_unset(list); break; case SKIPTO: - if (!handle_skipto(list, cmd->sequence, cmd->value.i)) { + if (!handle_skipto(list, cmd->sequence, cmd->arg0.i)) { // Target entry not loaded. return; } break; case SKIP: - if (!handle_skip(list, cmd->sequence, cmd->value.i)) { + if (!handle_skip(list, cmd->sequence, cmd->arg0.i)) { // Converted to skipto and entry not loaded. return; } break; case TOGGLE_PAUSE: - handle_toggle_pause(list, cmd->sequence, cmd->f); + handle_toggle_pause(list, cmd->sequence, cmd->argf); break; case SEEK: - handle_seek(list, cmd->sequence, cmd->value.u, cmd->f); + handle_seek(list, cmd->sequence, cmd->arg0.u, cmd->argf); break; case END: - if (!handle_end(list, cmd->value.u)) { + if (!handle_end(list, cmd->arg0.u, cmd->arg1.u)) { // End was converted to a skip. return; } @@ -702,12 +731,13 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, wstr *name) { struct lia_list_entry *entry = al_alloc_object(struct lia_list_entry); entry->opaque = opaque; - entry->id = list->increment; - list->increment = al_u32_inc_wrap(list->increment); + entry->id = get_incremental_id(list); + entry->start = LIANA_TIMESTAMP_INVALID; entry->paused_at = LIANA_TIMESTAMP_INVALID; - entry->held = false; entry->offset = 0; + entry->held = false; entry->ended = false; + entry->reset_id = get_incremental_id(list); entry->duration = duration; al_wstr_clone(&entry->name, name); entry->list = list; @@ -731,7 +761,7 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SKIPTO; cmd->sequence = sequence; - cmd->value.i = index; + cmd->arg0.i = index; al_array_push(list->queue, cmd); pump_queue(list); } @@ -741,7 +771,7 @@ void lia_list_skip(struct lia_list *list, s32 sequence, s32 n) struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SKIP; cmd->sequence = sequence; - cmd->value.i = n; + cmd->arg0.i = n; al_array_push(list->queue, cmd); pump_queue(list); } @@ -751,7 +781,7 @@ void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = TOGGLE_PAUSE; cmd->sequence = sequence; - cmd->f = pts; + cmd->argf = pts; al_array_push(list->queue, cmd); pump_queue(list); } @@ -761,17 +791,18 @@ void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent) struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SEEK; cmd->sequence = sequence; - cmd->value.u = id; - cmd->f = percent; + cmd->arg0.u = id; + cmd->argf = percent; al_array_push(list->queue, cmd); pump_queue(list); } -void lia_list_end(struct lia_list *list, u32 id) +void lia_list_end(struct lia_list *list, u32 id, u32 reset_id) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = END; - cmd->value.u = id; + cmd->arg0.u = id; + cmd->arg1.u = reset_id; al_array_push(list->queue, cmd); pump_queue(list); } diff --git a/src/liana/list.h b/src/liana/list.h index 4e5e343..bb4ac17 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -14,6 +14,8 @@ #define LIANA_BUFFER_AHEAD 2 +#define LIANA_LIST_SCUFFED_LOOP + enum { LIANA_SINK_SET = 0, LIANA_SINK_UNSET, @@ -68,9 +70,10 @@ struct lia_list_entry { u32 id; u64 start; u64 paused_at; - bool held; u64 offset; + bool held; bool ended; + u32 reset_id; u64 duration; wstr name; struct lia_list *list; @@ -89,8 +92,9 @@ struct lia_list_cmd { void *userdata; struct lia_list_entry *entry; s32 sequence; - union { s32 i; u32 u; } value; - f64 f; + union { s32 i; u32 u; } arg0; + union { s32 i; u32 u; } arg1; + f64 argf; }; struct lia_list { @@ -120,7 +124,7 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 i); void lia_list_skip(struct lia_list *list, s32 sequence, s32 n); void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts); void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent); -void lia_list_end(struct lia_list *list, u32 id); +void lia_list_end(struct lia_list *list, u32 id, u32 reset_id); void lia_list_reverse(struct lia_list *list); void lia_list_sort(struct lia_list *list); diff --git a/src/liana/server.c b/src/liana/server.c index f4a6efd..b7b8963 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -5,10 +5,16 @@ #include "handlers.h" #include "list.h" +static inline u32 get_incremental_id(struct lia_server *server) +{ + server->increment = al_u32_inc_wrap(server->increment); + return server->increment; +} + bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop) { server->loop = loop; - server->increment = 1; + server->increment = 0; al_array_init(server->nodes); al_array_init(server->zombies); return true; @@ -59,11 +65,16 @@ static void discard_packet_callback(void *userdata, struct nn_packet_stream *str al_assert(false); } -static void close_connection_internal(struct lia_node_connection *conn) +static void free_connection(struct lia_node_connection *conn) { al_assert(conn->handler); conn->handler->free(&conn->handler); cch_entry_return_handle(conn->node->entry, &conn->handle); + nn_packet_pool_free(&conn->pool); +} + +static void free_connection_stream(struct lia_node_connection *conn) +{ nn_packet_stream_free(conn->stream); al_free(conn->stream); conn->stream = NULL; @@ -73,12 +84,21 @@ static void data_connection_closed_callback(void *userdata, struct nn_packet_str { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; + nn_packet_pool_disable(&conn->pool); cch_handle_disable(&conn->handle); + // This is joining handler_thread(), we will never be here if init_thread() blocks or fails. nn_thread_join(&conn->thread); - close_connection_internal(conn); - nn_packet_pool_free(&conn->pool); + + if (conn->disconnected) { + free_connection(conn); + } else { + cch_handle_enable(&conn->handle); + nn_packet_pool_enable(&conn->pool); + } + + free_connection_stream(conn); } static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) @@ -103,8 +123,8 @@ static void subscribe_connection_closed_callback(void *userdata, struct nn_packe { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; (void)stream; - close_connection_internal(conn); - nn_packet_pool_free(&conn->pool); + free_connection(conn); + free_connection_stream(conn); } static void handle_connection(struct lia_node_connection *conn, struct nn_packet *packet) @@ -115,7 +135,9 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet u64 seek_pos = nn_packet_read_u64(packet); // Besides being wasteful, seeking to 0 on a new stream can skip data. - if (seek_pos > 0) conn->handler->seek(conn->handler, seek_pos); + if (mask != 0 || seek_pos > 0) { + conn->handler->seek(conn->handler, seek_pos); + } if (mask == 0) { struct nn_packet *rpacket = nn_packet_create(); @@ -152,20 +174,21 @@ static void signal_callback(void *userdata) nn_signal_stop(&conn->signal); nn_thread_join(&conn->thread); struct nn_packet *packet = conn->packet; + conn->packet = NULL; + 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. - close_connection_internal(conn); return; } if (!conn->errored) { - conn->id = server->increment; - server->increment = al_u32_inc_wrap(server->increment); + conn->id = get_incremental_id(server); nn_packet_pool_init(&conn->pool, 96, server->loop, packet_pool_callback, conn); al_array_push(node->connections, conn); handle_connection(conn, packet); } else { - conn->handler->free(&conn->handler); - cch_entry_return_handle(node->entry, &conn->handle); struct nn_packet_stream *stream = conn->stream; al_free(conn); conn = NULL; @@ -237,12 +260,15 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str cch_entry_get_handle(node->entry, &conn->handle); conn->handler = lia_handler_by_name(cch_entry_get_liana(node->entry))->create_server_handler(); conn->errored = false; + conn->disconnected = false; stream->packet_callback = discard_packet_callback; stream->connection_closed_callback = pre_init_connection_closed_callback; nn_thread_create(&conn->thread, init_thread, conn); } else { - // This is completely unused and connections never get removed from node->connection. + // Connections never get removed from node->connection. if ((conn = get_connection_from_id(node, connection_id))) { + stream->userdata = conn; + conn->stream = stream; handle_connection(conn, packet); } else { nn_packet_stream_disconnect(stream); @@ -278,8 +304,7 @@ void lia_server_add_socket(struct lia_server *server, struct nn_socket *sock) struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry) { struct lia_node *node = al_alloc_object(struct lia_node); - node->id = server->increment; - server->increment = al_u32_inc_wrap(server->increment); + node->id = get_incremental_id(server); node->entry = entry; al_array_init(node->connections); node->server = server; @@ -323,13 +348,17 @@ void lia_server_close(struct lia_server *server) { struct lia_node *node; al_array_foreach(server->nodes, i, node) { - struct lia_node_connection *connection; - al_array_foreach(node->connections, j, connection) { - if (connection->stream) { - nn_packet_stream_disconnect(connection->stream); + struct lia_node_connection *conn; + al_array_foreach(node->connections, j, conn) { + if (conn->stream) { + conn->disconnected = true; + nn_packet_stream_disconnect(conn->stream); + } else { + free_connection(conn); } } } + struct nn_packet_stream *zombie; al_array_foreach_rev(server->zombies, i, zombie) { nn_packet_stream_disconnect(zombie); @@ -340,14 +369,15 @@ void lia_server_free(struct lia_server *server) { struct lia_node *node; al_array_foreach(server->nodes, i, node) { - struct lia_node_connection *connection; - al_array_foreach(node->connections, j, connection) { - al_free(connection); + struct lia_node_connection *conn; + al_array_foreach(node->connections, j, conn) { + al_free(conn); } al_array_free(node->connections); al_free(node); } al_array_free(server->nodes); + struct nn_packet_stream *zombie; al_array_foreach(server->zombies, i, zombie) { al_free(zombie); diff --git a/src/liana/server.h b/src/liana/server.h index b4ada58..7c155d0 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -12,8 +12,9 @@ struct lia_node_connection { u32 id; struct nn_packet *packet; struct nn_packet_stream *stream; - struct lia_server_handler *handler; bool errored; + bool disconnected; + struct lia_server_handler *handler; struct cch_handle handle; struct nn_thread thread; struct nn_signal signal; |