diff options
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 163 | ||||
| -rw-r--r-- | src/liana/client.h | 14 | ||||
| -rw-r--r-- | src/liana/common.h | 28 | ||||
| -rw-r--r-- | src/liana/handlers/codec_client.c | 13 | ||||
| -rw-r--r-- | src/liana/handlers/codec_server.c | 13 | ||||
| -rw-r--r-- | src/liana/list.c | 151 | ||||
| -rw-r--r-- | src/liana/list.h | 4 | ||||
| -rw-r--r-- | src/liana/list_cmp.h | 64 | ||||
| -rw-r--r-- | src/liana/server.c | 160 | ||||
| -rw-r--r-- | src/liana/server.h | 2 | ||||
| -rw-r--r-- | src/liana/vcr.c | 219 | ||||
| -rw-r--r-- | src/liana/vcr.h | 11 |
12 files changed, 524 insertions, 318 deletions
diff --git a/src/liana/client.c b/src/liana/client.c index 1e873c6..ca98194 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -1,4 +1,5 @@ #define AL_LOG_SECTION "liana" +//#define AL_LOG_ENABLE_TRACE #include <al/log.h> #include <nnwt/multiplex.h> @@ -24,7 +25,14 @@ static void data_packet_callback(void *userdata, struct nn_packet_stream *stream { struct lia_client *client = (struct lia_client *)userdata; al_assert(client->vcr.data == stream); - lia_vcr_push_packet(&client->vcr, packet); + // The server might cut us off before sending any data, so don't attempt to + // recover unless we received at least one data packet. + if (!lia_vcr_push_packet(&client->vcr, packet)) { + nn_packet_stream_return_packet(stream, packet); + } else if (client->reconnect == RECONNECT_NONE) { + client->reconnect = RECONNECT_RECOVER; + lia_vcr_start(&client->vcr); + } } static s32 stream_compare(const void *a, const void *b) @@ -45,7 +53,7 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet) struct camu_codec_stream stream = { 0 }; str codec; nn_packet_read_str(packet, &codec); - stream.codec_info = camu_codec_info_by_name(&codec); + stream.codec_info = camu_codec_info_by_id(&codec); u8 mode = nn_packet_read_u8(packet); u8 type = nn_packet_read_u8(packet); u64 duration = nn_packet_read_u64(packet); @@ -76,13 +84,14 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet) AVFormatContext *format_context = avformat_alloc_context(); stream.av.stream = nn_packet_read_av_stream(format_context, av_codec, packet); switch (type) { - case CAMU_STREAM_ATTACHMENT: + case CAMU_STREAM_ATTACHMENT: { // Assume all the data we need is in the AVStream object. stream.type = CAMU_STREAM_ATTACHMENT; client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, &stream, NULL); avformat_free_context(format_context); continue; } + } stream.av.format_context = format_context; AVCodecParameters *codecpar = stream.av.stream->codecpar; switch (type) { @@ -125,11 +134,11 @@ static void collect_streams(struct lia_client *client, struct nn_packet *packet) } static const char *stream_type_to_str[] = { - [CAMU_STREAM_AUDIO] = "audio", - [CAMU_STREAM_VIDEO] = "video", - [CAMU_STREAM_SUBTITLE] = "subtitle", - [CAMU_STREAM_ATTACHMENT] = "attachment", - [CAMU_STREAM_UNKNOWN] = "unknown" + [CAMU_STREAM_AUDIO] = "Audio", + [CAMU_STREAM_VIDEO] = "Video", + [CAMU_STREAM_SUBTITLE] = "Subtitles", + [CAMU_STREAM_ATTACHMENT] = "Attachment", + [CAMU_STREAM_UNKNOWN] = "Unknown" }; static void parse_info_packet(struct lia_client *client, struct nn_packet *packet) @@ -149,49 +158,85 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe struct camu_codec_stream *stream; al_array_foreach_ptr(client->streams, i, stream) { u8 type = stream->type; + s32 index = stream->index; + switch (type) { + case CAMU_STREAM_AUDIO: + if (index < prefs->index.audio_min) prefs->index.audio_min = index; + if (index > prefs->index.audio_max) prefs->index.audio_max = index; + break; + case CAMU_STREAM_SUBTITLE: + if (index < prefs->index.subtitles_min) prefs->index.subtitles_min = index; + if (index > prefs->index.subtitles_max) prefs->index.subtitles_max = index; + break; + } if ((selected & (1 << type)) || !(prefs->enabled & (1 << type))) { continue; } const char *title = NULL; #ifdef CAMU_HAVE_FFMPEG + bool pref_is_index = false; + if (!accept_defaults) { + if (type == CAMU_STREAM_AUDIO && prefs->index.audio >= 0) { + if (index != prefs->index.audio) { + continue; + } + pref_is_index = true; + } + if (type == CAMU_STREAM_SUBTITLE && prefs->index.subtitles >= 0) { + if (index != prefs->index.subtitles) { + continue; + } + pref_is_index = true; + } + } if (stream->mode == CAMU_FFMPEG_COMPAT) { AVDictionary *metadata = stream->av.stream->metadata; const AVDictionaryEntry *title_entry = av_dict_get(metadata, "title", NULL, 0); - if (title_entry) { - title = title_entry->value; - } - const AVDictionaryEntry *lang_entry = av_dict_get(metadata, "language", NULL, 0); - if (lang_entry && !accept_defaults) { - u8 lang = CAMU_LANG_UNKNOWN; - str s = al_str_cr(lang_entry->value); - if (al_str_eq(&s, &al_str_c("eng"))) lang = CAMU_LANG_ENGLISH; - else if (al_str_eq(&s, &al_str_c("jpn"))) lang = CAMU_LANG_JAPANESE; - if ((type == CAMU_STREAM_AUDIO && lang != prefs->language.audio) || - (type == CAMU_STREAM_SUBTITLE && lang != prefs->language.subtitles)) { - continue; + if (title_entry) title = title_entry->value; + if (!pref_is_index && !accept_defaults) { + const AVDictionaryEntry *lang_entry = av_dict_get(metadata, "language", NULL, 0); + if (lang_entry) { + u8 lang = CAMU_LANG_UNKNOWN; + str s = al_str_cr(lang_entry->value); + if (al_str_eq(&s, &al_str_c("eng"))) lang = CAMU_LANG_ENGLISH; + else if (al_str_eq(&s, &al_str_c("jpn"))) lang = CAMU_LANG_JAPANESE; + if ((type == CAMU_STREAM_AUDIO && lang != prefs->language.audio) || + (type == CAMU_STREAM_SUBTITLE && lang != prefs->language.subtitles)) { + continue; + } } - break; } } #endif + switch (type) { + case CAMU_STREAM_AUDIO: + prefs->index.audio = index; + break; + case CAMU_STREAM_VIDEO: + prefs->index.video = index; + break; + case CAMU_STREAM_SUBTITLE: + prefs->index.subtitles = index; + break; + } if (title) { - log_info("Selected %s stream (index: %u, title: %s).", stream_type_to_str[type], stream->index, title); + log_info("%s selected (index: %u, title: %s).", stream_type_to_str[type], index, title); } else { - log_info("Selected %s stream (index: %u).", stream_type_to_str[type], stream->index); + log_info("%s selected (index: %u).", stream_type_to_str[type], index); } selected |= (1 << type); struct lia_vcr_track *track = al_alloc_object(struct lia_vcr_track); track->stream = stream; - track->client = lia_handler_by_name(&handler)->create_client_handler(); - track->client->callback = client->callback; - track->client->userdata = client->userdata; - if (!track->client->init(track->client, client->renderer, track->stream)) { - track->client->free(&track->client); + track->handler = lia_handler_by_name(&handler)->create_client_handler(); + track->handler->callback = client->callback; + track->handler->userdata = client->userdata; + if (!track->handler->init(track->handler, client->renderer, track->stream)) { + track->handler->free(&track->handler); al_free(track); // This will NOT attempt to select another stream. continue; } - client->mask |= 1 << stream->index; + client->mask |= 1 << index; client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, stream, track); lia_vcr_add_track(&client->vcr, track); } @@ -221,8 +266,12 @@ static void info_packet_callback(void *userdata, struct nn_packet_stream *stream struct nn_packet *rpacket = nn_packet_create(); nn_packet_write_u64(rpacket, client->mask); nn_packet_stream_send_packet(stream, rpacket); - client->reconnect = RECONNECT_RECOVER; - lia_vcr_start(&client->vcr); +} + +static void packet_dequeued_callback(void *userdata, struct nn_packet *packet) +{ + (void)userdata; + nn_packet_write_size(packet); } static void packet_sent_callback(void *userdata, struct nn_packet *packet) @@ -234,6 +283,21 @@ 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; + + struct nn_socket *sock = &client->data.sock; + if (sock->type == NNWT_SOCKET_TCP) { +#ifdef AL_LOG_ENABLE_TRACE + u32 rcvbuf = nn_socket_get_recv_buf(sock); + u32 sndbuf = nn_socket_get_send_buf(sock); +#endif + nn_socket_set_recv_buf(sock, MB(2)); + nn_socket_set_send_buf(sock, KB(8)); +#ifdef AL_LOG_ENABLE_TRACE + log_trace("rcvbuf: %u -> %u", rcvbuf, nn_socket_get_recv_buf(sock)); + log_trace("sndbuf: %u -> %u", sndbuf, nn_socket_get_send_buf(sock)); +#endif + } + if (client->reconnect == RECONNECT_SIGNAL_CLIENT) { client->reconnect = RECONNECT_NONE; // Even if client_seek() was called before the initial connection_callback(), @@ -254,20 +318,19 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) } else { al_assert(client->connection_id == 0 && client->reconnect == RECONNECT_NONE); } - stream->packet_sent_callback = packet_sent_callback; + struct nn_packet *packet = nn_packet_create(); nn_packet_write_u32(packet, client->node_id); nn_packet_write_u32(packet, client->connection_id); nn_packet_write_u64(packet, client->mask); nn_packet_write_u64(packet, client->pos); - if (client->mask == 0) { - stream->packet_callback = info_packet_callback; - } else { - stream->packet_callback = data_packet_callback; - client->reconnect = RECONNECT_RECOVER; - lia_vcr_start(&client->vcr); - } + + stream->packet_callback = (client->mask == 0) ? info_packet_callback : data_packet_callback; + stream->packet_dequeued_callback = packet_dequeued_callback; + stream->packet_sent_callback = packet_sent_callback; + nn_packet_stream_send_packet(stream, packet); + return true; } @@ -290,7 +353,8 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * lia_vcr_close_all(&client->vcr); } - // If reconnect = SIGNAL_CLIENT, we either never connected or recursed at the reconnect step. + // If reconnect = SIGNAL_CLIENT, we can guarantee that CLIENT_REMOVE_BUFFERS + // has been called after the most recent connection_callback(). if (client->reconnect != RECONNECT_SIGNAL_CLIENT) { client->rec = (struct lia_reconnect_info){ .reconnect = reconnect, @@ -298,8 +362,8 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * .mask = client->mask }; al_array_init(client->rec.detached); - // CLIENT_REMOVE_BUFFERS should be allowed to run the event loop to wait and should - // attempt to maintain the same state if called consecutively. + // CLIENT_REMOVE_BUFFERS should be allowed to run the event loop to wait. + // It should also attempt to maintain the same state if called consecutively. client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->rec); client->mask = client->rec.mask; struct camu_codec_stream *detached; @@ -316,9 +380,11 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * if (reconnect) { // If stream_reconnect() errors or is aborted, the client will be closed on recursion. client->reconnect = RECONNECT_SIGNAL_CLIENT; - // A client being seeked before an info packet is another reason mask may - // be unset here. In that case we don't want to forcefully close. - if (!client->rec.unconfigured && !client->mask) { + // It's possible the client was not configured yet. For example, if it was seeked + // before an info packet. In that case, we don't want to forcefully close it. + bool was_configured = !client->rec.unconfigured; + if (was_configured && client->mask == 0) { + // If mask was emptied by CLIENT_REMOVE_BUFFERS, close the client immediately. connection_closed_callback(userdata, stream); } else { #ifdef CAMU_DIRECT_MODE @@ -337,12 +403,13 @@ void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, { client->loop = loop; client->node_id = node_id; + client->duration = LIANA_TIMESTAMP_INVALID; + lia_vcr_init(&client->vcr, client->loop, &client->data, node_id); + al_array_init(client->streams); + client->mask = 0; client->pos = pos; client->at = LIANA_TIMESTAMP_INVALID; - client->mask = 0; - al_array_init(client->streams); client->reconnect = RECONNECT_NONE; - 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/client.h b/src/liana/client.h index 452150d..b1b4e61 100644 --- a/src/liana/client.h +++ b/src/liana/client.h @@ -29,6 +29,16 @@ struct lia_reconnect_info { struct lia_prefs { u8 enabled; + u32 node_id; + struct { + s32 audio; + s32 audio_min; + s32 audio_max; + s32 video; + s32 subtitles; + s32 subtitles_min; + s32 subtitles_max; + } index; struct { u8 audio; u8 subtitles; @@ -39,6 +49,8 @@ struct lia_client { struct nn_event_loop *loop; u32 node_id; struct lia_prefs prefs; + u64 duration; + struct lia_vcr vcr; array(struct camu_codec_stream) streams; u64 mask; u64 pos; @@ -49,8 +61,6 @@ struct lia_client { u16 port; u32 connection_id; struct nn_packet_stream data; - u64 duration; - struct lia_vcr vcr; struct camu_renderer *renderer; void (*callback)(void *, u8, struct camu_codec_stream *stream, void *); void *userdata; diff --git a/src/liana/common.h b/src/liana/common.h index 7776218..bdf4704 100644 --- a/src/liana/common.h +++ b/src/liana/common.h @@ -1,7 +1,35 @@ #pragma once +#include <nnwt/packet.h> + +#include "../server/common.h" +#include "../codec/codec.h" + #define LIANA_TIMESTAMP_INVALID ((u64)-1) #define LIANA_BASE_DELAY ((u64)1600000) // 1600ms #define LIANA_BASE_PING ((u64)600000) // 600ms #define LIANA_PAUSE_DELAY LIANA_BASE_PING + +static inline void disown_packet(struct nn_packet *packet) +{ +#ifdef CAMU_DIRECT_MODE + if (packet->opaque) { + switch (nn_packet_get_u8(packet, NNWT_PACKET_HEADER_LENGTH + 5)) { + case CAMU_NORMAL: + break; +#ifdef CAMU_HAVE_FFMPEG + case CAMU_FFMPEG_COMPAT: { + AVPacket *pkt = (AVPacket *)packet->opaque; + av_packet_free_side_data(pkt); + av_packet_free(&pkt); + break; + } +#endif + } + packet->opaque = NULL; + } +#else + (void)packet; +#endif +} diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index 5c8e94d..f7bde35 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -92,7 +92,9 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc case CAMU_NORMAL: { if (codec->dec) { if (packet->opaque) { - success = push_packet(codec, (struct nn_buffer *)packet->opaque); + struct nn_buffer *buffer = (struct nn_buffer *)packet->opaque; + success = push_packet(codec, buffer); + packet->opaque = NULL; } else { struct nn_buffer buffer; nn_packet_read_buffer(packet, &buffer); @@ -126,10 +128,10 @@ 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); + packet->opaque = NULL; #endif break; } @@ -141,12 +143,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc // Restore packet rindex in case we reuse it. packet->rindex = rindex; - if (!success) { - codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_ERRORED, codec->handler.stream, NULL); - return false; - } - - return true; + return success; } static void codec_client_flush(struct lia_client_handler *handler) diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c index ddba182..592f592 100644 --- a/src/liana/handlers/codec_server.c +++ b/src/liana/handlers/codec_server.c @@ -35,9 +35,7 @@ static bool codec_server_init(struct lia_server_handler *handler, struct cch_han break; #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: -#ifndef CAMU_DIRECT_MODE codec->packet.av.pkt = av_packet_alloc(); -#endif break; #endif } @@ -104,11 +102,6 @@ static bool codec_server_seek(struct lia_server_handler *handler, u64 pos) static void codec_server_step(struct lia_server_handler *handler) { struct lia_codec_server *codec = (struct lia_codec_server *)handler; -#ifdef CAMU_HAVE_FFMPEG -#ifdef CAMU_DIRECT_MODE - codec->packet.av.pkt = av_packet_alloc(); -#endif -#endif codec->handler.status = codec->demux->get_packet(codec->demux, &codec->packet); } @@ -139,11 +132,11 @@ static void codec_server_write_packet(struct lia_server_handler *handler, struct nn_packet_write_s32(packet, pkt->stream_index); nn_packet_write_u8(packet, codec->packet.mode); #ifdef CAMU_DIRECT_MODE - packet->opaque = pkt; + packet->opaque = av_packet_clone(pkt); #else nn_packet_write_av_packet(packet, pkt); - av_packet_unref(pkt); #endif + av_packet_unref(pkt); break; } #endif @@ -168,9 +161,7 @@ static void codec_server_free(struct lia_server_handler **handler) break; #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: -#ifndef CAMU_DIRECT_MODE av_packet_free(&codec->packet.av.pkt); -#endif break; #endif } diff --git a/src/liana/list.c b/src/liana/list.c index 01c2ff3..fc28c39 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -3,10 +3,11 @@ #include <al/log.h> #include <al/lib.h> #include <al/random.h> +#include <al/math.h> #include <nnwt/time.h> +#include <nnwt/sort.h> #include "list.h" -#include "list_cmp.h" enum { ADD_SINK = 0, @@ -44,6 +45,12 @@ void lia_list_init(struct lia_list *list, str *name) list->cmd = NULL; } +void lia_list_start_at(struct lia_list *list, s32 i) +{ + al_assert(list->current < 0); + list->current = -(i + 1); +} + // These small functions may seem excessive but their purpose is an attempt // to reduce noise in parts that are harder to understand. static inline void list_signal_meta(struct lia_list *list, struct lia_list_entry *entry, u8 meta) @@ -165,6 +172,11 @@ static inline void sink_set_entry(struct lia_list_sink *sink, struct lia_list_en sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, time); } +static inline void sink_sequence_change(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence) +{ + sink->callback(sink->userdata, LIANA_SINK_SEQUENCE, entry, sequence, NULL); +} + static inline void sink_seek_entry(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time) { sink->callback(sink->userdata, LIANA_SINK_SEEK, entry, sequence, time); @@ -181,14 +193,13 @@ static inline void sink_unset_entry(struct lia_list_sink *sink) sink->callback(sink->userdata, LIANA_SINK_UNSET, NULL, -1, NULL); } -static bool list_set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence) +static void list_set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence) { list->idle = false; entry_ref(list, entry); al_assert(list->current != sequence); list->current = sequence; list_signal_meta(list, entry, LIANA_META_CURRENT_CHANGED); - return true; } static void pump_queue(struct lia_list *list); @@ -201,7 +212,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) u8 pause; u64 at = LIANA_TIMESTAMP_INVALID; u64 pos = current->offset; - if (current->paused_at != LIANA_TIMESTAMP_INVALID) { + if (current->duration == 0 || current->paused_at != LIANA_TIMESTAMP_INVALID) { pause = LIANA_PAUSE_NONE; } else { at = nn_get_timestamp() + LIANA_BASE_DELAY; @@ -210,11 +221,7 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) } else { at = current->start; } - // This sink could have an entry set from a connection we no longer - // know about. In that case skipping to this entry with a pause_and_swap_to() - // would be better. If this sink is empty that's still okay because - // PAUSE_BOTH is required to handle that case sink-side. - pause = LIANA_PAUSE_BOTH; + pause = LIANA_PAUSE_RESUME; } struct lia_timing time = { .at = at, @@ -264,8 +271,8 @@ static void handle_remove_sink(struct lia_list *list, void *userdata) static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 sequence) { - s32 size = (s32)list->entries.count; - if (sequence < 0 || sequence >= size) { + s32 count = (s32)list->entries.count; + if (sequence < 0 || sequence >= count) { return NULL; } return al_array_at(list->entries, sequence); @@ -321,7 +328,7 @@ static u8 skipto_entry(struct lia_list_entry *current, struct lia_list_entry *ta } } - // Pause current to be resumed if it becomes the target of a skip (hold). + // Pause current to be resumed if it becomes the target of a skip (aka "hold" it). if (current->duration != 0 && current->paused_at == LIANA_TIMESTAMP_INVALID) { al_assert(current->start != LIANA_TIMESTAMP_INVALID); current->paused_at = at; @@ -355,11 +362,28 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) bool error; if (!entry_load_and_get_duration(list, target, index, &error)) { if (error) { - // On an error, sequence will address the same entry but index might not. - // entry_load_and_get_duration() may also adjust the current cmd's sequence, - // so use cmd->sequence here. + // On an error, sequence will always address the same entry but index might not. + // Keep index the same until sequence = index (search forward). Then, start + // decreasing index (search backward). entry_load_and_get_duration() may also + // adjust the current cmd's sequence, so use cmd->sequence here. struct lia_list_cmd *cmd = list->cmd; - return handle_skipto(list, cmd->sequence, index); + sequence = cmd->sequence; + if (sequence == LIANA_SEQUENCE_ANY) { + sequence = list->current; + } + if (sequence == index) { + if (index > 0) { + index--; + } else { + al_assert(sequence == 0); + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink_sequence_change(sink, current, sequence); + } + return true; // Current is now sequence 0. + } + } + return handle_skipto(list, sequence, index); } return false; } @@ -377,6 +401,8 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) .pause = pause }; + list_set_current(list, target, index); + struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { al_assert(sink->set != index); @@ -384,19 +410,26 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) sink_set_entry(sink, target, index, &time); } - return list_set_current(list, target, index); + return true; } static bool handle_skip(struct lia_list *list, s32 sequence, s32 n) { - if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; + if (sequence == LIANA_SEQUENCE_ANY) { + sequence = list->current; + } if (sequence < 0) return true; // list->current = -1 - return handle_skipto(list, sequence, sequence + n); + s32 i = sequence + n; + s32 max = (s32)list->entries.count - 1; + return handle_skipto(list, sequence, CLAMP(i, 0, max)); } static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) { - if (list->idle) { + s32 count = (s32)list->entries.count; + s32 start_at = abs(list->current) - 1; + bool on_start_index = list->current < 0 && count == start_at; + if (list->idle && on_start_index) { bool error; if (!entry_load_and_get_duration(list, entry, -1, &error)) { // We gave a sequence of -1 so, on error, don't add to list->entries. @@ -405,21 +438,22 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) } list_add_entry(list, entry); if (list->idle) { - if (list->current == -1) { // Start the list. + if (on_start_index) { // Start the list. entry->start = nn_get_timestamp() + LIANA_BASE_DELAY; struct lia_timing time = { .at = entry->start, .pos = entry->offset, - .pause = LIANA_PAUSE_RESUME + .pause = (entry->duration == 0) ? LIANA_PAUSE_NONE : LIANA_PAUSE_RESUME }; + list_set_current(list, entry, start_at); struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { al_assert(sink->set == -1); sink->set = 0; sink_set_entry(sink, entry, 0, &time); } - return list_set_current(list, entry, 0); - } else { // Skip to the added entry. + return true; + } else if (list->current >= 0) { // Skip to the added entry. // This is done via SKIPTO for consistency. Ended entries must still be held. struct lia_list_cmd *cmd = list->cmd; cmd->op = SKIPTO; @@ -428,7 +462,8 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) return handle_skipto(list, cmd->sequence, cmd->arg0.i); } } - al_assert(list->current != -1); + // count is the count before this entry was added. + al_assert(list->current >= 0 || count < start_at); return true; } @@ -480,7 +515,8 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) } } - log_trace("toggle_pause(#%u): pts: %f, pause: %hhu.", entry->id, pts, pause); + log_trace("toggle_pause(#%u): pts: %.4f, pause: %hhu.", entry->id, + (pause == LIANA_PAUSE_PAUSE) ? pts : NAN, pause); struct lia_timing time = { .at = at, @@ -488,6 +524,8 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) .pause = pause }; + list_signal_meta(list, entry, LIANA_META_ENTRY_PAUSED); + struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { if (sequence == list->current && sink->set != sequence) { @@ -497,8 +535,6 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) sink_pause_entry(sink, entry, sequence, &time); } } - - list_signal_meta(list, entry, LIANA_META_ENTRY_PAUSED); } static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos) @@ -542,6 +578,8 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos) .pause = pause }; + list_signal_meta(list, entry, LIANA_META_ENTRY_SEEKED); + struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { if (sequence == list->current && sink->set != sequence) { @@ -551,8 +589,6 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, u64 pos) sink_seek_entry(sink, entry, sequence, &time); } } - - list_signal_meta(list, entry, LIANA_META_ENTRY_SEEKED); } static bool handle_end(struct lia_list *list, u32 id, u32 reset_token) @@ -587,10 +623,10 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_token) return true; #endif - s32 size = (s32)list->entries.count; + s32 count = (s32)list->entries.count; if (sequence == list->current) { s32 next = sequence + 1; - if (next < size) { + if (next < count) { struct lia_list_cmd *cmd = list->cmd; cmd->op = SKIPTO; cmd->sequence = sequence; @@ -637,9 +673,9 @@ static bool handle_reverse(struct lia_list *list) { if (list->current < 0) 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++) { - u32 tail = size - (i + 1); + u32 count = list->entries.count; + for (u32 i = 0; i < count; i++) { + u32 tail = count - (i + 1); if (tail <= i) break; SWAP(al_array_at(list->entries, i), al_array_at(list->entries, tail)); } @@ -647,32 +683,39 @@ static bool handle_reverse(struct lia_list *list) return adjust_for_order_change(list, previous); } +static s32 list_entry_compare(const void *a, const void *b) +{ + struct lia_list_entry *entry1 = *(struct lia_list_entry **)a; + struct lia_list_entry *entry2 = *(struct lia_list_entry **)b; + return nn_numerical_compare(&entry1->brief, &entry2->brief); +} + static bool handle_sort(struct lia_list *list) { if (list->current < 0) 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); + al_array_sort(list->entries, struct lia_list_entry *, list_entry_compare); log_trace("sort()"); return adjust_for_order_change(list, previous); } static bool handle_shuffle(struct lia_list *list) { - u32 size = list->entries.count; - if (list->current < 0 || size < 2) return true; + u32 count = list->entries.count; + if (list->current < 0 || count < 2) return true; 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] */ - if (size == 2) { + if (count == 2) { if (al_random_int(0, 1) == 0) { SWAP(al_array_at(list->entries, 0), al_array_at(list->entries, 1)); } } else { - for (u32 i = 0; i < size - 2; i++) { - u32 j = al_random_int(i, size - 1); + for (u32 i = 0; i < count - 2; i++) { + u32 j = al_random_int(i, count - 1); SWAP(al_array_at(list->entries, i), al_array_at(list->entries, j)); } } @@ -717,12 +760,14 @@ static void run_queue(struct lia_list *list) } else { return; } + } else { + return; } struct lia_list_cmd *cmd = list->cmd; switch (cmd->op) { case ADD_SINK: if (!handle_add_sink(list, cmd->sink)) { - return; + goto retry; } break; case REMOVE_SINK: @@ -730,18 +775,18 @@ static void run_queue(struct lia_list *list) break; case ADD: if (!handle_add(list, cmd->entry)) { - return; + goto retry; } break; // Return on SKIPTO/SKIP: Target entry not loaded. case SKIPTO: if (!handle_skipto(list, cmd->sequence, cmd->arg0.i)) { - return; + goto retry; } break; case SKIP: if (!handle_skip(list, cmd->sequence, cmd->arg0.i)) { - return; + goto retry; } break; case TOGGLE_PAUSE: @@ -753,23 +798,23 @@ static void run_queue(struct lia_list *list) case END: if (!handle_end(list, cmd->arg0.u, cmd->arg1.u)) { // Converted to a SKIP and target entry not loaded. - return; + goto retry; } break; // Return on order change: Converted to SKIPTO and new current not loaded. case REVERSE: if (!handle_reverse(list)) { - return; + goto retry; } break; case SORT: if (!handle_sort(list)) { - return; + goto retry; } break; case SHUFFLE: if (!handle_shuffle(list)) { - return; + goto retry; } break; case UNSET: @@ -782,6 +827,14 @@ static void run_queue(struct lia_list *list) al_free(cmd); list->cmd = NULL; pump_queue(list); + return; +retry: + list->cmd = NULL; + al_array_insert(list->command_queue, 0, cmd); + // We could have swallowed a pump_queue() meant for a different + // command while list->cmd was set. That's ok because it doesn't + // change anything about finishing this command, which will call + // pump_queue() once complete. } void pump_queue(struct lia_list *list) diff --git a/src/liana/list.h b/src/liana/list.h index 580198f..0e4b83c 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -13,6 +13,7 @@ enum { LIANA_SINK_SET = 0, LIANA_SINK_UNSET, + LIANA_SINK_SEQUENCE, LIANA_SINK_BUFFER, LIANA_SINK_BUFFER_AND_QUEUE, LIANA_SINK_PAUSE, @@ -55,6 +56,8 @@ enum { }; struct lia_timing { + // @TODO: Seperate pause and resume tiemstamps. + //struct { u64 p; u64 r; } at; u64 at; u64 pos; u8 pause; @@ -118,6 +121,7 @@ static inline const char *lia_pause_op_name(u8 pause) } void lia_list_init(struct lia_list *list, str *name); +void lia_list_start_at(struct lia_list *list, s32 i); void lia_list_pump(struct lia_list *list); diff --git a/src/liana/list_cmp.h b/src/liana/list_cmp.h deleted file mode 100644 index 6bc7c69..0000000 --- a/src/liana/list_cmp.h +++ /dev/null @@ -1,64 +0,0 @@ -#pragma once - -#include <al/str.h> - -#include "list.h" - -AL_IGNORE_WARNING("-Wunused-function") - -static void camu_db_num_from_path(str *path, s64 *id, s64 *index) -{ - u32 last_slash = al_str_rfind(path, '/'); - if (last_slash == AL_WSTR_NO_POS) { - return; - } - str sub = al_str_substr(path, last_slash + 1, path->length); - - // Skip 2 '_' characters. - if (!(al_str_tok(&sub, '_') && al_str_tok(&sub, '_'))) { - return; - } - - u32 target = al_str_find(&sub, '_'); - if (target == AL_WSTR_NO_POS) { - return; - } - sub = al_str_substr(&sub, 0, target); - - bool error; - s64 num = al_str_to_long(&sub, 10, &error); - if (error) return; - *id = num; - - u32 ext_dot = al_str_rfind(path, '.'); - u32 a_of_media = al_str_rfind(path, 'a'); - if (ext_dot == AL_WSTR_NO_POS || a_of_media == AL_WSTR_NO_POS) { - return; - } - sub = al_str_substr(path, a_of_media + 1, ext_dot); - - num = al_str_to_long(&sub, 10, &error); - if (error) return; - *index = num; -} - -static s32 camu_db_compare(const void *a, const void *b) -{ - struct lia_list_entry *aa = *((struct lia_list_entry **)a); - struct lia_list_entry *bb = *((struct lia_list_entry **)b); - s64 a_id = -1, a_index = -1; - s64 b_id = -1, b_index = -1; - // Assuming brief is the file path. - camu_db_num_from_path(&aa->brief, &a_id, &a_index); - camu_db_num_from_path(&bb->brief, &b_id, &b_index); - if (a_id == b_id) { - if (a_index > b_index) return 1; - else if (a_index < b_index) return -1; - } else { - if (a_id > b_id) return 1; - else if (a_id < b_id) return -1; - } - return 0; -} - -AL_IGNORE_WARNING_END diff --git a/src/liana/server.c b/src/liana/server.c index 791fe0e..8305c61 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -1,6 +1,9 @@ #define AL_LOG_SECTION "liana" +//#define AL_LOG_ENABLE_TRACE #include <al/log.h> +#include "../server/common.h" + #include "server.h" #include "handlers.h" #include "list.h" @@ -29,6 +32,7 @@ static void data_packet_sent_callback(void *userdata, struct nn_packet *packet) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; nn_packet_pool_lock(&conn->pool); + disown_packet(packet); nn_packet_pool_return(&conn->pool, packet); nn_packet_pool_unlock(&conn->pool); } @@ -36,11 +40,12 @@ static void data_packet_sent_callback(void *userdata, struct nn_packet *packet) static void data_packets_sent_callback(void *userdata, struct nn_packet **packets, u32 count) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; + // If in DIRECT_MODE and the client disconnects first, free_connection_stream() has already run. nn_packet_pool_lock(&conn->pool); for (u32 i = 0; i < count; i++) { - if (packets[i]) { - nn_packet_pool_return(&conn->pool, packets[i]); - } + al_assert(packets[i]); + disown_packet(packets[i]); + nn_packet_pool_return(&conn->pool, packets[i]); } nn_packet_pool_unlock(&conn->pool); } @@ -50,11 +55,17 @@ static void packet_pool_callback(void *userdata, struct nn_packet *packet) struct lia_node_connection *conn = (struct lia_node_connection *)userdata; if (!nn_packet_stream_send_packet(conn->stream, packet)) { nn_packet_pool_lock(&conn->pool); + disown_packet(packet); nn_packet_pool_return(&conn->pool, packet); nn_packet_pool_unlock(&conn->pool); } } +//#define SPORADIC_ERROR_PACKET +#ifdef SPORADIC_ERROR_PACKET +#include <al/random.h> +#endif + static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; @@ -65,25 +76,59 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) conn->seek_pos = LIANA_TIMESTAMP_INVALID; } + // Set TCP_CORK with the intention to notify socket that we are likely + // about to write a lot of data and sending partial frames/small chunks + // won't be helpful. + struct nn_socket *sock = &conn->stream->sock; + if (sock->type == NNWT_SOCKET_TCP) { + nn_socket_set_cork(&conn->stream->sock, 1); +#ifdef AL_LOG_ENABLE_TRACE + u32 rcvbuf = nn_socket_get_recv_buf(sock); + u32 sndbuf = nn_socket_get_send_buf(sock); +#endif + nn_socket_set_recv_buf(sock, KB(8)); + nn_socket_set_send_buf(sock, MB(2)); +#ifdef AL_LOG_ENABLE_TRACE + log_trace("rcvbuf: %u -> %u", rcvbuf, nn_socket_get_recv_buf(sock)); + log_trace("sndbuf: %u -> %u", sndbuf, nn_socket_get_send_buf(sock)); +#endif + } + + // A possible throughput optimization here would be to bundle multiple + // AVPackets into a single packet up to a certain size. for (;;) { struct nn_packet *packet = nn_packet_pool_get(&conn->pool); if (!packet) { - return 0; + goto out; } - conn->handler->step(conn->handler); - conn->handler->write_packet(conn->handler, packet); - +#ifdef SPORADIC_ERROR_PACKET + bool error = al_random_int(0, 1000) == 17; + if (error) { + nn_packet_write_u8(packet, LIANA_PACKET_ERROR); + conn->handler->status = CAMU_ERR_ERROR; + } else { +#endif + 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; + if (conn->node->duration > 0 && conn->handler->status == CAMU_ERR_EOF) { + nn_packet_pool_lock(&conn->pool); + disown_packet(packet); + nn_packet_pool_return(&conn->pool, packet); + nn_packet_pool_unlock(&conn->pool); + continue; + } +#endif +#ifdef SPORADIC_ERROR_PACKET } #endif - nn_packet_pool_submit(&conn->pool, packet); + if (!nn_packet_pool_submit(&conn->pool, packet)) { + disown_packet(packet); + nn_packet_pool_return(&conn->pool, packet); + nn_packet_pool_unlock(&conn->pool); + } // Check status after submitting so the EOF packet gets sent. if (conn->handler->status != CAMU_OK) break; @@ -91,13 +136,17 @@ static nn_thread_result NNWT_THREADCALL handler_thread(void *userdata) nn_packet_pool_flush(&conn->pool); +out: + if (sock->type == NNWT_SOCKET_TCP) { + nn_socket_set_cork(sock, 0); + } + return 0; } static void discard_packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) { (void)userdata; - // @TODO: This should invalidate the connection instead of asserting. nn_packet_stream_return_packet(stream, packet); al_assert(false); } @@ -127,7 +176,7 @@ static void free_connection(struct lia_node_connection *conn) cch_entry_return_handle(node->entry, &conn->handle); nn_packet_pool_free(&conn->pool); bool removed = al_array_remove(node->connections, conn); - al_assert(removed); + al_assert(removed && !al_array_contains(node->connections, conn)); al_free(conn); if (should_free_node(node)) { free_node(node); @@ -145,6 +194,13 @@ static void free_connection_stream(struct lia_node_connection *conn) static void disable_connection_and_wait(struct lia_node_connection *conn) { nn_packet_pool_disable(&conn->pool); + nn_packet_pool_lock(&conn->pool); + struct nn_packet *packet; + while ((packet = nn_packet_pool_pop(&conn->pool))) { + disown_packet(packet); + nn_packet_pool_return(&conn->pool, packet); + } + nn_packet_pool_unlock(&conn->pool); cch_handle_disable(&conn->handle); nn_thread_join(&conn->thread); } @@ -188,12 +244,6 @@ static void subscribe_packet_callback(void *userdata, struct nn_packet_stream *s start_connection_handler(conn, mask); } -static void subscribe_packet_sent_callback(void *userdata, struct nn_packet *packet) -{ - (void)userdata; - nn_packet_free(packet); -} - static void subscribe_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; @@ -209,6 +259,8 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet u64 mask = nn_packet_read_u64(packet); u64 seek_pos = nn_packet_read_u64(packet); + nn_packet_stream_return_packet(stream, packet); + al_assert(!conn->ref); conn->ref = true; @@ -225,15 +277,12 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet nn_packet_write_str(rpacket, cch_entry_get_handler(conn->node->entry)); conn->handler->write_info(conn->handler, rpacket); stream->packet_callback = subscribe_packet_callback; - stream->packet_sent_callback = subscribe_packet_sent_callback; stream->connection_closed_callback = subscribe_connection_closed_callback; nn_packet_stream_send_packet(stream, rpacket); } else { start_connection_handler(conn, mask); conn->handler->subscribe(conn->handler, mask); } - - nn_packet_stream_return_packet(stream, packet); } static void connection_closed_callback(void *userdata, struct nn_packet_stream *stream) @@ -244,6 +293,12 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * al_free(stream); } +static void packet_dequeued_callback(void *userdata, struct nn_packet *packet) +{ + (void)userdata; + nn_packet_write_size(packet); +} + static void packet_sent_callback(void *userdata, struct nn_packet *packet) { (void)userdata; @@ -272,21 +327,18 @@ static void signal_callback(void *userdata) nn_thread_join(&conn->thread); nn_signal_stop(&conn->signal); + al_array_remove(node->requests, conn); struct nn_packet_stream *stream = conn->stream; struct nn_packet *packet = conn->packet; conn->packet = NULL; - if (packet) { - al_array_remove(node->requests, conn); - } - 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. + if (!packet) { // Connection was closed before init_thread() finished. al_free(conn); if (should_free_node(node)) { free_node(node); @@ -298,10 +350,14 @@ static void signal_callback(void *userdata) nn_packet_stream_return_packet(stream, packet); demote_and_disconnect_stream(server, stream); al_free(conn); + if (should_free_node(node)) { + free_node(node); + } } else { + al_assert(!should_free_node(node)); conn->id = get_incremental_id(server); al_array_push(node->connections, conn); - nn_packet_pool_init(&conn->pool, 1024, server->loop, packet_pool_callback, conn); + nn_packet_pool_init(&conn->pool, 3072, 2, server->loop, packet_pool_callback, conn); handle_connection(conn, packet); } } @@ -328,22 +384,22 @@ static struct lia_node *get_node_from_id(struct lia_server *server, u32 id) static struct lia_node_connection *get_connection_from_id(struct lia_node *node, u32 id) { - struct lia_node_connection *conn; + struct lia_node_connection *conn, *ret = NULL; al_array_foreach(node->connections, i, conn) { - if (conn->id == id) return conn; + if (conn->id == id) { + al_assert(!ret); + ret = conn; + } } - return NULL; + return ret; } static void pre_init_connection_closed_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_node_connection *conn = (struct lia_node_connection *)userdata; - struct lia_node *node = conn->node; struct nn_packet *packet = conn->packet; - // Checked in signal_callback and will signal to cleanup the connection. - conn->packet = NULL; + conn->packet = NULL; // Request cleanup in signal_callback(). nn_packet_stream_return_packet(stream, packet); - al_array_remove(node->requests, conn); } static void packet_callback(void *userdata, struct nn_packet_stream *stream, struct nn_packet *packet) @@ -357,7 +413,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str u32 connection_id = nn_packet_read_u32(packet); struct lia_node *node = get_node_from_id(server, node_id); - if (!node) { + if (!node || node->closed) { goto err; } @@ -403,7 +459,6 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str return; err: - // Return packet before disconnecting. nn_packet_stream_return_packet(stream, packet); nn_packet_stream_disconnect(stream); } @@ -412,6 +467,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream) { struct lia_server *server = (struct lia_server *)userdata; stream->packet_callback = packet_callback; + stream->packet_dequeued_callback = packet_dequeued_callback; stream->packet_sent_callback = packet_sent_callback; al_array_push(server->dormant_connections, stream); return true; @@ -474,7 +530,7 @@ static void duration_signal_callback(void *userdata) } } -void lia_node_get_duration(struct lia_node *node) +void lia_node_probe_duration(struct lia_node *node) { // It is vital we don't block the loop during handler->init(). nn_signal_init(&node->signal, node->server->loop, duration_signal_callback, node); @@ -492,12 +548,14 @@ void lia_node_close(struct lia_node *node) } else { struct lia_node_connection *conn; al_array_foreach_rev(node->requests, i, conn) { - nn_packet_stream_disconnect(conn->stream); + // We cannot disconnect the stream here because pre_init_connection_closed_callback() + // doesn't remove it from node->requests. + conn->errored = true; } al_array_foreach_rev(node->connections, i, conn) { if (conn->stream) { conn->disconnected = true; - //nn_packet_stream_disconnect(conn->stream); + nn_packet_stream_disconnect(conn->stream); } else { free_connection(conn); } @@ -513,6 +571,24 @@ void lia_server_close(struct lia_server *server) } } +void lia_server_abort(struct lia_server *server) +{ + struct lia_node *node; + al_array_foreach(server->nodes, i, node) { + struct lia_node_connection *conn; + al_array_foreach_rev(node->requests, j, conn) { + // We can't assert(conn->errored) here as we haven't fully waited on the loop. +#ifdef NAUNET_HAS_THREAD_CANCEL + // This could result in a mutex_destroy() on a locked mutex. + nn_thread_cancel(&conn->thread); + nn_signal_send(&conn->signal); +#else + (void)conn; +#endif + } + } +} + void lia_server_force_disconnect_nodes(struct lia_server *server) { struct lia_node *node; diff --git a/src/liana/server.h b/src/liana/server.h index ec6f82e..6c85371 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -62,7 +62,7 @@ struct lia_server { bool lia_server_init(struct lia_server *server, struct nn_event_loop *loop); void lia_server_add_stream(struct lia_server *server, struct nn_packet_stream *stream); struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_entry *entry); -void lia_node_get_duration(struct lia_node *node); +void lia_node_probe_duration(struct lia_node *node); void lia_node_close(struct lia_node *node); void lia_server_close(struct lia_server *server); void lia_server_force_disconnect_nodes(struct lia_server *server); diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 3b4866c..2612f36 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -6,8 +6,24 @@ #include "handlers/handler.h" #include "vcr.h" +#include "client.h" #include "common.h" +// packet_cache_v2 +// - Each track has it's own cond. +// - Tracks get signaled in order of stream index. +// - Cache has to handle case where the next packet in order is +// from a track that has yet to call packet_cache_wait() (Immediate return). +// - Once track has read as much as it can: +// - If it read all avaliable packets, call packet_cache_wait() again and that entire +// section will be consumed. +// - If it was partial, call packet_cache_yield(<number_of_packets>) for that many packets +// to be consumed. +// Considerations: +// - Works after client disconnect. +// - Optimizes seek and reseek. +// - Supports multiple AVPackets per nn_packet. + #define VCR_BUFFER_BUFFERED MB((u64)6) #ifdef CAMU_HUGE_VIDEO_BUFFER #define VCR_BUFFER_GROW_FACTOR ((u64)24) @@ -30,6 +46,12 @@ enum { VCR_TRACK_CLOSED }; +enum { + VCR_NOT_EOF = 0, + VCR_EOF, + VCR_EOF_ERRORED +}; + #define VCR_TRACK_THREADED(track) \ (track->stream->type == CAMU_STREAM_AUDIO || track->stream->type == CAMU_STREAM_VIDEO) @@ -50,6 +72,7 @@ static void signal_callback(void *userdata) vcr->metrics.last_cork_ts = LIANA_TIMESTAMP_INVALID; } } +#endif static void reset_metrics(struct lia_vcr *vcr) { @@ -78,7 +101,7 @@ static void update_metrics(struct lia_vcr *vcr, u64 size) f32 kbps = (frame / 125.f) / (mark / 1000000.f); f32 average_kbps = vcr->metrics.average_kbps; average_kbps = average_kbps == 0.f ? kbps : (average_kbps + kbps) / 2.f; - f32 buffered = atomic_load(u64)(&vcr->count, AL_ATOMIC_RELAXED) / (f32)MB(1); + f32 buffered = atomic_load(u64)(&vcr->size, AL_ATOMIC_RELAXED) / (f32)MB(1); f32 capacity = vcr->mark.buffered / (f32)MB(1); log_info("Receiving packets at %.2fkbps (%.2f/%.2fMB).", average_kbps, buffered, capacity); vcr->metrics.average_kbps = average_kbps; @@ -87,46 +110,45 @@ static void update_metrics(struct lia_vcr *vcr, u64 size) vcr->metrics.current_frame = 0; } } -#endif 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); - atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); + atomic_store(u64)(&vcr->size, 0, AL_ATOMIC_RELAXED); vcr->mark.buffered = VCR_BUFFER_BUFFERED; atomic_store(u64)(&vcr->mark.low, 0, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_UNTOUCHED; vcr->started = false; -#ifndef CAMU_DIRECT_MODE vcr->corked = false; +#ifndef CAMU_DIRECT_MODE nn_signal_init(&vcr->signal, loop, signal_callback, vcr); - reset_metrics(vcr); #else (void)loop; #endif + reset_metrics(vcr); } +#define count_minus_eof(packets, count) (packets[count - 1] ? count : count - 1) + static void return_entire_cache(struct lia_vcr_track *track) { - struct lia_vcr *vcr = track->vcr; u32 count = track->cache.cache.count; - nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), count); - track->cache.cache.count = 0; + if (count > 0) { + struct lia_vcr *vcr = track->vcr; + struct nn_packet **packets = al_array_offset(track->cache.cache, 0); + count = count_minus_eof(packets, count); +#ifdef VCR_BUFFER_WHOLE_FILE + for (u32 i = 0; i < count; i++) { + disown_packet(packets[i]); + } +#endif + nn_packet_stream_return_packets(vcr->data, packets, count); + track->cache.cache.count = 0; + } } -// packet_cache_v2 -// - Each track has it's own cond. -// - Tracks get signaled in order of stream index. -// - Cache has to handle case where the next packet in order is -// from a track that has yet to call packet_cache_wait() (Immediate return). -// - Once track has read as much as it can. -// - If it read all avaliable packets, call packet_cache_wait() again and that entire -// section will be consumed. -// - If it was partial, call packet_cache_yield(<number_of_packets>) for that many packets -// to be consumed. - static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) { nn_thread_set_priority(NNWT_THREAD_SCHED_FIFO, 32); @@ -137,16 +159,16 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) nn_thread_set_name(thread_name); bool corked; - u32 packets, index = 0; + u32 count, index = 0; struct nn_packet *packet = NULL; - while (nn_packet_cache_wait(&track->cache, &packets)) { - al_assert(packets >= index); + while (nn_packet_cache_wait(&track->cache, &count)) { + al_assert(count >= index); corked = false; - for (; index < packets; index++) { + for (; index < count; index++) { packet = nn_packet_cache_at(&track->cache, index); #ifdef VCR_BUFFER_WHOLE_FILE - if (!packet && packets > 2) { // Loop. + if (!packet && count > 2) { // Loop. index = 0; corked = false; break; @@ -156,7 +178,7 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) if (packet) { // Check if we should uncork the packet stream. u32 size = nn_packet_get_size(packet); - u64 buffer = atomic_sub(u64)(&vcr->count, size, AL_ATOMIC_RELAXED); + u64 buffer = atomic_sub(u64)(&vcr->size, size, AL_ATOMIC_RELAXED); #ifndef CAMU_DIRECT_MODE bool buffered = atomic_load(bool)(&track->buffered, AL_ATOMIC_RELAXED); u64 low = atomic_load(u64)(&vcr->mark.low, AL_ATOMIC_RELAXED); @@ -170,11 +192,11 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) nn_mutex_lock(&track->lock); - // NULL packet means flush. - bool success = track->client->handle_packet(track->client, packet); - if (!success) { - // In the case of codec_client, an EOF will be sent on an error. + struct lia_client_handler *handler = track->handler; + bool success = handler->handle_packet(handler, packet); // NULL packet = flush. + if (!success || (!packet && track->eof == VCR_EOF_ERRORED)) { log_error("Error handling packet, exiting track thread."); + handler->callback(handler->userdata, LIANA_CLIENT_ERRORED, handler->stream, NULL); return_entire_cache(track); track->cache.disabled = true; nn_packet_cache_unlock(&track->cache); @@ -199,7 +221,8 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) index++; // Count of packets consumed, also increment for WHOLE_FILE mode. #ifndef VCR_BUFFER_WHOLE_FILE - nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), index); + struct nn_packet **packets = al_array_offset(track->cache.cache, 0); + nn_packet_stream_return_packets(vcr->data, packets, count_minus_eof(packets, index)); al_array_remove_range(track->cache.cache, 0, index); #endif @@ -231,10 +254,13 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata) #endif } else { #ifndef VCR_BUFFER_WHOLE_FILE - al_assert(index == packets); + al_assert(index == count); index = 0; - nn_packet_stream_return_packets(vcr->data, al_array_offset(track->cache.cache, 0), packets); - al_array_remove_range(track->cache.cache, 0, packets); + if (count > 0) { + struct nn_packet **packets = al_array_offset(track->cache.cache, 0); + nn_packet_stream_return_packets(vcr->data, packets, count_minus_eof(packets, count)); + al_array_remove_range(track->cache.cache, 0, count); + } #endif nn_packet_cache_unlock(&track->cache); } @@ -250,10 +276,10 @@ void lia_vcr_start(struct lia_vcr *vcr) #endif struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - al_assert(!track->running); + al_assert(!track->started); if (VCR_TRACK_THREADED(track)) { nn_thread_create(&track->thread, vcr_track_thread, track); - track->running = true; + track->started = true; } } vcr->started = true; @@ -265,10 +291,13 @@ void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track) nn_cond_init(&track->cond); nn_mutex_init(&track->lock); atomic_store(bool)(&track->buffered, !VCR_TRACK_THREADED(track), AL_ATOMIC_RELAXED); - track->running = false; - nn_packet_cache_init(&track->cache, 256); + track->started = false; + if (VCR_TRACK_THREADED(track)) { + nn_packet_cache_init(&track->cache, 256); + track->state = VCR_TRACK_RUNNING; + } + track->eof = VCR_NOT_EOF; al_array_push(vcr->tracks, track); - track->state = VCR_TRACK_RUNNING; } bool lia_vcr_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_stream *stream) @@ -277,8 +306,10 @@ bool lia_vcr_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_strea al_array_foreach(vcr->tracks, i, track) { if (track->stream == stream) { al_array_remove_at(vcr->tracks, i); - nn_packet_cache_free(&track->cache); - track->client->free(&track->client); + if (VCR_TRACK_THREADED(track)) { + nn_packet_cache_free(&track->cache); + } + track->handler->free(&track->handler); al_free(track); return true; } @@ -302,15 +333,17 @@ static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index return NULL; } -#ifndef CAMU_DIRECT_MODE static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) { bool buffered = true; struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { + // The audio/video buffers that a track supplies call vcr_set_buffered(track) + // when they decide that they're buffered. buffered &= atomic_load(bool)(&track->buffered, AL_ATOMIC_RELAXED); } if (buffered) { +#ifndef CAMU_DIRECT_MODE if (vcr->expand == VCR_EXPAND_UNTOUCHED) { vcr->mark.buffered = buffer * VCR_BUFFER_GROW_FACTOR; vcr->expand = VCR_EXPAND_GROWN; @@ -330,22 +363,24 @@ static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) vcr->expand = VCR_EXPAND_COMPLETE; } al_array_foreach(vcr->tracks, i, track) { - nn_packet_cache_flush(&track->cache); + if (VCR_TRACK_THREADED(track)) { + nn_packet_cache_flush(&track->cache); + } } +#else + (void)buffer; +#endif } } -#endif -void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) +bool lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) { struct lia_vcr_track *track; u8 op = nn_packet_read_u8(packet); switch (op) { case LIANA_PACKET_DATA: { u32 size = nn_packet_get_size(packet); -#ifndef CAMU_DIRECT_MODE update_metrics(vcr, size); -#endif s32 index = nn_packet_read_s32(packet); if (!(track = get_track_from_index(vcr, index))) { log_error("Received data from an errored or unknown track (index: %d).", index); @@ -353,46 +388,47 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) } if (VCR_TRACK_THREADED(track)) { u64 buffer = 0; - bool can_send = nn_packet_cache_available(&track->cache); - if (can_send) { - buffer = atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED); + bool send_to_cache = !nn_packet_cache_disabled(&track->cache); + if (send_to_cache) { + buffer = atomic_add(u64)(&vcr->size, size, AL_ATOMIC_RELAXED); nn_packet_cache_send_packet(&track->cache, packet); } nn_packet_cache_unlock(&track->cache); - if (can_send) { + if (send_to_cache) { if (buffer >= vcr->mark.buffered) { -#ifndef CAMU_DIRECT_MODE cork_if_buffered(vcr, buffer); -#endif } - return; // Keep packet. - } - } else { - if (!track->client->handle_packet(track->client, packet)) { - log_warn("Error handling non-buffered packet."); + return true; // Keep packet. } + } else if (!track->handler->handle_packet(track->handler, packet)) { + log_warn("Error handling non-buffered packet."); } break; } case LIANA_PACKET_EOF: case LIANA_PACKET_ERROR: { - // @TODO: Should ERROR be passed down to LIANA_CLIENT_ERRORED? -#ifndef CAMU_DIRECT_MODE + bool error = op == LIANA_PACKET_ERROR; update_metrics(vcr, 0); // Flush. -#endif al_array_foreach(vcr->tracks, i, track) { - if (nn_packet_cache_available(&track->cache)) { - nn_packet_cache_send_packet(&track->cache, NULL); + al_assert(track->eof == VCR_NOT_EOF); + track->eof = error ? VCR_EOF_ERRORED : VCR_EOF; + if (VCR_TRACK_THREADED(track)) { + if (!nn_packet_cache_disabled(&track->cache)) { + nn_packet_cache_send_packet(&track->cache, NULL); + } + nn_packet_cache_unlock(&track->cache); } - nn_packet_cache_unlock(&track->cache); } #ifndef CAMU_DIRECT_MODE + // A corked stream won't close after an unexpected disconnect. This results in better + // behavior for the sink (e.g., an image buffer won't immediately be removed). + nn_packet_stream_cork(vcr->data, true); nn_signal_stop(&vcr->signal); #endif - if (op == LIANA_PACKET_EOF) { - log_info("Received EOF."); - } else if (op == LIANA_PACKET_ERROR) { + if (error) { log_warn("Forcing EOF due to an error packet."); + } else { + log_info("Received EOF."); } break; } @@ -400,7 +436,7 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) log_warn("Erroneous packet."); break; } - nn_packet_stream_return_packet(vcr->data, packet); + return false; } void lia_vcr_set_buffered(struct lia_vcr_track *track) @@ -420,14 +456,18 @@ void lia_vcr_uncork(struct lia_vcr_track *track) { // We need to avoid a race with cork() _and_ flush(). // Possible race with flush() if locking after checking if state != STOPPED: - // - cork() -> flush() -> uncork(). - // - In uncork() we evaluate track->state to be STOPPED then wait on the lock - // being held by flush(). flush() sets the state to CLOSED and signals the - // cond. Now nn_cond_is_waiting() is false at the point uncork() acquires the lock. + // -> cork() -> flush() -> uncork(). + // -> In uncork() we evaluate track->state to be STOPPED then wait on the lock + // being held by flush(). + // -> flush() sets the state to CLOSED and signals the cond. + // -> Now nn_cond_is_waiting() is false at the point uncork() acquires the lock. + // Possible race with vcr_track_thread(): + // -> Track corks while holding lock. + // -> Buffer goes under MARK_LOW, tries to uncork(), evaluates track->state = STOPPED + // then waits on the lock. + // -> handle_packet() errors and sets state to ERRORED. + // -> uncork() aquires the lock and causes an invalid state. nn_mutex_lock(&track->lock); - // If a buffer is running under MARK_LOW, it will be continuously trying to uncork() - // the track. Meaning an erroring handle_packet() could happen at the same time as - // an uncork(). Causing track->state to be ERRORED after we acquire the lock here. if (track->state != VCR_TRACK_STOPPED) { nn_mutex_unlock(&track->lock); return; @@ -438,7 +478,7 @@ void lia_vcr_uncork(struct lia_vcr_track *track) nn_mutex_unlock(&track->lock); } -static void vcr_track_close_internal(struct lia_vcr_track *track) +static void vcr_threaded_track_close(struct lia_vcr_track *track) { struct lia_vcr *vcr = track->vcr; // Calling packet_cache_disable() while holding track->lock can very possibly deadlock. @@ -453,44 +493,45 @@ static void vcr_track_close_internal(struct lia_vcr_track *track) track->state = VCR_TRACK_CLOSED; nn_mutex_unlock(&track->lock); if (vcr->started) { - al_assert(track->running); + al_assert(track->started); nn_thread_join(&track->thread); - track->running = false; + track->started = false; } } -static void vcr_close_all_internal(struct lia_vcr *vcr) +static void vcr_close_all(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { - vcr_track_close_internal(track); + vcr_threaded_track_close(track); return_entire_cache(track); } } vcr->started = false; -#ifndef CAMU_DIRECT_MODE vcr->corked = false; +#ifndef CAMU_DIRECT_MODE nn_signal_stop(&vcr->signal); - reset_metrics(vcr); #endif + reset_metrics(vcr); } void lia_vcr_flush(struct lia_vcr *vcr) { - vcr_close_all_internal(vcr); + vcr_close_all(vcr); struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { if (VCR_TRACK_THREADED(track)) { atomic_store(bool)(&track->buffered, false, AL_ATOMIC_RELAXED); - track->client->flush(track->client); + track->handler->flush(track->handler); nn_packet_cache_enable(&track->cache); track->state = VCR_TRACK_RUNNING; } else { - track->client->flush(track->client); + track->handler->flush(track->handler); } + track->eof = VCR_NOT_EOF; } - atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); + atomic_store(u64)(&vcr->size, 0, AL_ATOMIC_RELAXED); if (vcr->expand == VCR_EXPAND_COMPLETE) { atomic_store(u64)(&vcr->mark.low, 0, AL_ATOMIC_RELAXED); vcr->expand = VCR_EXPAND_GROWN; @@ -499,15 +540,17 @@ void lia_vcr_flush(struct lia_vcr *vcr) void lia_vcr_close_all(struct lia_vcr *vcr) { - vcr_close_all_internal(vcr); + vcr_close_all(vcr); } void lia_vcr_free(struct lia_vcr *vcr) { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - nn_packet_cache_free(&track->cache); - track->client->free(&track->client); + if (VCR_TRACK_THREADED(track)) { + nn_packet_cache_free(&track->cache); + } + track->handler->free(&track->handler); al_free(track); } al_array_free(vcr->tracks); diff --git a/src/liana/vcr.h b/src/liana/vcr.h index 5b8b9df..42d6ff9 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -12,11 +12,12 @@ struct lia_vcr_track { struct camu_codec_stream *stream; - struct lia_client_handler *client; struct nn_packet_cache cache; + struct lia_client_handler *handler; u32 state; atomic(bool) buffered; - bool running; + u8 eof; + bool started; struct nn_cond cond; struct nn_mutex lock; struct nn_thread thread; @@ -27,15 +28,15 @@ struct lia_vcr { struct nn_packet_stream *data; u16 node_id; array(struct lia_vcr_track *) tracks; - atomic(u64) count; + atomic(u64) size; struct { u64 buffered; atomic(u64) low; } mark; u8 expand; bool started; -#ifndef CAMU_DIRECT_MODE bool corked; +#ifndef CAMU_DIRECT_MODE struct nn_signal signal; #endif struct { @@ -52,7 +53,7 @@ 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_remove_track_by_stream(struct lia_vcr *vcr, struct camu_codec_stream *stream); bool lia_vcr_is_empty(struct lia_vcr *vcr); -void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet); +bool lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet); void lia_vcr_set_buffered(struct lia_vcr_track *track); void lia_vcr_cork(struct lia_vcr_track *track); void lia_vcr_uncork(struct lia_vcr_track *track); |