summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c2
-rw-r--r--src/liana/list.c75
-rw-r--r--src/liana/list.h12
-rw-r--r--src/liana/server.c74
-rw-r--r--src/liana/server.h3
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;