diff options
Diffstat (limited to 'src/liana')
| -rw-r--r-- | src/liana/client.c | 64 | ||||
| -rw-r--r-- | src/liana/handler.h | 1 | ||||
| -rw-r--r-- | src/liana/handlers/cdio_server.c | 8 | ||||
| -rw-r--r-- | src/liana/handlers/codec_client.c | 35 | ||||
| -rw-r--r-- | src/liana/handlers/codec_server.c | 4 | ||||
| -rw-r--r-- | src/liana/list.c | 6 | ||||
| -rw-r--r-- | src/liana/vcr.c | 98 | ||||
| -rw-r--r-- | src/liana/vcr.h | 12 |
8 files changed, 161 insertions, 67 deletions
diff --git a/src/liana/client.c b/src/liana/client.c index 527376a..0edc262 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -22,14 +22,15 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack u32 count = aki_packet_read_u32(packet); for (u32 i = 0; i < count; i++) { u8 mode = aki_packet_read_u8(packet); + u8 type = aki_packet_read_u8(packet); + u64 duration = aki_packet_read_u64(packet); + s32 index = aki_packet_read_s32(packet); + al_assert(index < 32); struct lia_vcr_track *track = NULL; switch (mode) { case CAMU_NORMAL: { - client->mask |= 1 << 0; + client->mask |= 1 << index; track = al_alloc_object(struct lia_vcr_track); - track->stream.mode = mode; - u8 type = aki_packet_read_u8(packet); - track->stream.type = type; if (type == CAMU_STREAM_AUDIO) { struct camu_audio_format *fmt = &track->stream.audio.fmt; fmt->format = aki_packet_read_s32(packet); @@ -44,49 +45,63 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack fmt->height = aki_packet_read_s32(packet); fmt->format = aki_packet_read_s32(packet); } - track->index = 0; break; } #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: { - const AVCodec *codec = avcodec_find_decoder(aki_packet_read_av_codec_id(packet)); - AVFormatContext *format_context = avformat_alloc_context(); - AVStream *stream = aki_packet_read_av_stream(format_context, codec, packet); - s32 index = stream->index; - al_assert(index < 32); - switch (stream->codecpar->codec_type) { - case AVMEDIA_TYPE_AUDIO: + switch (type) { + case CAMU_STREAM_AUDIO: client->mask |= 1 << index; break; - case AVMEDIA_TYPE_VIDEO: + case CAMU_STREAM_VIDEO: client->mask |= 1 << index; - //continue; break; - case AVMEDIA_TYPE_SUBTITLE: + case CAMU_STREAM_SUBTITLE: + client->mask |= 1 << index; + break; + case CAMU_STREAM_ATTACHMENT: + // Assume we have all the data we need in the AVStream object. + break; default: continue; } + enum AVCodecID codec_id = aki_packet_read_av_codec_id(packet); + const AVCodec *codec = avcodec_find_decoder(codec_id); + AVFormatContext *format_context = avformat_alloc_context(); + AVStream *stream = aki_packet_read_av_stream(format_context, codec, packet); + if (type == CAMU_STREAM_SUBTITLE && codec_id != AV_CODEC_ID_ASS) { + client->mask &= ~(1 << index); + continue; + } + if (type == 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); + avformat_free_context(format_context); + continue; + } track = al_alloc_object(struct lia_vcr_track); - track->stream.mode = CAMU_FFMPEG_COMPAT; - track->stream.type = stream->codecpar->codec_type; track->stream.av.format_context = format_context; track->stream.av.stream = stream; - if (track->stream.type == CAMU_STREAM_AUDIO) { + 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; } - track->index = index; break; } #endif } + track->index = index; + track->stream.mode = mode; + track->stream.type = type; + track->stream.duration = duration; track->client = lia_handler_by_name(&liana)->create_client_handler(); track->client->callback = client->callback; track->client->userdata = client->userdata; - track->stream.mode = mode; if (!track->client->init(track->client, client->renderer, &track->stream)) { track->client->free(&track->client); #ifdef CAMU_HAVE_FFMPEG @@ -107,8 +122,12 @@ static void info_packet_callback(void *userdata, struct aki_packet_stream *strea struct lia_client *client = (struct lia_client *)userdata; client->connection_id = aki_packet_read_u16(packet); parse_info_packet(client, packet); - //al_assert(client->mask != 0); aki_packet_free(packet); + if (client->mask == 0 || lia_vcr_is_empty(&client->vcr)) { + client->reconnect = false; + aki_packet_stream_disconnect(&client->data); + return; + } stream->packet_callback = data_packet_callback; struct aki_packet *rpacket = aki_packet_create(); aki_packet_write_s32(rpacket, client->mask); @@ -124,6 +143,7 @@ static void packet_sent_callback(void *userdata, struct aki_packet *packet) static void connection_callback(void *userdata, struct aki_packet_stream *stream) { struct lia_client *client = (struct lia_client *)userdata; + lia_vcr_start(&client->vcr, client->loop); stream->packet_sent_callback = packet_sent_callback; struct aki_packet *packet = aki_packet_create(); aki_packet_write_u16(packet, client->id); @@ -173,8 +193,6 @@ void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop, client->mask = 0; client->reconnect = false; lia_vcr_init(&client->vcr, &client->data); - // TODO: Starting and stopping of vcr could be more clear. - lia_vcr_start(&client->vcr, client->loop); al_str_clone(&client->addr, addr); client->port = port; if (!aki_packet_stream_init(&client->data, type, connection_callback, connection_closed_callback, client)) { diff --git a/src/liana/handler.h b/src/liana/handler.h index f994069..84aa1de 100644 --- a/src/liana/handler.h +++ b/src/liana/handler.h @@ -26,6 +26,7 @@ enum { enum { LIANA_CLIENT_CONFIGURE = 0, LIANA_CLIENT_DATA, + LIANA_CLIENT_SUBTITLE, LIANA_CLIENT_REMOVE_BUFFERS, LIANA_CLIENT_RESUME_AT, LIANA_CLIENT_EOF, diff --git a/src/liana/handlers/cdio_server.c b/src/liana/handlers/cdio_server.c index d271829..641ed92 100644 --- a/src/liana/handlers/cdio_server.c +++ b/src/liana/handlers/cdio_server.c @@ -32,6 +32,8 @@ static void cdio_server_write_info(struct lia_server_handler *handler, struct ak aki_packet_write_u32(packet, 1); aki_packet_write_u8(packet, CAMU_NORMAL); aki_packet_write_u8(packet, CAMU_STREAM_AUDIO); + aki_packet_write_u64(packet, cdio->handler.get_duration(&cdio->handler)); + aki_packet_write_s32(packet, 0); aki_packet_write_s32(packet, cdio->fmt.format); aki_packet_write_s32(packet, cdio->fmt.sample_rate); aki_packet_write_s32(packet, cdio->fmt.channel_count); @@ -48,9 +50,9 @@ static u64 cdio_server_get_duration(struct lia_server_handler *handler) struct lia_cdio_server *cdio = (struct lia_cdio_server *)handler; struct cch_chapter *first = &al_array_at(cdio->handle->entry->chapters, 0); struct cch_chapter *last = &al_array_last(cdio->handle->entry->chapters); - f64 seconds = camu_audio_format_bytes_to_sec(&cdio->fmt, (last->end - first->start) * CDIO_CD_FRAMESIZE_RAW); - al_log_info("cdio", "Length: %.2fs.", seconds); - return (u64)(seconds * 1000000.0); + u64 length = camu_audio_format_bytes_to_usec(&cdio->fmt, (last->end - first->start) * CDIO_CD_FRAMESIZE_RAW); + al_log_info("cdio", "Length: %.2fs.", length / 1000000.0); + return length; } static bool cdio_server_seek(struct lia_server_handler *handler, u64 pos) diff --git a/src/liana/handlers/codec_client.c b/src/liana/handlers/codec_client.c index d420c59..7c57d03 100644 --- a/src/liana/handlers/codec_client.c +++ b/src/liana/handlers/codec_client.c @@ -21,13 +21,15 @@ static bool codec_client_init(struct lia_client_handler *handler, struct camu_re struct camu_codec_stream *stream) { struct lia_codec_client *codec = (struct lia_codec_client *)handler; - codec->dec = camu_ff_decoder_create(); - //codec->dec = camu_stbi_decoder_create(); - //codec->dec = camu_spng_decoder_create(); - //codec->dec = camu_wuffs_decoder_create(); codec->handler.stream = stream; - if (!codec->dec->init(codec->dec, renderer, stream, data_callback, codec)) { - return false; + if (stream->type == CAMU_STREAM_AUDIO || stream->type == CAMU_STREAM_VIDEO) { + codec->dec = camu_ff_decoder_create(); + //codec->dec = camu_stbi_decoder_create(); + //codec->dec = camu_spng_decoder_create(); + //codec->dec = camu_wuffs_decoder_create(); + if (!codec->dec->init(codec->dec, renderer, stream, data_callback, codec)) { + return false; + } } return true; } @@ -38,8 +40,6 @@ static bool push_av_packet(struct lia_codec_client *codec, AVPacket *pkt) struct camu_codec_packet packet; packet.av.pkt = pkt; s32 ret = codec->dec->push(codec->dec, &packet); - av_packet_unref(pkt); - av_packet_free(&pkt); return ret == CAMU_OK; } #endif @@ -56,7 +56,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc { struct lia_codec_client *codec = (struct lia_codec_client *)handler; if (!packet) { - struct lia_codec_client *codec = (struct lia_codec_client *)handler; + if (!codec->dec) return true; s32 ret = codec->dec->push(codec->dec, NULL); // Flush returns success. ret = codec->dec->process(codec->dec); @@ -68,6 +68,7 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc u8 type = aki_packet_read_u8(packet); switch (type) { case CAMU_NORMAL: { + if (!codec->dec) return true; struct aki_buffer buffer; aki_packet_read_buffer(packet, &buffer); success = push_packet(codec, &buffer); @@ -75,12 +76,20 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc } #ifdef CAMU_HAVE_FFMPEG case CAMU_FFMPEG_COMPAT: { - success = push_av_packet(codec, aki_packet_read_av_packet(packet)); + AVPacket *pkt = aki_packet_read_av_packet(packet); + if (!codec->dec) { + codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_SUBTITLE, codec->handler.stream, pkt); + return true; + } else { + success = push_av_packet(codec, pkt); + } + av_packet_unref(pkt); + av_packet_free(&pkt); break; } #endif } - // Forcing in EOF on errors is not necessary but should be a better experience client-side. + // Forcing in EOF on an error is not necessary but should be a better experience client-side. if (!success) { codec->handler.callback(codec->handler.userdata, LIANA_CLIENT_EOF, codec->handler.stream, NULL); return false; @@ -96,13 +105,13 @@ static bool codec_client_handle_packet(struct lia_client_handler *handler, struc static void codec_client_flush(struct lia_client_handler *handler) { struct lia_codec_client *codec = (struct lia_codec_client *)handler; - codec->dec->flush(codec->dec); + if (codec->dec) codec->dec->flush(codec->dec); } static void codec_client_free(struct lia_client_handler **handler) { struct lia_codec_client *codec = (struct lia_codec_client *)*handler; - codec->dec->free(&codec->dec); + if (codec->dec) codec->dec->free(&codec->dec); al_free(codec); *handler = NULL; } diff --git a/src/liana/handlers/codec_server.c b/src/liana/handlers/codec_server.c index 58da2e6..94b18df 100644 --- a/src/liana/handlers/codec_server.c +++ b/src/liana/handlers/codec_server.c @@ -42,9 +42,11 @@ static void codec_server_write_info(struct lia_server_handler *handler, struct a struct camu_codec_stream *stream; al_array_foreach_ptr(codec->demux->streams, i, stream) { aki_packet_write_u8(packet, stream->mode); + aki_packet_write_u8(packet, stream->type); + aki_packet_write_u64(packet, stream->duration); + aki_packet_write_s32(packet, i); switch (stream->mode) { case CAMU_NORMAL: { - aki_packet_write_u8(packet, stream->type); struct camu_video_format *fmt = &stream->video.fmt; aki_packet_write_s32(packet, fmt->width); aki_packet_write_s32(packet, fmt->height); diff --git a/src/liana/list.c b/src/liana/list.c index 65efc9d..09720b6 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -151,7 +151,7 @@ static bool handle_add(struct lia_list *list, struct lia_list_entry *entry) sink->set = list->current; sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time); } - // meta playing + al_log_info("list", "Now playing: %ls.\n", AL_WSTR_PRINTF(&entry->name)); } else { /* if (list->queued == -1) { @@ -309,7 +309,7 @@ static bool handle_skipto(struct lia_list *list, s32 sequence, s32 index) sink->callback(sink->userdata, LIANA_SINK_SET, target, index, &time); } - // meta playing + al_log_info("list", "Now playing: %ls.\n", AL_WSTR_PRINTF(&target->name)); return true; } @@ -403,7 +403,7 @@ static void handle_end(struct lia_list *list, s32 sequence) al_array_foreach(list->sinks, i, sink) { sink->queued = -1; } - // meta playing + al_log_info("list", "Now playing: %ls.\n", AL_WSTR_PRINTF(¤t->name)); } else if (next < size) { struct lia_list_cmd *cmd = al_alloc_object(struct lia_list_cmd); cmd->op = SKIPTO; diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 213d5cb..c9e7975 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -3,9 +3,18 @@ #include "vcr.h" #include "handler.h" -#define VCR_BUFFER_INIT 256 -#define VCR_BUFFER_BUFFERED (VCR_BUFFER_INIT - 16) -#define VCR_BUFFER_LOW (VCR_BUFFER_INIT - 64) +#define VCR_BUFFER_BUFFERED MB(24) + +enum { + VCR_EXPAND_UNTOUCHED = 0, + VCR_EXPAND_GROWN, + VCR_EXPAND_COMPLETE +}; + +// Only track the buffered state of audio and video streams as we don't +// expect any other type of stream to ever call cork(). +#define TRACK_IGNORE_BUFFERED(track) \ + (!(track->stream.type == CAMU_STREAM_AUDIO || track->stream.type == CAMU_STREAM_VIDEO)) static void signal_callback(void *userdata) { @@ -13,14 +22,22 @@ static void signal_callback(void *userdata) aki_packet_stream_cork(vcr->data, false); } +static void reset_metrics(struct lia_vcr *vcr) +{ + vcr->metric.current_frame = 0; + vcr->metric.last_report_ts = 0; +} + void lia_vcr_init(struct lia_vcr *vcr, struct aki_packet_stream *data) { al_array_init(vcr->tracks); - al_atomic_store(s32)(&vcr->count, 0, AL_ATOMIC_RELAXED); - vcr->mark.low = VCR_BUFFER_LOW; + al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); vcr->mark.buffered = VCR_BUFFER_BUFFERED; + vcr->mark.low = 0; + vcr->expand = VCR_EXPAND_UNTOUCHED; vcr->data = data; aki_signal_init(&vcr->signal, signal_callback, vcr); + reset_metrics(vcr); } void lia_vcr_start(struct lia_vcr *vcr, struct aki_event_loop *loop) @@ -45,9 +62,10 @@ static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata) // We were signaled to close, exit thread. goto out; } + u32 size = aki_packet_get_size(packet); u8 buffered = al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED); - s32 count = al_atomic_sub(s32)(&vcr->count, 1, AL_ATOMIC_RELAXED); - if (count <= vcr->mark.low && buffered) { + u64 buffer = al_atomic_sub(u64)(&vcr->count, size, AL_ATOMIC_RELAXED); + if (buffered && buffer <= vcr->mark.low) { aki_signal_send(&vcr->signal); } } @@ -90,12 +108,17 @@ void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track) track->vcr = vcr; aki_cond_init(&track->cond); aki_mutex_init(&track->mutex); - al_atomic_store(u8)(&track->buffered, 0, AL_ATOMIC_RELAXED); - aki_packet_cache_init(&track->cache, VCR_BUFFER_INIT); + al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED); + aki_packet_cache_init(&track->cache, 256); al_array_push(vcr->tracks, track); al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); } +bool lia_vcr_is_empty(struct lia_vcr *vcr) +{ + return vcr->tracks.size == 0; +} + static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index) { struct lia_vcr_track *track; @@ -105,7 +128,7 @@ static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index return NULL; } -static void cork_if_buffered(struct lia_vcr *vcr) +static void cork_if_buffered(struct lia_vcr *vcr, u64 buffer) { u8 buffered = 1; struct lia_vcr_track *track; @@ -113,6 +136,15 @@ static void cork_if_buffered(struct lia_vcr *vcr) buffered &= al_atomic_load(u8)(&track->buffered, AL_ATOMIC_RELAXED); } if (buffered) { + if (vcr->expand == VCR_EXPAND_UNTOUCHED) { + vcr->mark.buffered = buffer * 2; + vcr->expand = VCR_EXPAND_GROWN; + al_log_info("vcr", "Expanded buffer to size %.2fMB.", vcr->mark.buffered / (f32)MB(1)); + return; + } else if (vcr->expand == VCR_EXPAND_GROWN) { + vcr->mark.low = vcr->mark.buffered - MB(2); + vcr->expand = VCR_EXPAND_COMPLETE; + } aki_packet_stream_cork(vcr->data, true); al_array_foreach(vcr->tracks, i, track) { aki_packet_cache_flush(&track->cache); @@ -120,7 +152,24 @@ static void cork_if_buffered(struct lia_vcr *vcr) } } -bool lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet) +static void update_metrics(struct lia_vcr *vcr, u32 size) +{ + vcr->metric.current_frame += size; + u64 now = aki_get_timestamp(); + if (!vcr->metric.last_report_ts) { + vcr->metric.last_report_ts = now; + return; + } + u64 diff; + if ((diff = now - vcr->metric.last_report_ts) > 1000000Lu) { + f32 kbps = (vcr->metric.current_frame / 125.f) / (diff / 1000000.f); + al_log_info("vcr", "Receiving packets at %.2fkbps.", kbps); + vcr->metric.current_frame = 0; + vcr->metric.last_report_ts = now; + } +} + +void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet) { struct lia_vcr_track *track; u8 op = aki_packet_read_u8(packet); @@ -130,34 +179,37 @@ bool lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet) if (!track) { al_log_warn("liana", "Received data from errored or unknown track."); aki_packet_free(packet); - return false; + return; } if (!track->running) { aki_thread_create(&track->thread, vcr_track_thread, track); track->running = true; } - s32 count; + u32 size = aki_packet_get_size(packet); if (!aki_packet_cache_send_packet(&track->cache, packet)) { aki_packet_free(packet); - } else if ((count = al_atomic_add(s32)(&vcr->count, 1, AL_ATOMIC_RELAXED)) > vcr->mark.buffered) { - vcr->mark.buffered = count; - cork_if_buffered(vcr); + return; + } + u64 buffer; + if ((buffer = al_atomic_add(u64)(&vcr->count, size, AL_ATOMIC_RELAXED)) >= vcr->mark.buffered) { + cork_if_buffered(vcr, buffer); } + update_metrics(vcr, size); break; case LIANA_PACKET_EOF: al_array_foreach(vcr->tracks, i, track) { aki_packet_cache_send_packet(&track->cache, NULL); } aki_packet_free(packet); + aki_signal_stop(&vcr->signal); break; case LIANA_PACKET_ERROR: al_log_warn("liana", "Unhandled error packet."); aki_packet_free(packet); break; default: - al_assert_and_return(false); + al_assert(false); } - return true; } void lia_vcr_cork(struct lia_vcr_track *track) @@ -206,11 +258,15 @@ void lia_vcr_flush(struct lia_vcr *vcr) vcr_track_close_internal(track); track->client->flush(track->client); aki_packet_cache_enable(&track->cache); - al_atomic_store(u8)(&track->buffered, 0, AL_ATOMIC_RELAXED); + al_atomic_store(u8)(&track->buffered, TRACK_IGNORE_BUFFERED(track), AL_ATOMIC_RELAXED); al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); } - aki_signal_send(&vcr->signal); - al_atomic_store(s32)(&vcr->count, 0, AL_ATOMIC_RELAXED); + aki_signal_stop(&vcr->signal); + al_atomic_store(u64)(&vcr->count, 0, AL_ATOMIC_RELAXED); + if (vcr->expand == VCR_EXPAND_COMPLETE) { + vcr->mark.low = 0; + vcr->expand = VCR_EXPAND_GROWN; + } } void lia_vcr_close_all(struct lia_vcr *vcr) diff --git a/src/liana/vcr.h b/src/liana/vcr.h index f471c4e..fc7fb32 100644 --- a/src/liana/vcr.h +++ b/src/liana/vcr.h @@ -29,16 +29,22 @@ struct lia_vcr_track { struct lia_vcr { array(struct lia_vcr_track *) tracks; - atomic(s32) count; - struct { s32 low, buffered; } mark; + atomic(u64) count; + struct { u64 buffered, low; } mark; + u8 expand; struct aki_packet_stream *data; struct aki_signal signal; + struct { + u64 current_frame; + u64 last_report_ts; + } metric; }; void lia_vcr_init(struct lia_vcr *vcr, struct aki_packet_stream *data); void lia_vcr_start(struct lia_vcr *vcr, struct aki_event_loop *loop); void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track); -bool lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet); +bool lia_vcr_is_empty(struct lia_vcr *vcr); +void lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet); void lia_vcr_cork(struct lia_vcr_track *track); void lia_vcr_uncork(struct lia_vcr_track *track); void lia_vcr_flush(struct lia_vcr *vcr); |