summaryrefslogtreecommitdiff
path: root/src/liana
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-01-29 16:45:30 -0500
committerAndrew Opalach <andrew@akon.city> 2025-01-29 16:45:30 -0500
commit6134465ddc10f8b43bddf74209f0ec8654595041 (patch)
tree49eb55da3d08abdc7dea9628c3ae9674dfe7afde /src/liana
parent457a3cc1a04e45e31370d9083186436b0d12ab1d (diff)
downloadcamu-6134465ddc10f8b43bddf74209f0ec8654595041.tar.gz
camu-6134465ddc10f8b43bddf74209f0ec8654595041.tar.bz2
camu-6134465ddc10f8b43bddf74209f0ec8654595041.zip
Fixes and optimizations around seek and loop
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
-rw-r--r--src/liana/client.c33
-rw-r--r--src/liana/client.h2
-rw-r--r--src/liana/handlers/codec.h4
-rw-r--r--src/liana/handlers/codec_server.c22
-rw-r--r--src/liana/list.c108
-rw-r--r--src/liana/list.h2
-rw-r--r--src/liana/server.c21
-rw-r--r--src/liana/server.h16
8 files changed, 154 insertions, 54 deletions
diff --git a/src/liana/client.c b/src/liana/client.c
index c847ac4..1da27c1 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -10,6 +10,12 @@
#include "handlers.h"
#include "list.h"
+enum {
+ RECONNECT_NONE = 0,
+ RECONNECT_ON_CONNECTION_CLOSED,
+ RECONNECT_SIGNAL_CLIENT
+};
+
static void data_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet)
{
struct lia_client *client = (struct lia_client *)userdata;
@@ -142,7 +148,7 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream
parse_info_packet(client, packet);
nn_packet_stream_return_packet(stream, packet);
if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) {
- client->reconnect = false;
+ al_assert(client->reconnect == RECONNECT_NONE);
nn_packet_stream_disconnect(&client->data);
return;
}
@@ -162,7 +168,8 @@ static void packet_sent_callback(void *userdata, struct nn_packet *packet)
static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_client *client = (struct lia_client *)userdata;
- if (client->reconnect) {
+ if (client->reconnect == RECONNECT_SIGNAL_CLIENT) {
+ client->reconnect = RECONNECT_NONE;
// Even if seek() was called before the initial connection_callback(),
// we still want to call RESUME_AT here.
struct lia_timing time = {
@@ -177,7 +184,6 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
} else {
client->callback(client->userdata, LIANA_CLIENT_RECONNECTED, NULL, NULL);
}
- client->reconnect = false;
} else {
al_assert(client->connection_id == 0);
}
@@ -200,7 +206,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream)
{
struct lia_client *client = (struct lia_client *)userdata;
- if (client->reconnect) {
+ if (client->reconnect == RECONNECT_ON_CONNECTION_CLOSED) {
lia_vcr_flush(&client->vcr);
} else {
lia_vcr_close_all(&client->vcr);
@@ -209,7 +215,9 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream *
// This should be accounted for in the client code here to not cause
// any unexpected behavior.
client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect);
- if (client->reconnect) {
+ if (client->reconnect == RECONNECT_ON_CONNECTION_CLOSED) {
+ // If reconnect() errors, this will close the client on recursion.
+ client->reconnect = RECONNECT_SIGNAL_CLIENT;
#ifdef CAMU_DIRECT_MODE
nn_multiplex_direct_reconnect(stream);
#else
@@ -227,7 +235,7 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u
client->node_id = node_id;
client->pos = pos;
client->mask = 0;
- client->reconnect = false;
+ client->reconnect = RECONNECT_NONE;
lia_vcr_init(&client->vcr, client->loop, &client->data);
al_str_clone(&client->addr, addr);
client->port = port;
@@ -244,10 +252,12 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u
void lia_client_seek(struct lia_client *client, u64 pos, u64 at)
{
+ // If reconnect = ON_CONNECTION_CLOSED or SIGNAL_CLIENT, we are safe to edit pos
+ // and at inplace because they aren't evaluated until connection_callback().
client->pos = pos;
client->at = at;
- if (!client->reconnect) {
- client->reconnect = true;
+ if (client->reconnect == RECONNECT_NONE) {
+ client->reconnect = RECONNECT_ON_CONNECTION_CLOSED;
nn_packet_stream_disconnect(&client->data);
}
}
@@ -259,8 +269,11 @@ void lia_client_reseek(struct lia_client *client)
void lia_client_disconnect(struct lia_client *client)
{
- client->reconnect = false;
- nn_packet_stream_disconnect(&client->data);
+ u8 reconnect = client->reconnect;
+ client->reconnect = RECONNECT_NONE;
+ if (reconnect != RECONNECT_ON_CONNECTION_CLOSED) {
+ nn_packet_stream_disconnect(&client->data);
+ }
}
void lia_client_free(struct lia_client *client)
diff --git a/src/liana/client.h b/src/liana/client.h
index eca47b9..9fe676b 100644
--- a/src/liana/client.h
+++ b/src/liana/client.h
@@ -12,7 +12,7 @@ struct lia_client {
s32 mask;
u64 pos;
u64 at;
- bool reconnect;
+ u8 reconnect;
str addr;
u16 port;
u32 connection_id;
diff --git a/src/liana/handlers/codec.h b/src/liana/handlers/codec.h
index 8c5528f..29d513b 100644
--- a/src/liana/handlers/codec.h
+++ b/src/liana/handlers/codec.h
@@ -2,11 +2,15 @@
#include "../../codec/codec.h"
+#include "../server.h"
#include "../handler.h"
struct lia_codec_server {
struct lia_server_handler handler;
struct camu_demuxer *demux;
+#ifdef LIANA_SERVER_LOOP
+ u64 pts_offset;
+#endif
struct camu_codec_packet packet;
};
diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c
index ff62d47..dcbfe17 100644
--- a/src/liana/handlers/codec_server.c
+++ b/src/liana/handlers/codec_server.c
@@ -17,13 +17,20 @@
static bool codec_server_init(struct lia_server_handler *handler, struct cch_handle *handle)
{
struct lia_codec_server *codec = (struct lia_codec_server *)handler;
+
codec->demux = camu_ff_demuxer_create();
//codec->demux = camu_stbi_demuxer_create();
//codec->demux = camu_spng_demuxer_create();
//codec->demux = camu_wuffs_demuxer_create();
+
if (!codec->demux->init(codec->demux, handle)) {
return false;
}
+
+#ifdef LIANA_SERVER_LOOP
+ codec->pts_offset = 0;
+#endif
+
switch (codec->demux->mode) {
case CAMU_NORMAL:
break;
@@ -35,14 +42,17 @@ static bool codec_server_init(struct lia_server_handler *handler, struct cch_han
break;
#endif
}
+
return true;
}
static void codec_server_write_info(struct lia_server_handler *handler, struct nn_packet *packet)
{
struct lia_codec_server *codec = (struct lia_codec_server *)handler;
+
u64 duration = codec->demux->get_duration(codec->demux);
nn_packet_write_u64(packet, duration);
+
nn_packet_write_u32(packet, codec->demux->streams.count);
struct camu_codec_stream *stream;
al_array_foreach_ptr(codec->demux->streams, i, stream) {
@@ -83,6 +93,9 @@ static u64 codec_server_get_duration(struct lia_server_handler *handler)
static bool codec_server_seek(struct lia_server_handler *handler, u64 pos)
{
struct lia_codec_server *codec = (struct lia_codec_server *)handler;
+#ifdef LIANA_SERVER_LOOP
+ codec->pts_offset = 0;
+#endif
return codec->demux->seek(codec->demux, pos);
}
@@ -114,6 +127,11 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
AVPacket *pkt = codec->packet.av.pkt;
+#ifdef LIANA_SERVER_LOOP
+ struct camu_codec_stream *stream = &al_array_at(codec->demux->streams, pkt->stream_index);
+ AVRational time_base = stream->av.stream->time_base;
+ pkt->pts += av_rescale_q(codec->pts_offset, AV_TIME_BASE_Q, time_base);
+#endif
nn_packet_write_s32(packet, pkt->stream_index);
nn_packet_write_u8(packet, codec->packet.mode);
#ifdef CAMU_DIRECT_MODE
@@ -128,6 +146,10 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct
}
} else if (codec->handler.status == CAMU_ERR_EOF) {
nn_packet_write_u8(packet, LIANA_PACKET_EOF);
+#ifdef LIANA_SERVER_LOOP
+ codec->pts_offset += codec->demux->get_duration(codec->demux);
+ codec->demux->seek(codec->demux, 0);
+#endif
} else {
nn_packet_write_u8(packet, LIANA_PACKET_ERROR);
}
diff --git a/src/liana/list.c b/src/liana/list.c
index db7d54c..04ed879 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -209,6 +209,7 @@ static void handle_remove_sink(struct lia_list *list, void *userdata)
static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
{
+ // The list being idle doesn't mean list->current/sink->set isn't set.
if (list->idle) {
bool error;
if (!entry_load_and_get_duration(list, entry, -1, &error)) {
@@ -225,7 +226,6 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
};
struct lia_list_sink *sink;
al_array_foreach(list->sinks, i, sink) {
- //al_assert(sink->set == -1);
sink->set = list->current;
sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time);
}
@@ -239,27 +239,6 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry)
return true;
}
-static void adjust_current(struct lia_list *list, struct lia_list_entry *previous)
-{
- al_assert(list->current >= 0);
- struct lia_list_cmd *cmd = list->cmd;
- struct lia_list_entry *entry;
- al_array_foreach(list->entries, i, entry) {
- if (entry->opaque == previous->opaque) {
- if (i == (u32)list->current) {
- return;
- }
- cmd->op = SKIPTO;
- cmd->sequence = i;
- cmd->arg0.i = list->current;
- list->current = i;
- break;
- }
- }
- al_assert(cmd->op == SKIPTO && !entry->held);
- pump_queue(list);
-}
-
static void unset_all(struct lia_list *list)
{
list->current = -1;
@@ -579,9 +558,34 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id)
return true;
}
-static void handle_reverse(struct lia_list *list)
+static bool adjust_current(struct lia_list *list, struct lia_list_entry *previous)
+{
+ al_assert(list->current >= 0);
+ struct lia_list_cmd *cmd = list->cmd;
+ struct lia_list_entry *entry;
+ al_array_foreach(list->entries, i, entry) {
+ if (entry->opaque == previous->opaque) {
+ if (i == (u32)list->current) {
+ return true;
+ }
+ cmd->op = SKIPTO;
+ cmd->sequence = i;
+ cmd->arg0.i = list->current;
+ list->current = i;
+ break;
+ }
+ }
+ al_assert(cmd->op == SKIPTO && !entry->held);
+ pump_queue(list);
+ return false;
+}
+
+static bool handle_reverse(struct lia_list *list)
{
- if (list->current == -1) return;
+ if (list->current == -1) {
+ return true;
+ }
+
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
u32 size = list->entries.count;
for (u32 i = 0; i < size; i++) {
@@ -589,35 +593,48 @@ static void handle_reverse(struct lia_list *list)
if (tail <= i) break;
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, tail));
}
- adjust_current(list, previous);
+
+ return adjust_current(list, previous);
}
-static void handle_sort(struct lia_list *list)
+static bool handle_sort(struct lia_list *list)
{
- if (list->current == -1) return;
+ if (list->current == -1) {
+ return true;
+ }
+
struct lia_list_entry *previous = al_array_at(list->entries, list->current);
al_array_sort(list->entries, struct lia_list_entry *, camu_db_compare);
- adjust_current(list, previous);
+
+ return adjust_current(list, previous);
}
-static void handle_shuffle(struct lia_list *list)
+static bool handle_shuffle(struct lia_list *list)
{
- if (list->current == -1) return;
- struct lia_list_entry *previous = al_array_at(list->entries, list->current);
+ if (list->current == -1) {
+ return true;
+ }
+
u32 size = list->entries.count;
- if (size <= 1) return;
+ if (size <= 1) return false;
+
+ struct lia_list_entry *previous = al_array_at(list->entries, list->current);
/* https://en.wikipedia.org/wiki/Fisher%E2%80%93Yates_shuffle
for i from 0 to nāˆ’2 do
j ← random integer such that i ≤ j ≤ n-1
exchange a[i] and a[j]
*/
- for (u32 i = 0; i < size - 2; i++) {
+ // Changed size - 2 to size - 1 to support 2 entry lists.
+ // I assume this alters the algorithm, not sure how badly.
+ for (u32 i = 0; i < size - 1; i++) {
u32 j = i + (al_rand() % (size - i));
SWAP(al_array_at(list->entries, i), al_array_at(list->entries, j));
}
- adjust_current(list, previous);
+
+ return adjust_current(list, previous);
}
+/*
static void handle_clear(struct lia_list *list)
{
unset_all(list);
@@ -629,6 +646,7 @@ static void handle_clear(struct lia_list *list)
}
list->entries.count = 0;
}
+*/
static void run_queue(struct lia_list *list)
{
@@ -679,17 +697,25 @@ static void run_queue(struct lia_list *list)
}
break;
case REVERSE:
- handle_reverse(list);
- return;
+ if (!handle_reverse(list)) {
+ return;
+ }
+ break;
case SORT:
- handle_sort(list);
- return;
+ if (!handle_sort(list)) {
+ return;
+ }
+ break;
case SHUFFLE:
- handle_shuffle(list);
- return;
+ if (!handle_shuffle(list)) {
+ return;
+ }
+ break;
case CLEAR:
+ /*
handle_clear(list);
- return;
+ */
+ break;
}
al_free(cmd);
list->cmd = NULL;
diff --git a/src/liana/list.h b/src/liana/list.h
index 125dba6..0ffe652 100644
--- a/src/liana/list.h
+++ b/src/liana/list.h
@@ -7,7 +7,7 @@
#define LIANA_SEQUENCE_ANY -1
#define LIANA_TIMESTAMP_INVALID ((u64)-1)
-#define LIANA_BASE_DELAY 575000Lu // 575ms
+#define LIANA_BASE_DELAY 450000Lu // 450ms
#define LIANA_BASE_PING 150000Lu // 150ms
#define LIANA_PAUSE_DELAY LIANA_BASE_PING
#define LIANA_DELAY_IGNORE 0Lu
diff --git a/src/liana/server.c b/src/liana/server.c
index 45200d1..6459a0e 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -54,19 +54,38 @@ 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;
+
if (conn->seek_pos != LIANA_TIMESTAMP_INVALID) {
conn->handler->seek(conn->handler, conn->seek_pos);
conn->seek_pos = LIANA_TIMESTAMP_INVALID;
}
+
for (;;) {
struct nn_packet *packet = nn_packet_pool_get(&conn->pool);
- if (!packet) break;
+ if (!packet) {
+ return 0;
+ }
+
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;
+ }
+#endif
+
nn_packet_pool_submit(&conn->pool, packet);
+
+ // Check status after submitting so the EOF packet gets sent.
if (conn->handler->status != CAMU_OK) break;
}
+
nn_packet_pool_flush(&conn->pool);
+
return 0;
}
diff --git a/src/liana/server.h b/src/liana/server.h
index 7951ada..7d0b561 100644
--- a/src/liana/server.h
+++ b/src/liana/server.h
@@ -7,6 +7,22 @@
#include "../cache/entry.h"
+//#define LIANA_SERVER_LOOP
+
+// OLD: cache entry -> [connection handler -> connection]
+// NEW: cache entry -> node handler -> [connection handler -> connection]
+// [] = list of
+
+// liana supported server resources:
+// - file
+// - http
+// - hls
+// - cd
+// - dvd
+// - bluray
+// - archive
+// - live source (radio)
+
struct lia_node_connection {
u32 id;
struct nn_packet *packet;