summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-06-07 11:54:55 -0400
committerAndrew Opalach <andrew@akon.city> 2025-06-07 11:54:55 -0400
commite9475ce94ba69bd437d8cf0cf3f78062be928568 (patch)
tree63efef03577f8471b904295fc6878d7d960c00aa /src/liana
parent1e53cb62651b62ad7c03cbeb64e76546e8978ff7 (diff)
downloadcamu-e9475ce94ba69bd437d8cf0cf3f78062be928568.tar.gz
camu-e9475ce94ba69bd437d8cf0cf3f78062be928568.tar.bz2
camu-e9475ce94ba69bd437d8cf0cf3f78062be928568.zip
Clarify and fix various sink behaviors
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c2
-rw-r--r--src/liana/handlers/codec_client.c2
-rw-r--r--src/liana/list.c13
-rw-r--r--src/liana/server.c17
-rw-r--r--src/liana/vcr.c7
-rw-r--r--src/liana/vcr.h3
6 files changed, 26 insertions, 18 deletions
diff --git a/src/liana/client.c b/src/liana/client.c
index 00eaca6..5222d7a 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -290,7 +290,7 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop,
client->mask = 0;
al_array_init(client->streams);
client->reconnect = RECONNECT_NONE;
- lia_vcr_init(&client->vcr, client->loop, &client->data);
+ lia_vcr_init(&client->vcr, client->loop, &client->data, node_id);
al_str_clone(&client->addr, addr);
client->port = port;
client->connection_id = 0;
diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c
index 38b5d3b..b85fc43 100644
--- a/src/liana/handlers/codec_client.c
+++ b/src/liana/handlers/codec_client.c
@@ -4,7 +4,6 @@
#include "../../codec/ffmpeg/decoder.h"
#include "../../codec/ffmpeg/packet_ext.h"
#endif
-
#include "../../codec/stb_image/decoder.h"
#include "../../codec/spng/decoder.h"
#include "../../codec/wuffs/decoder.h"
@@ -126,6 +125,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc
struct camu_codec_stream *stream = codec->handler.stream;
AVRational time_base = stream->av.stream->time_base;
pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, time_base);
+ // @TODO: Shift av_packet_free() to vcr. This leaks right now.
#else
av_packet_free_side_data(pkt);
av_packet_free(&pkt);
diff --git a/src/liana/list.c b/src/liana/list.c
index fd58ea1..2a1458f 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -51,11 +51,17 @@ static void signal_meta(struct lia_list *list, struct lia_list_entry *entry, u8
list->callback(list->userdata, LIANA_LIST_META, entry, META_OPAQUE(meta));
}
+static void entry_unload(struct lia_list *list, struct lia_list_entry *entry)
+{
+ list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL);
+}
+
static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, bool *error)
{
u8 status;
list->callback(list->userdata, LIANA_LOAD_ENTRY, entry, &status);
if (status == LIANA_ENTRY_ERRORED) {
+ // If sequence is <0 that must mean entry is not contained in list->entries.
if (sequence >= 0) {
al_array_remove_at(list->entries, (u32)sequence);
if (list->current > sequence) {
@@ -74,6 +80,8 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e
}
*error = true;
signal_meta(list, entry, LIANA_META_ENTRY_ERRORED);
+ al_assert(!al_array_contains(list->entries, entry));
+ entry_unload(list, entry);
return false;
}
*error = false;
@@ -84,11 +92,6 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e
return false;
}
-static void entry_unload(struct lia_list *list, struct lia_list_entry *entry)
-{
- list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL);
-}
-
static void entry_ref(struct lia_list *list, struct lia_list_entry *entry)
{
list->callback(list->userdata, LIANA_REF_ENTRY, entry, NULL);
diff --git a/src/liana/server.c b/src/liana/server.c
index cc4ceed..d07cb8d 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -54,6 +54,7 @@ static void packet_pool_callback(void *userdata, struct nn_packet *packet)
static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
+ nn_thread_set_name("liana_handler");
if (conn->seek_pos != LIANA_TIMESTAMP_INVALID) {
conn->handler->seek(conn->handler, conn->seek_pos);
@@ -107,6 +108,10 @@ static void free_node(struct lia_node *node)
{
struct lia_server *server = node->server;
cch_entry_free(&node->entry);
+ al_assert(!node->requests.count);
+ al_array_free(node->requests);
+ al_assert(!node->connections.count);
+ al_array_free(node->connections);
al_array_remove(server->nodes, node);
al_free(node);
}
@@ -298,6 +303,7 @@ static void signal_callback(void *userdata)
static nn_thread_result NNWT_THREADCALL init_thread(void *userdata)
{
struct lia_node_connection *conn = (struct lia_node_connection *)userdata;
+ nn_thread_set_name("liana_node_init");
if (!conn->handler->init(conn->handler, &conn->handle)) {
conn->errored = true;
}
@@ -427,6 +433,7 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en
static nn_thread_result NNWT_THREADCALL init_duration_thread(void *userdata)
{
struct lia_node *node = (struct lia_node *)userdata;
+ nn_thread_set_name("liana_init_dur");
if (!node->handler->init(node->handler, &node->handle)) {
node->errored = true;
} else {
@@ -496,14 +503,8 @@ void lia_server_close(struct lia_server *server)
void lia_server_free(struct lia_server *server)
{
- struct lia_node *node;
- al_array_foreach(server->nodes, i, node) {
- al_assert(!node->requests.count);
- al_array_free(node->requests);
- al_assert(!node->connections.count);
- al_array_free(node->connections);
- al_free(node);
- }
+ // Assuming we joined on the event loop, server->nodes should be empty.
+ al_assert(!server->nodes.count);
al_array_free(server->nodes);
al_assert(!server->dormant_connections.count);
al_array_free(server->dormant_connections);
diff --git a/src/liana/vcr.c b/src/liana/vcr.c
index 270bc38..228d758 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -35,9 +35,10 @@ static void reset_metrics(struct lia_vcr *vcr)
vcr->metric.last_report_ts = 0;
}
-void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data)
+void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data, u16 node_id)
{
vcr->data = data;
+ vcr->node_id = node_id;
al_array_init(vcr->tracks);
al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED);
vcr->mark.buffered = VCR_BUFFER_BUFFERED;
@@ -63,9 +64,11 @@ static void return_entire_cache(struct lia_vcr_track *track)
static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata)
{
nn_thread_set_priority(NNWT_THREAD_SCHED_FIFO, 32);
-
struct lia_vcr_track *track = (struct lia_vcr_track *)userdata;
struct lia_vcr *vcr = track->vcr;
+ const char thread_name[16]; // 16 = limit.
+ al_snprintf((char *)thread_name, sizeof(thread_name), "vcr:%hu", vcr->node_id);
+ nn_thread_set_name(thread_name);
s32 state;
bool corked;
diff --git a/src/liana/vcr.h b/src/liana/vcr.h
index 94d0c47..167ce7f 100644
--- a/src/liana/vcr.h
+++ b/src/liana/vcr.h
@@ -25,6 +25,7 @@ struct lia_vcr_track {
struct lia_vcr {
struct nn_packet_stream *data;
+ u16 node_id;
array(struct lia_vcr_track *) tracks;
atomic(u64) count;
struct { u64 buffered, low; } mark;
@@ -39,7 +40,7 @@ struct lia_vcr {
} metric;
};
-void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data);
+void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data, u16 node_id);
void lia_vcr_start(struct lia_vcr *vcr);
void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track);
bool lia_vcr_is_empty(struct lia_vcr *vcr);