diff options
| author | 2025-02-19 13:40:38 -0500 | |
|---|---|---|
| committer | 2025-02-19 13:40:38 -0500 | |
| commit | 2bee71a7e032c0972418e324bb1d7e6b02330b18 (patch) | |
| tree | fa5819e9d9efe75cbc6f50d462dd0eb6bd19b6a2 /src/liana | |
| parent | b36f022defd8d4ec5a8c29578bb583bea05dfbb6 (diff) | |
| download | camu-2bee71a7e032c0972418e324bb1d7e6b02330b18.tar.gz camu-2bee71a7e032c0972418e324bb1d7e6b02330b18.tar.bz2 camu-2bee71a7e032c0972418e324bb1d7e6b02330b18.zip | |
Server resource unload, many tweaks and fixes
- Initial liana client preferences.
- Hook up libplacebo dx11 backend.
- Make usage of FFmpeg hardware decoding api make some sense.
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 139 | ||||
| -rw-r--r-- | src/liana/client.h | 12 | ||||
| -rw-r--r-- | src/liana/handlers/codec_client.c | 14 | ||||
| -rw-r--r-- | src/liana/list.c | 422 | ||||
| -rw-r--r-- | src/liana/list.h | 15 | ||||
| -rw-r--r-- | src/liana/server.c | 94 | ||||
| -rw-r--r-- | src/liana/server.h | 1 | ||||
| -rw-r--r-- | src/liana/vcr.c | 29 | ||||
| -rw-r--r-- | src/liana/vcr.h | 3 |
9 files changed, 425 insertions, 304 deletions
diff --git a/src/liana/client.c b/src/liana/client.c index 44cba88..d83ef48 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -28,26 +28,20 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe str liana; nn_packet_read_str(packet, &liana); client->duration = nn_packet_read_u64(packet); - // @TODO: This should be made into 2 steps. - // 1. Collect all streams into an array - // 2. Perform selection based on prefrences. - bool have_audio = false; - bool have_video = false; - bool have_subs = false; + u32 count = nn_packet_read_u32(packet); for (u32 i = 0; i < count; i++) { + // Very important zero-initialization. + struct camu_codec_stream stream = { 0 }; u8 mode = nn_packet_read_u8(packet); u8 type = nn_packet_read_u8(packet); u64 duration = nn_packet_read_u64(packet); s32 index = nn_packet_read_s32(packet); al_assert(index < 32); - struct lia_vcr_track *track = NULL; switch (mode) { case CAMU_NORMAL: { - client->mask |= 1 << index; - track = al_alloc_object(struct lia_vcr_track); if (type == CAMU_STREAM_AUDIO) { - struct camu_audio_format *fmt = &track->stream.audio.fmt; + struct camu_audio_format *fmt = &stream.audio.fmt; fmt->format = nn_packet_read_s32(packet); fmt->sample_rate = nn_packet_read_s32(packet); fmt->channel_count = nn_packet_read_s32(packet); @@ -55,7 +49,7 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe av_channel_layout_default(&fmt->channel_layout, fmt->channel_count); #endif } else if (type == CAMU_STREAM_VIDEO) { - struct camu_video_format *fmt = &track->stream.video.fmt; + struct camu_video_format *fmt = &stream.video.fmt; fmt->width = nn_packet_read_u32(packet); fmt->height = nn_packet_read_u32(packet); fmt->format = nn_packet_read_s32(packet); @@ -67,78 +61,81 @@ static void parse_info_packet(struct lia_client *client, struct nn_packet *packe enum AVCodecID codec_id = nn_packet_read_av_codec_id(packet); const AVCodec *codec = avcodec_find_decoder(codec_id); AVFormatContext *format_context = avformat_alloc_context(); - AVStream *stream = nn_packet_read_av_stream(format_context, codec, packet); + stream.av.stream = nn_packet_read_av_stream(format_context, codec, packet); switch (type) { - case CAMU_STREAM_AUDIO: - if (have_audio) { - goto skip; - } - client->mask |= 1 << index; - have_audio = true; - break; - case CAMU_STREAM_VIDEO: - if (have_video) { - goto skip; - } - client->mask |= 1 << index; - have_video = true; - break; - case CAMU_STREAM_SUBTITLE: - if (have_subs || codec_id != AV_CODEC_ID_ASS) { - goto skip; - } - client->mask |= 1 << index; - have_subs = true; - break; - case CAMU_STREAM_ATTACHMENT: { - struct camu_codec_stream attachment; - attachment.type = CAMU_STREAM_ATTACHMENT; - attachment.av.stream = stream; - client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, &attachment, track); + case CAMU_STREAM_ATTACHMENT: // Assume all the data we need is in the AVStream object. - // fallthrough + stream.type = CAMU_STREAM_ATTACHMENT; + client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, &stream, NULL); + avformat_free_context(format_context); + continue; } - default: - goto skip; - } - track = al_alloc_object(struct lia_vcr_track); - track->stream.av.format_context = format_context; - track->stream.av.stream = stream; + stream.av.format_context = format_context; if (type == CAMU_STREAM_AUDIO) { - struct camu_audio_format *fmt = &track->stream.audio.fmt; - fmt->format = stream->codecpar->format; - fmt->sample_rate = stream->codecpar->sample_rate; - av_channel_layout_copy(&fmt->channel_layout, &stream->codecpar->ch_layout); - fmt->channel_count = stream->codecpar->ch_layout.nb_channels; + struct camu_audio_format *fmt = &stream.audio.fmt; + AVCodecParameters *codecpar = stream.av.stream->codecpar; + fmt->format = codecpar->format; + fmt->sample_rate = codecpar->sample_rate; + av_channel_layout_copy(&fmt->channel_layout, &codecpar->ch_layout); + fmt->channel_count = codecpar->ch_layout.nb_channels; } break; -skip: - avformat_free_context(format_context); - continue; } #endif } - al_assert(track); - track->index = index; - track->stream.mode = mode; - track->stream.type = type; - track->stream.duration = duration; + stream.mode = mode; + stream.type = type; + stream.duration = duration; + stream.index = index; + al_array_push(client->streams, stream); + } + + bool have_audio = false; + bool have_video = false; + bool have_subs = false; + + struct camu_codec_stream *stream; + al_array_foreach_ptr(client->streams, i, stream) { + switch (stream->type) { + case CAMU_STREAM_AUDIO: + if (have_audio || !(client->prefs.enabled_mask & CAMU_MASK_AUDIO)) { + continue; + } + have_audio = true; + break; + case CAMU_STREAM_VIDEO: + if (have_video || !(client->prefs.enabled_mask & CAMU_MASK_VIDEO)) { + continue; + } + have_video = true; + break; + case CAMU_STREAM_SUBTITLE: + if (have_subs || !(client->prefs.enabled_mask & CAMU_MASK_SUBTITLE)) { + continue; + } + have_subs = true; + break; + } + + client->mask |= 1 << stream->index; + + struct lia_vcr_track *track = al_alloc_object(struct lia_vcr_track); + track->stream = stream; track->client = lia_handler_by_name(&liana)->create_client_handler(); track->client->callback = client->callback; track->client->userdata = client->userdata; - if (!track->client->init(track->client, client->renderer, &track->stream)) { + + if (!track->client->init(track->client, client->renderer, track->stream)) { track->client->free(&track->client); -#ifdef CAMU_HAVE_FFMPEG - if (track->stream.mode == CAMU_FFMPEG_COMPAT) { - avformat_free_context(track->stream.av.format_context); - } -#endif al_free(track); continue; } - client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, track->client->stream, track); + + client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, stream, track); + lia_vcr_add_track(&client->vcr, track); } + client->callback(client->userdata, LIANA_CLIENT_CONFIGURE_COMPLETE, NULL, NULL); } @@ -229,13 +226,14 @@ static void connection_closed_callback(void *userdata, struct nn_packet_stream * } } -void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type, - str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer) +void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, + u8 type, str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer) { client->loop = loop; client->node_id = node_id; client->pos = pos; client->mask = 0; + al_array_init(client->streams); client->reconnect = RECONNECT_NONE; lia_vcr_init(&client->vcr, client->loop, &client->data); al_str_clone(&client->addr, addr); @@ -281,5 +279,12 @@ void lia_client_free(struct lia_client *client) { lia_vcr_free(&client->vcr); nn_packet_stream_free(&client->data); + struct camu_codec_stream *stream; + al_array_foreach_ptr(client->streams, i, stream) { + if (stream->mode == CAMU_FFMPEG_COMPAT) { + avformat_free_context(stream->av.format_context); + } + } + al_array_free(client->streams); al_str_free(&client->addr); } diff --git a/src/liana/client.h b/src/liana/client.h index 5d5fc90..9e49d0e 100644 --- a/src/liana/client.h +++ b/src/liana/client.h @@ -18,9 +18,17 @@ enum { LIANA_CLIENT_CLOSED }; +struct lia_prefs { + u8 enabled_mask; + s8 audio_lang; + s8 subtitle_lang; +}; + struct lia_client { struct nn_event_loop *loop; u32 node_id; + struct lia_prefs prefs; + array(struct camu_codec_stream) streams; u32 mask; u64 pos; u64 at; @@ -36,8 +44,8 @@ struct lia_client { void *userdata; }; -void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, u8 type, - str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer); +void lia_client_connect(struct lia_client *client, struct nn_event_loop *loop, + u8 type, str *addr, u16 port, u32 node_id, u64 pos, struct camu_renderer *renderer); void lia_client_seek(struct lia_client *client, u64 pos, u64 at); void lia_client_reseek(struct lia_client *client); void lia_client_disconnect(struct lia_client *client); diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index 36e475e..8a86799 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -43,6 +43,13 @@ static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt) s32 ret = codec->dec->push_av_packet(codec->dec, pkt); return ret == CAMU_OK; } + +static void passthrough_subtitle(struct lia_codec_client *codec, AVPacket *pkt) +{ + struct camu_codec_packet packet = { .av.pkt = pkt }; + struct camu_codec_stream *stream = codec->handler.stream; + codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, stream, &packet); +} #endif static bool push_packet(struct lia_codec_client *codec, struct nn_buffer *buffer) @@ -109,15 +116,16 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc nn_packet_read_av_packet(packet, pkt); packet->opaque = pkt; } - struct camu_codec_stream *stream = codec->handler.stream; if (codec->dec) { success = push_av_packet(codec, pkt); } else { - codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, stream, pkt); + passthrough_subtitle(codec, pkt); success = true; } #ifdef VCR_BUFFER_WHOLE_FILE - pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, stream->av.stream->time_base); + struct camu_codec_stream *stream = codec->handler.stream; + AVRational time_base = stream->av.stream->time; + pkt->pts += av_rescale_q(stream->duration, AV_TIME_BASE_Q, time_base); #else av_packet_unref(pkt); av_packet_free(&pkt); diff --git a/src/liana/list.c b/src/liana/list.c index 04ed879..ca8aaa5 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -6,28 +6,7 @@ #include "list.h" #include "list_cmp.h" -/* -static void buffer_ahead(struct lia_list *list) -{ - s32 size = (s32)list->entries.count; - if (list->current >= 0 && list->current + 1 < size) { - s32 ahead = list->current + 1; - for (s32 i = ahead; i < MIN(ahead + LIANA_BUFFER_AHEAD, size); i++) { - struct lia_list_entry *entry = al_array_at(list->entries, i); - struct lia_timing time = { - .at = LIANA_TIMESTAMP_INVALID, - .seek_pos = entry->offset, - .pause = LIANA_PAUSE_NONE, - .ended = false - }; - struct lia_list_sink *sink; - al_array_foreach(list->sinks, i, sink) { - sink->callback(sink->userdata, LIANA_SINK_BUFFER, entry, i, &time); - } - } - } -} -*/ +//#define LIANA_LIST_TRACE enum { ADD_SINK = 0, @@ -60,12 +39,10 @@ void lia_list_init(struct lia_list *list, str *name) list->increment = 0; al_array_init(list->entries); al_array_init(list->sinks); - al_array_init(list->queue); - list->cmd = NULL; + al_array_init(list->command_queue); + list->active_cmd = NULL; } -static void pump_queue(struct lia_list *list); - static bool assume_ended(struct lia_list_entry *entry, u64 at) { if (entry->duration == LIANA_TIMESTAMP_INVALID) return false; @@ -82,6 +59,13 @@ static bool assume_ended(struct lia_list_entry *entry, u64 at) return false; } +#define META_OPAQUE(op) ((u8[]){ op }) + +static void signal_meta(struct lia_list *list, struct lia_list_entry *entry, u8 meta) +{ + list->callback(list->userdata, LIANA_LIST_META, entry, META_OPAQUE(meta)); +} + static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, bool *error) { u8 status; @@ -96,65 +80,85 @@ static bool entry_load_and_get_duration(struct lia_list *list, struct lia_list_e } } *error = true; - list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_ENTRY_ERRORED }); + signal_meta(list, entry, LIANA_META_ENTRY_ERRORED); return false; } *error = false; if (status == LIANA_ENTRY_LOADED) { - list->callback(list->userdata, LIANA_GET_DURATION, entry, &entry->duration); + list->callback(list->userdata, LIANA_GET_ENTRY_DURATION, entry, &entry->duration); return true; } return false; } -static void set_queued(struct lia_list *list) +static void entry_unload(struct lia_list *list, struct lia_list_entry *entry) { - if (list->queued < 0) - return; - - struct lia_list_entry *queued = al_array_at(list->entries, list->queued); + list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL); +} - bool error; - struct lia_list_entry *current = al_array_at(list->entries, list->current); - if (!entry_load_and_get_duration(list, current, list->current, &error)) - return; +static void entry_ref(struct lia_list *list, struct lia_list_entry *entry) +{ + list->callback(list->userdata, LIANA_REF_ENTRY, entry, NULL); +} - queued->start = current->start + (current->duration - current->offset); - u8 pause = queued->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE; - struct lia_timing time = { - .at = queued->start, - .seek_pos = queued->offset, - .pause = pause - }; +static void entry_unref(struct lia_list *list, struct lia_list_entry *entry) +{ + list->callback(list->userdata, LIANA_UNREF_ENTRY, entry, NULL); +} - struct lia_list_sink *sink; - al_array_foreach(list->sinks, i, sink) { - if (sink->queued != list->queued) { - sink->queued = list->queued; - sink->callback(sink->userdata, LIANA_SINK_BUFFER_AND_QUEUE, queued, list->queued, &time); +static void unref_all_entries(struct lia_list *list) +{ + struct lia_list_entry *entry; + al_array_foreach(list->entries, i, entry) { + if (list->current >= 0 && i != (u32)list->current) { + entry_unref(list, entry); } } } -static void evaluate_queued(struct lia_list *list) +static void sink_set(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time) { - s32 size = (s32)list->entries.count; - s32 next = list->current + 1; - if (next >= size || next == list->queued) - return; + sink->set = sequence; + sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, time); +} - bool error; - struct lia_list_entry *queued = al_array_at(list->entries, next); - if (!entry_load_and_get_duration(list, queued, next, &error)) - return; +static void sink_seek(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); +} - list->queued = next; +static void sink_toggle_pause(struct lia_list_sink *sink, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time) +{ + sink->callback(sink->userdata, LIANA_SINK_PAUSE, entry, sequence, time); +} + +static void sink_unset(struct lia_list_sink *sink) +{ + sink->callback(sink->userdata, LIANA_SINK_UNSET, NULL, -1, NULL); } +static void set_current(struct lia_list *list, struct lia_list_entry *entry, s32 sequence, struct lia_timing *time) +{ + entry_ref(list, entry); + list->current = sequence; + list->idle = false; + // @TODO: This is really wrong. Switching back and forth between 2 entries + // will unload everything else. + unref_all_entries(list); + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink_set(sink, entry, sequence, time); + } + signal_meta(list, entry, LIANA_META_CURRENT_CHANGED); +} + +static void pump_queue(struct lia_list *list); + static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) { struct lia_list_entry *current; - if (list->current >= 0 && !list->idle && !(current = al_array_at(list->entries, list->current))->ended) { + if (list->current >= 0 && !list->idle) { + current = al_array_at(list->entries, list->current); bool error; if (!entry_load_and_get_duration(list, current, list->current, &error)) { if (error) pump_queue(list); @@ -185,12 +189,8 @@ static bool handle_add_sink(struct lia_list *list, struct lia_list_sink *sink) .pause = pause, .ended = ended }; - sink->set = list->current; - sink->callback(sink->userdata, LIANA_SINK_SET, current, list->current, &time); - } else { - sink->set = -1; + sink_set(sink, current, list->current, &time); } - sink->queued = -1; al_array_push(list->sinks, sink); return true; } @@ -209,33 +209,27 @@ 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. + // The list being idle doesn't mean list->current or sink->set aren't set. if (list->idle) { bool error; if (!entry_load_and_get_duration(list, entry, -1, &error)) { return error; } - list->current++; - list->idle = false; entry->start = nn_get_timestamp() + LIANA_BASE_DELAY; + } else { + entry->start = LIANA_TIMESTAMP_INVALID; + } + al_array_push(list->entries, entry); + signal_meta(list, entry, LIANA_META_ADDED_ENTRY); + if (list->idle) { struct lia_timing time = { .at = entry->start, .seek_pos = entry->offset, .pause = LIANA_PAUSE_RESUME, .ended = false }; - struct lia_list_sink *sink; - al_array_foreach(list->sinks, i, sink) { - sink->set = list->current; - sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time); - } - list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_ADDED_ENTRY }); - list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_CURRENT_CHANGED }); - } else { - entry->start = LIANA_TIMESTAMP_INVALID; - list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_ADDED_ENTRY }); + set_current(list, entry, list->current + 1, &time); } - al_array_push(list->entries, entry); return true; } @@ -253,21 +247,35 @@ static void unset_all(struct lia_list *list) static void handle_unset(struct lia_list *list) { unset_all(list); - list->current = list->entries.count - 1; - list->idle = true; struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { - sink->callback(sink->userdata, LIANA_SINK_UNSET, NULL, -1, NULL); + sink_unset(sink); } + list->current = list->entries.count - 1; + list->idle = true; } 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) return NULL; + if (sequence < 0 || sequence >= size) { + return NULL; + } return al_array_at(list->entries, sequence); } +static s32 get_sequence_from_entry_id(struct lia_list *list, u32 id) +{ + // Not returning the entry pointer here seems wasteful but it's a meaningful simplification. + struct lia_list_entry *entry; + al_array_foreach(list->entries, i, entry) { + if (entry->id == id) { + return (s32)i; + } + } + return -1; +} + static struct lia_list_entry *get_entry_from_id(struct lia_list *list, u32 id, s32 *sequence) { struct lia_list_entry *entry; @@ -284,18 +292,18 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) { if (index == list->current) return true; if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; - if (sequence < 0) return true; + if (sequence < 0) return true; // list->current = -1 if (sequence != list->current) { // Skipping from an entry other than current is not handled and will cause very - // confusing errors. On top of likely resulting in unexpected behavior. + // confusing errors. To handle it wouldn't make sense anyway because the outcome + // would likely be unexpected to the user. return true; } struct lia_list_entry *current = get_entry_from_sequence(list, sequence); struct lia_list_entry *target = get_entry_from_sequence(list, index); - al_assert(current && !current->held); + al_assert(current && !current->held && current != target); if (!target) return true; - al_assert(current != target); bool error; if (!entry_load_and_get_duration(list, target, index, &error)) { // index might point to a different entry after an error. @@ -320,7 +328,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) // An ended entry may never have been paused, but a non-ended entry that wasn't set // cannot be unpaused. Checking assume_ended(target) should be safe here as long // as it can't go from true to false (consideration for seek?). - //al_assert(assume_ended(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID); + //al_assert(target->paused_at != LIANA_TIMESTAMP_INVALID); } // These are not equivalent to current/target->ended. @@ -350,6 +358,9 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) // Pause current and set the held flag indicating it should be resumed // if it becomes the target of a skip. + // + // @TODO: Why not hold an ended entry? The idea of ended entries being + // unpaused but never held seems like a over-complication. if (!current_ended && current->paused_at == LIANA_TIMESTAMP_INVALID) { current->paused_at = at; if (current->paused_at < current->start) { @@ -360,11 +371,10 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) current->held = true; } - al_log_debug("list", "skipto [#%u-#%u]: pause: %hhu, held: %s, current_ended: %s, target_ended: %s.", +#ifdef LIANA_LIST_TRACE + al_log_info("list", "skipto [#%u-#%u]: pause: %hhu, held: %s, current_ended: %s, target_ended: %s.", current->id, target->id, pause, BOOLSTR(current->held), BOOLSTR(current_ended), BOOLSTR(target_ended)); - - list->current = index; - list->idle = false; +#endif struct lia_timing time = { .at = at, @@ -373,13 +383,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) .ended = target_ended }; - struct lia_list_sink *sink; - al_array_foreach(list->sinks, i, sink) { - sink->set = index; - sink->callback(sink->userdata, LIANA_SINK_SET, target, index, &time); - } - - list->callback(list->userdata, LIANA_LIST_META, target, (u8[]){ LIANA_META_CURRENT_CHANGED }); + set_current(list, target, index, &time); return true; } @@ -387,7 +391,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) static bool handle_skip(struct lia_list *list, s32 sequence, s32 n) { if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; - struct lia_list_cmd *cmd = list->cmd; + struct lia_list_cmd *cmd = list->active_cmd; cmd->op = SKIPTO; cmd->sequence = sequence; cmd->arg0.i = sequence + n; @@ -398,11 +402,12 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) { // pts should be treated as a hint. if (sequence == LIANA_SEQUENCE_ANY) sequence = list->current; - if (sequence < 0) return; + if (sequence < 0) return; // list->current = -1 if (sequence != list->current) { // A non-current entry should always be paused. return; } + struct lia_list_entry *entry = get_entry_from_sequence(list, sequence); al_assert(entry && !entry->held); @@ -430,7 +435,11 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) break; } - al_log_debug("list", "toggle_pause [#%u]: pts: %f, pause: %hhu.", entry->id, pts, pause); +#ifdef LIANA_LIST_TRACE + al_log_info("list", "toggle_pause [#%u]: pts: %f, pause: %hhu.", entry->id, pts, pause); +#else + (void)pts; +#endif struct lia_timing time = { .at = at, @@ -441,7 +450,7 @@ static void handle_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { - sink->callback(sink->userdata, LIANA_SINK_PAUSE, entry, sequence, &time); + sink_toggle_pause(sink, entry, sequence, &time); } } @@ -450,14 +459,15 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent struct lia_list_entry *entry; if (sequence == LIANA_SEQUENCE_ANY) { sequence = list->current; - if (sequence < 0) return; - entry = get_entry_from_sequence(list, sequence); - al_assert(entry); } else { - entry = get_entry_from_id(list, id, &sequence); - if (!entry) return; + sequence = get_sequence_from_entry_id(list, id); } + if (sequence < 0) return; + + entry = get_entry_from_sequence(list, sequence); + al_assert(entry); + if (entry->duration == LIANA_TIMESTAMP_INVALID) { al_log_warn("list", "Skipping seek on entry with no duration."); return; @@ -477,7 +487,9 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent list->idle = false; - al_log_debug("list", "seek [#%u]: pos: %f.", entry->id, pos / 1000000.0); +#ifdef LIANA_LIST_TRACE + al_log_info("list", "seek [#%u]: pos: %f.", entry->id, pos / 1000000.0); +#endif struct lia_timing time = { .at = at, @@ -489,22 +501,19 @@ static void handle_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { if (sequence == list->current && sink->set != sequence) { - sink->set = sequence; - sink->callback(sink->userdata, LIANA_SINK_SET, entry, sequence, &time); + sink_set(sink, entry, sequence, &time); } - sink->callback(sink->userdata, LIANA_SINK_SEEK, entry, sequence, &time); + sink_seek(sink, entry, sequence, &time); } - list->callback(list->userdata, LIANA_LIST_META, entry, (u8[]){ LIANA_META_ENTRY_SEEKED }); + signal_meta(list, entry, LIANA_META_ENTRY_SEEKED); } static bool handle_end(struct lia_list *list, u32 id, u32 reset_id) { - s32 sequence; - struct lia_list_entry *entry = get_entry_from_id(list, id, &sequence); - if (!entry) { - return true; - } + s32 sequence = get_sequence_from_entry_id(list, id); + if (sequence < 0) return true; + struct lia_list_entry *entry = get_entry_from_sequence(list, sequence); if (reset_id != entry->reset_id) { al_log_warn("list", "Got end() with out of order or incorrect reset id, ignoring."); @@ -516,7 +525,10 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id) return true; } - al_log_debug("list", "end [#%u].", entry->id); +#ifdef LIANA_LIST_TRACE + al_log_info("list", "end [#%u].", entry->id); +#endif + entry->ended = true; entry->offset = entry->duration; @@ -534,17 +546,17 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id) if (sequence == list->current) { s32 next = sequence + 1; if (list->queued >= 0) { - list->current = list->queued; - list->queued = -1; + struct lia_list_entry *queued = al_array_at(list->entries, list->queued); struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { sink->queued = -1; - sink->set = list->current; + sink->set = list->queued; } - struct lia_list_entry *current = al_array_at(list->entries, list->current); - list->callback(list->userdata, LIANA_LIST_META, current, (u8[]){ LIANA_META_CURRENT_CHANGED }); + list->current = list->queued; + list->queued = -1; + signal_meta(list, queued, LIANA_META_CURRENT_CHANGED); } else if (next < size) { - struct lia_list_cmd *cmd = list->cmd; + struct lia_list_cmd *cmd = list->active_cmd; cmd->op = SKIPTO; cmd->sequence = sequence; cmd->arg0.i = next; @@ -561,7 +573,7 @@ static bool handle_end(struct lia_list *list, u32 id, u32 reset_id) 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_cmd *cmd = list->active_cmd; struct lia_list_entry *entry; al_array_foreach(list->entries, i, entry) { if (entry->opaque == previous->opaque) { @@ -582,11 +594,10 @@ static bool adjust_current(struct lia_list *list, struct lia_list_entry *previou static bool handle_reverse(struct lia_list *list) { - if (list->current == -1) { - return true; - } + 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); @@ -599,11 +610,10 @@ static bool handle_reverse(struct lia_list *list) static bool handle_sort(struct lia_list *list) { - if (list->current == -1) { - return true; - } + 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); return adjust_current(list, previous); @@ -611,14 +621,13 @@ static bool handle_sort(struct lia_list *list) static bool handle_shuffle(struct lia_list *list) { - if (list->current == -1) { - return true; - } + if (list->current < 0) return true; u32 size = list->entries.count; 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 @@ -634,27 +643,14 @@ static bool handle_shuffle(struct lia_list *list) return adjust_current(list, previous); } -/* -static void handle_clear(struct lia_list *list) -{ - unset_all(list); - list->idle = true; - struct lia_list_entry *entry; - al_array_foreach(list->entries, i, entry) { - al_wstr_free(&entry->name); - al_free(entry); - } - list->entries.count = 0; -} -*/ - static void run_queue(struct lia_list *list) { - if (!list->cmd) { - if (!list->queue.count) return; - al_array_pop_at(list->queue, 0, list->cmd); + // @TODO: What does current = -1/currentless really mean. + if (!list->active_cmd) { + if (!list->command_queue.count) return; + al_array_pop_at(list->command_queue, 0, list->active_cmd); } - struct lia_list_cmd *cmd = list->cmd; + struct lia_list_cmd *cmd = list->active_cmd; switch (cmd->op) { case ADD_SINK: if (!handle_add_sink(list, cmd->sink)) { @@ -680,7 +676,7 @@ static void run_queue(struct lia_list *list) break; case SKIP: if (!handle_skip(list, cmd->sequence, cmd->arg0.i)) { - // Converted to skipto and entry not loaded. + // Converted to skipto and target entry not loaded. return; } break; @@ -712,21 +708,16 @@ static void run_queue(struct lia_list *list) } break; case CLEAR: - /* - handle_clear(list); - */ break; } al_free(cmd); - list->cmd = NULL; + list->active_cmd = NULL; pump_queue(list); } void pump_queue(struct lia_list *list) { run_queue(list); - //evaluate_queued(list); - //set_queued(list); } void lia_list_pump(struct lia_list *list) @@ -737,19 +728,21 @@ void lia_list_pump(struct lia_list *list) void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struct lia_list_entry *, s32, struct lia_timing *), void *userdata) { struct lia_list_sink *sink = al_alloc_object(struct lia_list_sink); + sink->set = -1; + sink->queued = -1; sink->callback = callback; sink->userdata = userdata; struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = ADD_SINK; cmd->sink = sink; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } void lia_list_remove_sink(struct lia_list *list, void *userdata) { - // Don't queue remove sink because we can't let any currently queued commands - // touch this sink. + // Don't queue remove sink because we can't let any currently queued + // commands touch this sink. handle_remove_sink(list, userdata); } @@ -770,7 +763,7 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, wstr *name) struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = ADD; cmd->entry = entry; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -778,7 +771,7 @@ void lia_list_unset(struct lia_list *list) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = UNSET; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -788,7 +781,7 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) cmd->op = SKIPTO; cmd->sequence = sequence; cmd->arg0.i = index; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -798,7 +791,7 @@ void lia_list_skip(struct lia_list *list, s32 sequence, s32 n) cmd->op = SKIP; cmd->sequence = sequence; cmd->arg0.i = n; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -808,7 +801,7 @@ void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) cmd->op = TOGGLE_PAUSE; cmd->sequence = sequence; cmd->argf = pts; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -819,7 +812,7 @@ void lia_list_seek(struct lia_list *list, s32 sequence, u32 id, f64 percent) cmd->sequence = sequence; cmd->arg0.u = id; cmd->argf = percent; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -829,7 +822,7 @@ void lia_list_end(struct lia_list *list, u32 id, u32 reset_id) cmd->op = END; cmd->arg0.u = id; cmd->arg1.u = reset_id; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -837,7 +830,7 @@ void lia_list_reverse(struct lia_list *list) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = REVERSE; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -845,7 +838,7 @@ void lia_list_sort(struct lia_list *list) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SORT; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } @@ -853,35 +846,37 @@ void lia_list_shuffle(struct lia_list *list) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SHUFFLE; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } +/* void lia_list_clear(struct lia_list *list) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = CLEAR; - al_array_push(list->queue, cmd); + al_array_push(list->command_queue, cmd); pump_queue(list); } +*/ void lia_list_close(struct lia_list *list) { // @TODO: Consider sinks being in use. Wait for list->sinks to be empty? struct lia_list_entry *entry; al_array_foreach(list->entries, i, entry) { - list->callback(list->userdata, LIANA_UNLOAD_ENTRY, entry, NULL); + entry_unload(list, entry); } } void lia_list_free(struct lia_list *list) { struct lia_list_cmd *cmd; - al_array_foreach(list->queue, i, cmd) { + al_array_foreach(list->command_queue, i, cmd) { al_free(cmd); } - al_array_free(list->queue); - if (list->cmd) al_free(list->cmd); + al_array_free(list->command_queue); + if (list->active_cmd) al_free(list->active_cmd); struct lia_list_entry *entry; al_array_foreach(list->entries, i, entry) { @@ -898,3 +893,74 @@ void lia_list_free(struct lia_list *list) al_str_free(&list->name); } + +/* +static void buffer_ahead(struct lia_list *list) +{ + s32 size = (s32)list->entries.count; + if (list->current >= 0 && list->current + 1 < size) { + s32 ahead = list->current + 1; + for (s32 i = ahead; i < MIN(ahead + LIANA_BUFFER_AHEAD, size); i++) { + struct lia_list_entry *entry = al_array_at(list->entries, i); + struct lia_timing time = { + .at = LIANA_TIMESTAMP_INVALID, + .seek_pos = entry->offset, + .pause = LIANA_PAUSE_NONE, + .ended = false + }; + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + sink->callback(sink->userdata, LIANA_SINK_BUFFER, entry, i, &time); + } + } + } +} + +static void set_queued(struct lia_list *list) +{ + if (list->queued < 0) { + return; + } + + struct lia_list_entry *queued = al_array_at(list->entries, list->queued); + + bool error; + struct lia_list_entry *current = al_array_at(list->entries, list->current); + if (!entry_load_and_get_duration(list, current, list->current, &error)) { + return; + } + + queued->start = current->start + (current->duration - current->offset); + u8 pause = queued->paused_at == LIANA_TIMESTAMP_INVALID ? LIANA_PAUSE_RESUME : LIANA_PAUSE_NONE; + struct lia_timing time = { + .at = queued->start, + .seek_pos = queued->offset, + .pause = pause + }; + + struct lia_list_sink *sink; + al_array_foreach(list->sinks, i, sink) { + if (sink->queued != list->queued) { + sink->queued = list->queued; + sink->callback(sink->userdata, LIANA_SINK_BUFFER_AND_QUEUE, queued, list->queued, &time); + } + } +} + +static void evaluate_queued(struct lia_list *list) +{ + s32 size = (s32)list->entries.count; + s32 next = list->current + 1; + if (next >= size || next == list->queued) { + return; + } + + bool error; + struct lia_list_entry *queued = al_array_at(list->entries, next); + if (!entry_load_and_get_duration(list, queued, next, &error)) { + return; + } + + list->queued = next; +} +*/ diff --git a/src/liana/list.h b/src/liana/list.h index 0ffe652..4b1ae02 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 450000Lu // 450ms +#define LIANA_BASE_DELAY 750000Lu // 750ms #define LIANA_BASE_PING 150000Lu // 150ms #define LIANA_PAUSE_DELAY LIANA_BASE_PING #define LIANA_DELAY_IGNORE 0Lu @@ -32,13 +32,16 @@ enum { enum { LIANA_LOAD_ENTRY = 0, - LIANA_GET_DURATION, + LIANA_GET_ENTRY_DURATION, + LIANA_REF_ENTRY, + LIANA_UNREF_ENTRY, LIANA_UNLOAD_ENTRY, LIANA_LIST_META }; enum { - LIANA_ENTRY_PREPARING = 0, + LIANA_ENTRY_UNLOADED = 0, + LIANA_ENTRY_PREPARING, LIANA_ENTRY_PREPARED, LIANA_ENTRY_LOADING, LIANA_ENTRY_LOADED, @@ -109,8 +112,8 @@ struct lia_list { u32 increment; array(struct lia_list_entry *) entries; array(struct lia_list_sink *) sinks; - array(struct lia_list_cmd *) queue; - struct lia_list_cmd *cmd; + array(struct lia_list_cmd *) command_queue; + struct lia_list_cmd *active_cmd; void (*callback)(void *, u8, struct lia_list_entry *, void *); void *userdata; }; @@ -133,7 +136,7 @@ void lia_list_end(struct lia_list *list, u32 id, u32 reset_id); void lia_list_reverse(struct lia_list *list); void lia_list_sort(struct lia_list *list); void lia_list_shuffle(struct lia_list *list); -void lia_list_clear(struct lia_list *list); +//void lia_list_clear(struct lia_list *list); void lia_list_close(struct lia_list *list); void lia_list_free(struct lia_list *list); diff --git a/src/liana/server.c b/src/liana/server.c index 24c4bfc..13ca0aa 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -98,6 +98,19 @@ static void discard_packet_callback(void *userdata, struct nn_packet_stream *str al_assert(false); } +static bool should_free_node(struct lia_node *node) +{ + return node->closed && !node->handler && !node->requests.count && !node->connections.count; +} + +static void free_node(struct lia_node *node) +{ + struct lia_server *server = node->server; + cch_entry_free(&node->entry); + al_array_remove(server->nodes, node); + al_free(node); +} + static void free_connection(struct lia_node_connection *conn) { struct lia_node *node = conn->node; @@ -109,8 +122,8 @@ static void free_connection(struct lia_node_connection *conn) al_array_remove_checked(node->connections, conn, removed); al_assert(removed); al_free(conn); - if (node->closed && !node->connections.count) { - cch_entry_free(&node->entry); + if (should_free_node(node)) { + free_node(node); } } @@ -213,6 +226,8 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet 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) @@ -232,7 +247,7 @@ static void packet_sent_callback(void *userdata, struct nn_packet *packet) static void demote_and_disconnect_stream(struct lia_server *server, struct nn_packet_stream *stream) { // Discard queue based on the currently set packet_sent_callback. - // This should always be the expected behavior but here it's mainly to + // This should always be the expected behavior but here it's important to // not lose packets that belong to the packet pool. nn_packet_stream_discard_queue(stream); stream->userdata = server; @@ -248,34 +263,39 @@ static void signal_callback(void *userdata) struct lia_node_connection *conn = (struct lia_node_connection *)userdata; struct lia_node *node = conn->node; struct lia_server *server = node->server; + 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 || conn->errored) { conn->handler->free(&conn->handler); cch_entry_return_handle(node->entry, &conn->handle); } + if (!packet) { // Connection was closed before init was done. - nn_packet_stream_free(conn->stream); - al_free(conn->stream); + if (should_free_node(node)) { + free_node(node); + } + nn_packet_stream_free(stream); + al_free(stream); al_free(conn); - return; - } - struct nn_packet_stream *stream = conn->stream; - if (!conn->errored) { - conn->id = get_incremental_id(server); - nn_packet_pool_init(&conn->pool, 1024, server->loop, packet_pool_callback, conn); - al_array_push(node->connections, conn); - handle_connection(conn, packet); - } - nn_packet_stream_return_packet(stream, packet); - // conn->errored will not have changed but we want to return the packet before disconnecting. - if (conn->errored) { + } else if (conn->errored) { + // We must return the packet before disconnecting. + nn_packet_stream_return_packet(stream, packet); al_free(conn); conn = NULL; demote_and_disconnect_stream(server, stream); + } else { + 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); + handle_connection(conn, packet); } } @@ -302,9 +322,7 @@ static struct lia_node_connection *get_connection_from_id(struct lia_node *node, { struct lia_node_connection *conn; al_array_foreach(node->connections, i, conn) { - if (conn->id == id) { - return conn; - } + if (conn->id == id) return conn; } return NULL; } @@ -329,6 +347,10 @@ 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) { + goto err; + } + struct lia_node_connection *conn = NULL; if (connection_id == 0) { conn = al_alloc_object(struct lia_node_connection); @@ -346,6 +368,7 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str conn->seek_pos = LIANA_TIMESTAMP_INVALID; stream->packet_callback = discard_packet_callback; stream->connection_closed_callback = pre_init_connection_closed_callback; + al_array_push(node->requests, conn); nn_thread_create(&conn->thread, init_thread, conn); } else { if ((conn = get_connection_from_id(node, connection_id))) { @@ -363,13 +386,16 @@ static void packet_callback(void *userdata, struct nn_packet_stream *stream, str al_assert(conn->node == node); conn->stream = stream; handle_connection(conn, packet); - } - nn_packet_stream_return_packet(stream, packet); - // Return packet before possibly disconnecting. - if (!conn) { - nn_packet_stream_disconnect(stream); + } else { + goto err; } } + + return; +err: + // Return packet before disconnecting. + nn_packet_stream_return_packet(stream, packet); + nn_packet_stream_disconnect(stream); } static bool connection_callback(void *userdata, struct nn_packet_stream *stream) @@ -393,8 +419,10 @@ struct lia_node *lia_server_create_node(struct lia_server *server, struct cch_en struct lia_node *node = al_alloc_object(struct lia_node); node->id = get_incremental_id(server); node->entry = entry; + al_array_init(node->requests); al_array_init(node->connections); node->closed = false; + node->duration = LIANA_TIMESTAMP_INVALID; node->server = server; al_array_push(server->nodes, node); return node; @@ -405,7 +433,6 @@ static nn_thread_result NNWT_THREADCALL init_duration_thread(void *userdata) struct lia_node *node = (struct lia_node *)userdata; if (!node->handler->init(node->handler, &node->handle)) { node->errored = true; - node->duration = LIANA_TIMESTAMP_INVALID; } else { node->duration = node->handler->get_duration(node->handler); } @@ -420,7 +447,11 @@ static void duration_signal_callback(void *userdata) nn_thread_join(&node->thread); node->handler->free(&node->handler); cch_entry_return_handle(node->entry, &node->handle); - node->callback(node->userdata, LIANA_NODE_DURATION, node->duration); + if (should_free_node(node)) { + free_node(node); + } else { + node->callback(node->userdata, LIANA_NODE_DURATION, node->duration); + } } void lia_node_get_duration(struct lia_node *node) @@ -437,10 +468,13 @@ void lia_node_get_duration(struct lia_node *node) void lia_node_close(struct lia_node *node) { node->closed = true; - if (!node->connections.count) { - cch_entry_free(&node->entry); + if (should_free_node(node)) { + free_node(node); } else { struct lia_node_connection *conn; + al_array_foreach_rev(node->requests, i, conn) { + nn_packet_stream_disconnect(conn->stream); + } al_array_foreach_rev(node->connections, i, conn) { if (conn->stream) { conn->disconnected = true; @@ -464,6 +498,8 @@ void lia_server_free(struct lia_server *server) { struct lia_node *node; al_array_foreach(server->nodes, i, node) { + al_assert(!node->requests.count); + al_array_free(node->requests); al_assert(!node->connections.count); al_array_free(node->connections); al_free(node); diff --git a/src/liana/server.h b/src/liana/server.h index 7d0b561..364008e 100644 --- a/src/liana/server.h +++ b/src/liana/server.h @@ -46,6 +46,7 @@ enum { struct lia_node { u32 id; struct cch_entry *entry; + array(struct lia_node_connection *) requests; array(struct lia_node_connection *) connections; bool closed; u64 duration; diff --git a/src/liana/vcr.c b/src/liana/vcr.c index c6f6d93..70fb7b6 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -3,7 +3,7 @@ #include "vcr.h" #include "handler.h" -#define VCR_BUFFER_BUFFERED MB(6) +#define VCR_BUFFER_BUFFERED MB(24) enum { VCR_EXPAND_UNTOUCHED = 0, @@ -18,7 +18,7 @@ enum { }; #define VCR_TRACK_THREADED(track) \ - (track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO) + (track->stream->type == CAMU_STREAM_AUDIO || track->stream->type == CAMU_STREAM_VIDEO) #ifndef CAMU_DIRECT_MODE static void signal_callback(void *userdata) @@ -30,8 +30,8 @@ static void signal_callback(void *userdata) static void reset_metrics(struct lia_vcr *vcr) { - vcr->metric.current_frame = 0Lu; - vcr->metric.last_report_ts = 0Lu; + vcr->metric.current_frame = 0; + vcr->metric.last_report_ts = 0; } void lia_vcr_init(struct lia_vcr *vcr, struct nn_event_loop *loop, struct nn_packet_stream *data) @@ -198,7 +198,7 @@ static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index { struct lia_vcr_track *track; al_array_foreach(vcr->tracks, i, track) { - if (track->index == index) return track; + if (track->stream->index == index) return track; } return NULL; } @@ -238,17 +238,17 @@ static void update_metrics(struct lia_vcr *vcr, u32 size) return; } u64 diff; - if ((diff = now - vcr->metric.last_report_ts) > 1000000Lu) { + if ((diff = now - vcr->metric.last_report_ts) > 1000000) { vcr->metric.last_report_ts = now; u64 frame = vcr->metric.current_frame; - vcr->metric.current_frame = 0Lu; - if (diff > 2500000Lu) { - // We are buffering fast enough for it to not matter. - al_log_debug("vcr", "Ignoring %llu bytes in metrics.", frame); + vcr->metric.current_frame = 0; + if (diff > 3000000) { return; } f32 kbps = (frame / 125.f) / (diff / 1000000.f); - al_log_info("vcr", "Receiving packets at %.2fkbps.", kbps); + f32 size = vcr->mark.buffered / (f32)MB(1); + f32 buffered = al_atomic_load(u64)(&vcr->count, AL_ATOMIC_RELAXED) / (f32)MB(1); + al_log_info("vcr", "Receiving packets at %.2fkbps (%.2f/%.2fMB).", kbps, buffered, size); } } @@ -262,7 +262,7 @@ void lia_vcr_push_packet(struct lia_vcr *vcr, struct nn_packet *packet) update_metrics(vcr, size); track = get_track_from_index(vcr, nn_packet_read_s32(packet)); if (!track) { - al_log_warn("liana", "Received data from errored or unknown track."); + al_log_debug("liana", "Received data from errored or unknown track."); break; } if (VCR_TRACK_THREADED(track)) { @@ -395,11 +395,6 @@ void lia_vcr_free(struct lia_vcr *vcr) al_array_foreach(vcr->tracks, i, track) { nn_packet_cache_free(&track->cache); track->client->free(&track->client); -#ifdef CAMU_HAVE_FFMPEG - if (track->stream.mode == CAMU_FFMPEG_COMPAT) { - avformat_free_context(track->stream.av.format_context); - } -#endif al_free(track); } al_array_free(vcr->tracks); diff --git a/src/liana/vcr.h b/src/liana/vcr.h index 8330365..94d0c47 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -11,8 +11,7 @@ //#define VCR_BUFFER_WHOLE_FILE struct lia_vcr_track { - s32 index; - struct camu_codec_stream stream; + struct camu_codec_stream *stream; struct lia_client_handler *client; atomic(s32) state; atomic(bool) buffered; |