diff options
Diffstat (limited to 'src')
| -rw-r--r-- | src/buffer/clock.c | 1 | ||||
| -rw-r--r-- | src/buffer/video.c | 5 | ||||
| -rw-r--r-- | src/cache/meson.build | 19 | ||||
| -rw-r--r-- | src/codec/ffmpeg/decoder.c | 2 | ||||
| -rw-r--r-- | src/liana/client.c | 22 | ||||
| -rw-r--r-- | src/liana/client.h | 3 | ||||
| -rw-r--r-- | src/liana/handler.h | 4 | ||||
| -rw-r--r-- | src/liana/handlers/cdio_server.c | 1 | ||||
| -rw-r--r-- | src/liana/list.c | 4 | ||||
| -rw-r--r-- | src/liana/list.h | 21 | ||||
| -rw-r--r-- | src/liana/meson.build | 13 | ||||
| -rw-r--r-- | src/liana/vcr.c | 6 | ||||
| -rw-r--r-- | src/libsink/common.h | 2 | ||||
| -rw-r--r-- | src/libsink/sink.c | 26 | ||||
| -rw-r--r-- | src/libsink/sink.h | 1 | ||||
| -rw-r--r-- | src/server/local_compat.c | 10 | ||||
| -rw-r--r-- | src/sink/desktop.c | 1 |
17 files changed, 83 insertions, 58 deletions
diff --git a/src/buffer/clock.c b/src/buffer/clock.c index c4a07a6..a04bb08 100644 --- a/src/buffer/clock.c +++ b/src/buffer/clock.c @@ -59,7 +59,6 @@ void camu_clock_resume(struct camu_clock *clock, u64 target) //} else if (target > 0) { if (target > 0) { tick = calc_tick_offset(tick, aki_get_timestamp(), target); - al_printf("%f\n", tick); } if (clock->paused_at == 0.0) { if (target == 0) { diff --git a/src/buffer/video.c b/src/buffer/video.c index b2381b2..a7a5594 100644 --- a/src/buffer/video.c +++ b/src/buffer/video.c @@ -149,6 +149,7 @@ void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_fra // Not thread-safe, must be called while the buffer is not being read from or written to. void camu_video_buffer_reset(struct camu_video_buffer *buf) { + buf->pts = -1.0; buf->queue->reset(buf->queue); buf->buffered = false; al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED); @@ -172,7 +173,7 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out) f64 pts = camu_clock_get_pts(buf->clock, buf->latency); if (pts > buf->pts) buf->pts = pts; } - u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_RELAXED); + u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE); if (flow == SIGNALED && !buf->single_frame) { buf->callback(buf->userdata, CAMU_BUFFER_EOF); return false; @@ -180,7 +181,7 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out) u8 ret = buf->queue->read(buf->queue, buf->pts, out); if (flow == FLUSHED && (ret == CAMU_QUEUE_EOF || (buf->single_frame && ret == CAMU_QUEUE_OK))) { buf->callback(buf->userdata, CAMU_BUFFER_EOF); - al_atomic_store(u8)(&buf->flow, SIGNALED, AL_ATOMIC_RELAXED); + al_atomic_store(u8)(&buf->flow, SIGNALED, AL_ATOMIC_RELEASE); al_log_debug("video_buffer", "Flushed."); } else if (flow == FLOWING && buf->queue->count(buf->queue) <= BUFFER_WATERMARK_LOW) { buf->callback(buf->userdata, CAMU_BUFFER_UNCORK); diff --git a/src/cache/meson.build b/src/cache/meson.build index 8684ccd..5fb47ae 100644 --- a/src/cache/meson.build +++ b/src/cache/meson.build @@ -3,9 +3,24 @@ cache_src = [ 'handle.c', 'threaded_waits.c', 'handlers/file.c', - 'handlers/cdio.c', 'backings/memory.c' ] +cache_deps = [] + +cache_have_cdio = false +libcdio_paranoia = dependency('libcdio_paranoia', required: false, allow_fallback: true) +libcdio_cdda = dependency('libcdio_cdda', required: false, allow_fallback: true) +if libcdio_paranoia.found() and libcdio_cdda.found() + cache_src += ['handlers/cdio.c'] + cache_deps += [libcdio_paranoia, libcdio_cdda] + cache_have_cdio = true +endif + +#libdvdcss = dependency('libdvdcss', required: false, allow_fallback: true) +#libdvdread = dependency('libdvdread', required: false, allow_fallback: true) +#libdvdnav = dependency('libdvdnav', required: false, allow_fallback: true) +#if libdvdcss.found() and libdvdread.found() and libdvdnav.found() +#endif if akiyo_has_mmap cache_src += ['backings/file_mapped.c'] @@ -17,4 +32,4 @@ if akiyo_has_curl cache_src += ['handlers/http.c'] endif -cache = declare_dependency(sources: cache_src) +cache = declare_dependency(sources: cache_src, dependencies: cache_deps) diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c index d231b2f..3aafd35 100644 --- a/src/codec/ffmpeg/decoder.c +++ b/src/codec/ffmpeg/decoder.c @@ -14,6 +14,8 @@ static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *rend { struct camu_ff_decoder *av = (struct camu_ff_decoder *)dec; + av->codec_context = NULL; + AVCodecParameters *codecpar = stream->av.stream->codecpar; const AVCodec *codec = avcodec_find_decoder(codecpar->codec_id); diff --git a/src/liana/client.c b/src/liana/client.c index d556453..5e33086 100644 --- a/src/liana/client.c +++ b/src/liana/client.c @@ -5,6 +5,7 @@ #include "client.h" #include "handlers.h" +#include "list.h" static void data_packet_callback(void *userdata, struct aki_packet_stream *stream, struct aki_packet *packet) { @@ -84,9 +85,17 @@ static void parse_info_packet(struct lia_client *client, struct aki_packet *pack track->client->callback = client->callback; track->client->userdata = client->userdata; track->stream.mode = mode; - if (track->client->init(track->client, client->renderer, &track->stream)) { - client->callback(client->userdata, LIANA_CLIENT_CONFIGURE, track->client->stream, track); + if (!track->client->init(track->client, client->renderer, &track->stream)) { + track->client->free(&track->client); +#if 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); lia_vcr_add_track(&client->vcr, track); } } @@ -141,6 +150,12 @@ static void connection_closed_callback(void *userdata, struct aki_packet_stream client->callback(client->userdata, LIANA_CLIENT_REMOVE_BUFFERS, NULL, &client->reconnect); if (client->reconnect) { client->reconnect = false; + struct lia_timing time = { + .at = client->at, + .seek_pos = client->pos, + .pause = LIANA_PAUSE_NONE + }; + client->callback(client->userdata, LIANA_CLIENT_RESUME_AT, NULL, &time); aki_packet_stream_reconnect(stream, &client->addr, client->port); } else { client->callback(client->userdata, LIANA_CLIENT_CLOSED, NULL, NULL); @@ -173,10 +188,11 @@ void lia_client_set_renderer(struct lia_client *client, struct camu_renderer *re client->renderer = renderer; } -void lia_client_seek(struct lia_client *client, u64 pos) +void lia_client_seek(struct lia_client *client, u64 pos, u64 at) { client->reconnect = true; client->pos = pos; + client->at = at; aki_packet_stream_disconnect(&client->data); } diff --git a/src/liana/client.h b/src/liana/client.h index 5c16cd8..24526af 100644 --- a/src/liana/client.h +++ b/src/liana/client.h @@ -11,6 +11,7 @@ struct lia_client { u16 id; s32 mask; u64 pos; + u64 at; bool reconnect; str addr; u16 port; @@ -25,7 +26,7 @@ struct lia_client { void lia_client_connect(struct lia_client *client, struct aki_event_loop *loop, u8 type, str *addr, u16 port, u16 id, u64 pos, struct camu_renderer *renderer); -void lia_client_seek(struct lia_client *client, u64 pos); +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); void lia_client_free(struct lia_client *client); diff --git a/src/liana/handler.h b/src/liana/handler.h index 6c23ac5..f994069 100644 --- a/src/liana/handler.h +++ b/src/liana/handler.h @@ -25,11 +25,9 @@ enum { enum { LIANA_CLIENT_CONFIGURE = 0, - LIANA_CLIENT_SET, - LIANA_CLIENT_PAUSE, - LIANA_CLIENT_RESUME, LIANA_CLIENT_DATA, LIANA_CLIENT_REMOVE_BUFFERS, + LIANA_CLIENT_RESUME_AT, LIANA_CLIENT_EOF, LIANA_CLIENT_CLOSED }; diff --git a/src/liana/handlers/cdio_server.c b/src/liana/handlers/cdio_server.c index 1eaf6f3..118f6b5 100644 --- a/src/liana/handlers/cdio_server.c +++ b/src/liana/handlers/cdio_server.c @@ -1,5 +1,4 @@ #include <al/log.h> -#include <cdio/paranoia/cdda.h> #include "../../cache/entry.h" #include "../../cache/handlers/cdio.h" diff --git a/src/liana/list.c b/src/liana/list.c index 37a144f..18bfd6c 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -59,7 +59,7 @@ void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struc u64 seek_pos = current->offset; if (current->paused_at == LIANA_TIMESTAMP_INVALID) { now += LIANA_BASE_DELAY; - if (now > current->start) { + if (now > current->start && now - current->start > LIANA_BASE_DELAY) { at = now; seek_pos += now - current->start; } else { @@ -333,7 +333,7 @@ void lia_list_end(struct lia_list *list, s32 sequence) s32 next = sequence + 1; list->previous = sequence; struct lia_list_entry *current = al_array_at(list->entries, list->current); - //current->offset = current->duration; + current->offset = current->duration; if (list->queued >= 0) { list->current = list->queued; list->queued = -1; diff --git a/src/liana/list.h b/src/liana/list.h index 8adea17..eb9a3cb 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -25,25 +25,10 @@ enum { LIANA_META_QUEUED }; -// Pause: -// - Set `paused_at` to now() + PAUSE_DELAY. -// - Increment `offset` by how long the entry will have been -// playing when paused at the requested timestamp (`paused_at` - `start`). -// Resume: -// - Unset `paused_at` -// - Set `start` to now() + PAUSE_DELAY. -// Next/Prev: -// Current Playing, Target Playing: -// - Set `current->held_at` to now() + PAUSE_DELAY. -// - Increment `current->offset` by the same logic as pause. -// - -// Current Playing, Target Paused: -// Current Paused, Target Playing: -// Current Paused, Target Paused: -// Seek: +// NOTE: To handle an entry being queued right before a skip, keep a global +// "max time until all sinks buffered" and used that instead of LIANA_PAUSE_DELAY (if greater). -// NOTE: keep global per list max time until all sinks _should_ by buffered -// possibly use that instead of LIANA_PAUSE_DELAY. +// TODO: Factor in LIANA_BASE_PING. enum { LIANA_PAUSE_NONE = 0, diff --git a/src/liana/meson.build b/src/liana/meson.build index 0859e99..da53214 100644 --- a/src/liana/meson.build +++ b/src/liana/meson.build @@ -13,23 +13,12 @@ liana_client_src = [ liana_deps = [codecs] liana_args = [] -libcdio_paranoia = dependency('libcdio_paranoia', required: false, allow_fallback: true) -libcdio_cdda = dependency('libcdio_cdda', required: false, allow_fallback: true) -if libcdio_paranoia.found() and libcdio_cdda.found() +if cache_have_cdio liana_server_src += ['handlers/cdio_server.c'] liana_client_src += ['handlers/cdio_client.c'] - liana_deps += [libcdio_paranoia, libcdio_cdda] liana_args += ['-DLIANA_HAVE_CDIO'] endif -#libdvdcss = dependency('libdvdcss', required: false, allow_fallback: true) -#libdvdread = dependency('libdvdread', required: false, allow_fallback: true) -#libdvdnav = dependency('libdvdnav', required: false, allow_fallback: true) -#if libdvdcss.found() and libdvdread.found() and libdvdnav.found() -# liana_server_src += ['handlers/dvd_server.c'] -# liana_server_deps += [libdvdcss] -#endif - liana_server = declare_dependency(sources: liana_server_src, dependencies: liana_deps, compile_args: [liana_args, '-DLIANA_SERVER']) liana_client = declare_dependency(sources: liana_client_src, dependencies: liana_deps, diff --git a/src/liana/vcr.c b/src/liana/vcr.c index fad8b78..39fe0fe 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -123,7 +123,11 @@ bool lia_vcr_push_packet(struct lia_vcr *vcr, struct aki_packet *packet) switch (op) { case LIANA_PACKET_DATA: track = get_track_from_index(vcr, aki_packet_read_s32(packet)); - al_assert(track); + if (!track) { + al_log_warn("liana", "Received data from errored or unknown track."); + aki_packet_free(packet); + return false; + } if (!track->running) { aki_thread_create(&track->thread, vcr_track_thread, track); track->running = true; diff --git a/src/libsink/common.h b/src/libsink/common.h index deae228..f8dd9fa 100644 --- a/src/libsink/common.h +++ b/src/libsink/common.h @@ -1,6 +1,6 @@ #pragma once -#define CAMU_SINK_LOCAL 0 +#define CAMU_SINK_LOCAL 1 enum { CAMU_SINK_SET = 0, diff --git a/src/libsink/sink.c b/src/libsink/sink.c index 859b717..990cd69 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -210,7 +210,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) case SKIP: { if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_str(packet, &sink->default_list); aki_packet_write_u8(packet, CAMU_LIST_SKIP); //s32 sequence = sink->current ? sink->current->sequence : LIANA_SEQUENCE_ANY; s32 sequence = LIANA_SEQUENCE_ANY; @@ -222,7 +222,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) case SHUFFLE: { if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_str(packet, &sink->default_list); aki_packet_write_u8(packet, CAMU_LIST_SHUFFLE); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); break; @@ -233,7 +233,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) #else if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_str(packet, &sink->default_list); aki_packet_write_u8(packet, CAMU_LIST_TOGGLE_PAUSE); s32 sequence = sink->current ? sink->current->sequence : LIANA_SEQUENCE_ANY; aki_packet_write_s32(packet, sequence); @@ -245,7 +245,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) case SEEK: { if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_str(packet, &sink->default_list); aki_packet_write_u8(packet, CAMU_LIST_SEEK); s32 sequence = sink->current ? sink->current->sequence : LIANA_SEQUENCE_ANY; aki_packet_write_s32(packet, sequence); @@ -256,7 +256,7 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd) case END: { if (!sink->conn) return; struct aki_packet *packet = aki_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION); - aki_packet_write_str(packet, al_str_c("default")); + aki_packet_write_str(packet, &sink->default_list); aki_packet_write_u8(packet, CAMU_LIST_END); aki_packet_write_s32(packet, cmd->value.i); aki_rpc_connection_command(sink->conn, packet, NULL, NULL); @@ -546,6 +546,12 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str camu_audio_buffer_reset(&entry->audio.buf); break; } + case LIANA_CLIENT_RESUME_AT: { + struct lia_timing *time = (struct lia_timing *)opaque; + camu_clock_set(&entry->clock, time->seek_pos / 1000000.0); + camu_clock_resume(&entry->clock, time->at); + break; + } case LIANA_CLIENT_EOF: { switch (stream->type) { case CAMU_STREAM_AUDIO: { @@ -917,10 +923,12 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con struct camu_sink_entry *current = sink->current; aki_mutex_unlock(&sink->mutex); if (current->sequence == sequence) { - // TODO: Thread-safety here. - lia_client_seek(¤t->client, pos); - camu_clock_set(¤t->clock, pos / 1000000.0); - camu_clock_resume(¤t->clock, at); +#if CAMU_SINK_LOCAL + (void)at; + lia_client_seek(¤t->client, pos, 0); +#else + lia_client_seek(¤t->client, pos, at); +#endif } aki_packet_free(packet); diff --git a/src/libsink/sink.h b/src/libsink/sink.h index b9249cd..6d560f6 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -72,6 +72,7 @@ struct camu_sink { struct aki_signal signal; queue(struct camu_sink_cmd) queue; struct aki_mutex mutex; + str default_list; struct camu_sink_entry *current; struct camu_sink_entry *queued; struct camu_sink_entry *target; diff --git a/src/server/local_compat.c b/src/server/local_compat.c index 9870ce5..39698c5 100644 --- a/src/server/local_compat.c +++ b/src/server/local_compat.c @@ -4,7 +4,9 @@ #include "../cache/handlers/http.h" #endif #include "../cache/handlers/file.h" +#ifdef LIANA_HAVE_CDIO #include "../cache/handlers/cdio.h" +#endif #include "local_compat.h" #include "common.h" @@ -84,14 +86,15 @@ static void worker_signal_callback(void *userdata) static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) { struct camu_local_compat *compat = (struct camu_local_compat *)userdata; + aki_thread_setcanceltype(AKI_THREAD_CANCEL_ASYNCHRONOUS); #ifdef CAMU_HAVE_PORTAL bool have_python = false; #endif u32 count; while (aki_packet_cache_wait(&compat->queue, &count)) { struct aki_packet *packet = aki_packet_cache_pop(&compat->queue); + aki_packet_cache_unlock(&compat->queue); if (!packet) { - aki_packet_cache_unlock(&compat->queue); break; } str local; @@ -99,6 +102,7 @@ static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) struct cch_entry *entry = NULL; struct camu_post *post = NULL; if (al_str_cmp(&local, al_str_c("cdda://"), 0, 7) == 0) { +#ifdef LIANA_HAVE_CDIO entry = cch_handler_cdio_create(); struct cch_chapter *chapter = &al_array_at(entry->chapters, 0); if (local.len > 7) { @@ -109,6 +113,7 @@ static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) } entry->chapter = chapter; entry->handler->maybe_spawn_worker(entry->handler, chapter->start); +#endif } else { #ifdef CAMU_HAVE_PORTAL #ifndef AKIYO_HAS_CURL @@ -141,7 +146,6 @@ static aki_thread_result AKI_THREADCALL queue_thread(void *userdata) #endif #endif } - aki_packet_cache_unlock(&compat->queue); if (entry) { struct camu_server_resource *resource = al_alloc_object(struct camu_server_resource); al_str_clone(&resource->unique_id, &local); @@ -203,7 +207,9 @@ void camu_local_compat_send(struct camu_local_compat *compat, struct aki_packet void camu_local_compat_stop(struct camu_local_compat *compat) { aki_packet_cache_disable(&compat->queue); + aki_thread_cancel(&compat->thread); aki_thread_join(&compat->thread); + aki_signal_stop(&compat->worker_signal); aki_signal_stop(&compat->result_signal); camu_queue_free(compat->results); struct aki_packet *packet; diff --git a/src/sink/desktop.c b/src/sink/desktop.c index 8e628fd..ddde86b 100644 --- a/src/sink/desktop.c +++ b/src/sink/desktop.c @@ -156,6 +156,7 @@ bool camu_desktop_connect(struct camu_desktop *c, u8 type, struct aki_event_loop { c->loop = loop; camu_sink_init(&c->sink, c->loop, &c->mixer, c->renderer); + al_str_from(&c->sink.default_list, "default"); c->sink.callback = sink_callback; c->sink.userdata = c; return camu_sink_connect(&c->sink, type, addr, port, al_str_c("desktop")); |