diff options
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 33 | ||||
| -rw-r--r-- | src/liana/client.h | 2 | ||||
| -rw-r--r-- | src/liana/handlers/codec.h | 4 | ||||
| -rw-r--r-- | src/liana/handlers/codec_server.c | 22 | ||||
| -rw-r--r-- | src/liana/list.c | 108 | ||||
| -rw-r--r-- | src/liana/list.h | 2 | ||||
| -rw-r--r-- | src/liana/server.c | 21 | ||||
| -rw-r--r-- | src/liana/server.h | 16 |
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; |