From 72eb0c6381f9406a4e45d23f65373d4936770433 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Sun, 27 Oct 2024 15:50:21 -0400 Subject: More synced list, another vcr fix Signed-off-by: Andrew Opalach --- src/buffer/audio.c | 25 +++---- src/buffer/clock.c | 10 +-- src/buffer/video.c | 10 +-- src/cache/handle.c | 4 - src/codec/ffmpeg/decoder.c | 5 +- src/liana/list.c | 115 ++++++++++++++++------------- src/liana/list.h | 20 +++-- src/liana/server.c | 2 +- src/liana/vcr.c | 13 ++-- src/libsink/common.h | 10 +-- src/libsink/sink.c | 155 +++++++++++++++++++++------------------ src/libsink/sink.h | 1 + src/mixer/mixer.c | 10 ++- src/render/renderer_libplacebo.c | 38 ++++++---- src/server/list.c | 58 +++++++++++++++ src/server/list.h | 5 ++ src/server/resource.h | 13 ++++ src/server/server.c | 98 ++++++------------------- src/server/server.h | 9 +-- 19 files changed, 324 insertions(+), 277 deletions(-) create mode 100644 src/server/resource.h (limited to 'src') diff --git a/src/buffer/audio.c b/src/buffer/audio.c index 2af24b2..141c589 100644 --- a/src/buffer/audio.c +++ b/src/buffer/audio.c @@ -9,11 +9,9 @@ #include "common_internal.h" #include "volume.h" -#define BUFFER_USEC (3 * 1000000L) -#define BUFFER_MARK_MIN (1.1 * 1000000L) // Must be a most half of the buffer size. -#define BUFFER_MARK_BUFFERED (1.0 * 1000000L) - -#define LARGE_DESYNC_PTS 0.322 +#define BUFFER_USEC (6 * 1000000L) +#define BUFFER_MARK_MIN (2.6 * 1000000L) // Must be a most half of the buffer size. +#define BUFFER_MARK_BUFFERED (1.5 * 1000000L) #ifdef CAMU_AUDIO_BUFFER_FADE #define FADE_STEP(fmt) (1.f / (fmt)->sample_rate) @@ -229,12 +227,18 @@ void camu_audio_buffer_flush(struct camu_audio_buffer *buf) // around in the buffer without worrying about pops. size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t req) { + u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE); + if (flow == SIGNALED || camu_clock_is_ended(buf->clock)) { + buf->callback(buf->userdata, CAMU_BUFFER_EOF); + return 0; + } + if (camu_clock_is_paused(buf->clock)) { #ifdef CAMU_AUDIO_BUFFER_FADE if (buf->pause == PAUSE_PLAYING) { buf->pause = PAUSE_FADING; buf->fade_offset = 0; - } else if (buf->volume == 0.f) { // Fade out done. + } else if (buf->volume == 0.f || buf->pause == PAUSE_PAUSED) { al_memset(data, 0, req); if (NOT_PAUSED(buf->pause)) { buf->pause = PAUSE_PAUSED; @@ -260,12 +264,6 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re al_atomic_store(u8)(&buf->unpause, 0, AL_ATOMIC_RELEASE); } - u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE); - if (flow == SIGNALED) { - buf->callback(buf->userdata, CAMU_BUFFER_EOF); - return 0; - } - f64 pts = camu_clock_get_pts(buf->clock, buf->latency); size_t ret, signal = req; size_t have = al_ring_buffer_occupied(&buf->rb); @@ -279,9 +277,6 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re // Attempt syncing to the clock. // For this to work the mixer must report a reasonably accurate value for latency. if (UNLIKELY(!buf->ignore_desync && buf->pause == PAUSE_PAUSED)) { - if (UNLIKELY(fabs(pts) >= LARGE_DESYNC_PTS)) { - al_log_warn("audio_buffer", "Abnormally large audio desync of %.5fs", pts); - } if (pts > 0.0) { // Skip. ret = camu_audio_format_sec_to_bytes(&buf->fmt.req, pts); ret = AL_MIN(ret, have); diff --git a/src/buffer/clock.c b/src/buffer/clock.c index a04bb08..cee0ec7 100644 --- a/src/buffer/clock.c +++ b/src/buffer/clock.c @@ -35,23 +35,19 @@ void camu_clock_offset(struct camu_clock *clock, f64 amount) void camu_clock_pause(struct camu_clock *clock, u64 target) { - f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE); - if (pause == ENDED) return; al_assert(clock->paused_at == -1.0); f64 tick = aki_get_tick(); if (target > 0) { tick = calc_tick_offset(tick, aki_get_timestamp(), target); - al_atomic_store(f64)(&clock->pause, tick, AL_ATOMIC_RELEASE); + al_atomic_store(f64)(&clock->pause, tick, AL_ATOMIC_RELAXED); } else { - al_atomic_store(f64)(&clock->pause, PAUSED, AL_ATOMIC_RELEASE); + al_atomic_store(f64)(&clock->pause, PAUSED, AL_ATOMIC_RELAXED); } clock->paused_at = tick; } void camu_clock_resume(struct camu_clock *clock, u64 target) { - f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE); - if (pause == ENDED) return; al_assert(clock->paused_at != -1.0); f64 tick = aki_get_tick(); //if (pause != PAUSED && tick < pause) { @@ -70,7 +66,7 @@ void camu_clock_resume(struct camu_clock *clock, u64 target) } else { clock->offset += tick - clock->paused_at; } - al_atomic_store(f64)(&clock->pause, RUNNING, AL_ATOMIC_RELEASE); + al_atomic_store(f64)(&clock->pause, RUNNING, AL_ATOMIC_RELAXED); clock->paused_at = -1.0; } diff --git a/src/buffer/video.c b/src/buffer/video.c index a7a5594..00902f3 100644 --- a/src/buffer/video.c +++ b/src/buffer/video.c @@ -169,15 +169,15 @@ void camu_video_buffer_flush(struct camu_video_buffer *buf) bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out) { + u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE); if (!buf->single_frame) { + if (flow == SIGNALED || camu_clock_is_ended(buf->clock)) { + buf->callback(buf->userdata, CAMU_BUFFER_EOF); + return false; + } 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_ACQUIRE); - if (flow == SIGNALED && !buf->single_frame) { - buf->callback(buf->userdata, CAMU_BUFFER_EOF); - return false; - } 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); diff --git a/src/cache/handle.c b/src/cache/handle.c index 5bc868c..206738e 100644 --- a/src/cache/handle.c +++ b/src/cache/handle.c @@ -61,7 +61,6 @@ off_t cch_handle_seek(struct cch_handle *handle, off_t offset, s32 whence) } if (filesize <= 0) return -1; if (whence == CAMU_SEEK_SIZE) { - al_log_debug("cache_handle", "seek_size"); return filesize; } switch (whence) { @@ -70,19 +69,16 @@ off_t cch_handle_seek(struct cch_handle *handle, off_t offset, s32 whence) return -1; } handle->pointer = offset; - al_log_debug("cache_handle", "seek_set %li", offset); break; case SEEK_CUR: if (handle->pointer + offset >= filesize || handle->pointer + offset < 0) { return -1; } handle->pointer += offset; - al_log_debug("cache_handle", "seek_cur %li", offset); break; case SEEK_END: handle->prev_pointer = handle->pointer; handle->pointer = filesize + offset; - al_log_debug("cache_handle", "seek_end %li", offset); break; default: return -1; diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c index 3aafd35..375bdfe 100644 --- a/src/codec/ffmpeg/decoder.c +++ b/src/codec/ffmpeg/decoder.c @@ -36,13 +36,12 @@ static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *rend goto err; } - /* s32 cpus = av_cpu_count(); if (cpus > 4) cpus = 2; - av->codec_context->thread_count = cpus; + av->codec_context->thread_count = cpus / 2; + // FF_THREAD_FRAME or FF_THREAD_SLICE. av->codec_context->thread_type = FF_THREAD_SLICE; al_log_debug("ff_decoder", "Using %i threads for decoder.", cpus); - */ #ifndef CAMU_SINK_NO_VIDEO if (codecpar->codec_type == AVMEDIA_TYPE_VIDEO && renderer && renderer->get_buffer2) { diff --git a/src/liana/list.c b/src/liana/list.c index a8d5f71..9698497 100644 --- a/src/liana/list.c +++ b/src/liana/list.c @@ -28,7 +28,7 @@ static void buffer_ahead(struct lia_list *list) struct lia_list_entry *buffered = al_array_at(list->entries, i); struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { - sink->callback(sink->userdata, CAMU_SINK_BUFFER, buffered, i, LIANA_TIMESTAMP_INVALID); + sink->callback(sink->userdata, LIANA_SINK_BUFFER, buffered, i, LIANA_TIMESTAMP_INVALID); } } } @@ -44,6 +44,21 @@ static void unset_all(struct lia_list *list) } } +static bool assume_ended(struct lia_list_entry *entry, u64 at) +{ + if (entry->duration == 0) return true; + if (entry->paused_at == LIANA_TIMESTAMP_INVALID && entry->start != LIANA_TIMESTAMP_INVALID) { + if (entry->offset > entry->duration) { + return true; + } + if (at < entry->start) return false; + if (at - entry->start >= entry->duration - entry->offset) { + return true; + } + } + return false; +} + 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); @@ -58,10 +73,9 @@ void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struc u64 at = LIANA_TIMESTAMP_INVALID; u64 seek_pos = current->offset; if (current->paused_at == LIANA_TIMESTAMP_INVALID) { - now += LIANA_BASE_DELAY; - if (now > current->start && now - current->start > LIANA_BASE_DELAY) { - at = now; - seek_pos += now - current->start; + at = now + LIANA_BASE_DELAY; + if (at > current->start && at - current->start > LIANA_BASE_DELAY) { + seek_pos += at - current->start; } else { at = current->start; } @@ -72,28 +86,16 @@ void lia_list_add_sink(struct lia_list *list, void (*callback)(void *, u8, struc struct lia_timing time = { .at = at, .seek_pos = seek_pos, - .pause = pause + .pause = pause, + .ended = assume_ended(current, now) }; - sink->callback(sink->userdata, CAMU_SINK_SET, current, list->current, &time); + sink->callback(sink->userdata, LIANA_SINK_SET, current, list->current, &time); } else { sink->set = -1; } sink->queued = -1; } -void lia_list_unset(struct lia_list *list) -{ - unset_all(list); - list->current = list->entries.size - 1; - list->previous = -1; - list->queued = -1; - list->idle = true; - struct lia_list_sink *sink; - al_array_foreach(list->sinks, i, sink) { - sink->callback(sink->userdata, CAMU_SINK_CLEAR, NULL, -1, NULL); - } -} - void lia_list_remove_sink(struct lia_list *list, void *userdata) { struct lia_list_sink *sink; @@ -113,6 +115,7 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) entry->paused_at = LIANA_TIMESTAMP_INVALID; entry->held = false; entry->offset = 0; + entry->ended = false; entry->duration = duration; al_str_clone(&entry->name, name); al_array_push(list->entries, entry); @@ -123,13 +126,14 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) struct lia_timing time = { .at = entry->start, .seek_pos = entry->offset, - .pause = LIANA_PAUSE_RESUME + .pause = LIANA_PAUSE_RESUME, + .ended = false }; struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { al_assert(sink->set == -1); sink->set = list->current; - sink->callback(sink->userdata, CAMU_SINK_SET, entry, list->current, &time); + sink->callback(sink->userdata, LIANA_SINK_SET, entry, list->current, &time); } if (list->callback) list->callback(list->userdata, LIANA_META_PLAYING, entry); } else { @@ -146,7 +150,7 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { sink->queued = list->queued; - sink->callback(sink->userdata, CAMU_SINK_BUFFER_AND_QUEUE, entry, list->queued, &time); + sink->callback(sink->userdata, LIANA_SINK_BUFFER_AND_QUEUE, entry, list->queued, &time); } } else { */ @@ -156,6 +160,19 @@ void lia_list_add(struct lia_list *list, void *opaque, u64 duration, str *name) } } +void lia_list_unset(struct lia_list *list) +{ + unset_all(list); + list->current = list->entries.size - 1; + list->previous = -1; + list->queued = -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); + } +} + static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 sequence) { s32 size = (s32)list->entries.size; @@ -163,20 +180,6 @@ static struct lia_list_entry *get_entry_from_sequence(struct lia_list *list, s32 return al_array_at(list->entries, sequence); } -static bool assume_done(struct lia_list_entry *entry, u64 at) -{ - if (entry->duration == 0) return true; - if (entry->paused_at == LIANA_TIMESTAMP_INVALID && entry->start != LIANA_TIMESTAMP_INVALID) { - if (entry->offset > entry->duration) { - return true; - } - if (at < entry->start) return false; - if (at - entry->start >= entry->duration - entry->offset) { - return true; - } - } - return false; -} void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) { @@ -222,29 +225,28 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) al_assert(!current->held); - bool done = assume_done(current, now); - if (done || current->paused_at == LIANA_TIMESTAMP_INVALID) { - if (!done) { + bool ended = assume_ended(current, now); + bool target_ended = assume_ended(target, now); + if (ended || current->paused_at == LIANA_TIMESTAMP_INVALID) { + if (!ended) { current->paused_at = at; current->offset += current->paused_at - current->start; current->start = LIANA_TIMESTAMP_INVALID; current->held = true; } - if (assume_done(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID) { - pause = !done ? LIANA_PAUSE_PAUSE : LIANA_PAUSE_NONE; + if (target_ended || target->paused_at != LIANA_TIMESTAMP_INVALID) { + pause = !ended ? LIANA_PAUSE_PAUSE : LIANA_PAUSE_NONE; } else { //if (target->paused_at == LIANA_TIMESTAMP_INVALID) { target->start = at; - pause = !done ? LIANA_PAUSE_BOTH : LIANA_PAUSE_RESUME; + pause = !ended ? LIANA_PAUSE_BOTH : LIANA_PAUSE_RESUME; } - al_log_debug("list", "%d: not paused, done: %d, pause: %d", sequence, done, pause); } else { - if (assume_done(target, now) || target->paused_at != LIANA_TIMESTAMP_INVALID) { + if (target_ended || target->paused_at != LIANA_TIMESTAMP_INVALID) { pause = LIANA_PAUSE_NONE; } else { //if (target->paused_at == LIANA_TIMESTAMP_INVALID) { target->start = at; pause = LIANA_PAUSE_RESUME; } - al_log_debug("list", "%d: paused, pause: %d", sequence, pause); } list->current = index; @@ -254,13 +256,14 @@ void lia_list_skipto(struct lia_list *list, s32 sequence, s32 index) struct lia_timing time = { .at = at, .seek_pos = target->offset, - .pause = pause + .pause = pause, + .ended = target_ended }; struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { sink->set = index; - sink->callback(sink->userdata, CAMU_SINK_SET, target, index, &time); + sink->callback(sink->userdata, LIANA_SINK_SET, target, index, &time); } list->callback(list->userdata, LIANA_META_PLAYING, target); @@ -299,12 +302,13 @@ void lia_list_toggle_pause(struct lia_list *list, s32 sequence, f64 pts) struct lia_timing time = { .at = at, .seek_pos = LIANA_TIMESTAMP_INVALID, - .pause = pause + .pause = pause, + .ended = false }; struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { - sink->callback(sink->userdata, CAMU_SINK_PAUSE, current, sequence, &time); + sink->callback(sink->userdata, LIANA_SINK_PAUSE, current, sequence, &time); } } @@ -317,14 +321,16 @@ void lia_list_seek(struct lia_list *list, s32 sequence, f64 percent) u64 now = aki_get_timestamp(); current->offset = pos; current->start = now + LIANA_BASE_DELAY; + current->ended = false; struct lia_timing time = { .at = current->start, .seek_pos = current->offset, - .pause = LIANA_PAUSE_NONE + .pause = LIANA_PAUSE_NONE, + .ended = false }; struct lia_list_sink *sink; al_array_foreach(list->sinks, i, sink) { - sink->callback(sink->userdata, CAMU_SINK_SEEK, current, sequence, &time); + sink->callback(sink->userdata, LIANA_SINK_SEEK, current, sequence, &time); } } @@ -335,7 +341,12 @@ void lia_list_end(struct lia_list *list, s32 sequence) s32 size = (s32)list->entries.size; s32 next = sequence + 1; struct lia_list_entry *current = al_array_at(list->entries, list->current); + if (current->ended) { + al_log_warn("list", "Got end() from an already ended resource, ignoring."); + return; + } current->offset = current->duration; + current->ended = true; if (list->queued >= 0) { list->current = list->queued; list->previous = sequence; diff --git a/src/liana/list.h b/src/liana/list.h index eb9a3cb..bbb2ce4 100644 --- a/src/liana/list.h +++ b/src/liana/list.h @@ -7,17 +7,25 @@ #define LIANA_SEQUENCE_ANY -1 #define LIANA_TIMESTAMP_INVALID UINT64_MAX -#define LIANA_BASE_PING 200000Lu // 100ms -#define LIANA_BASE_DELAY 450000Lu // 250ms +#define LIANA_BASE_DELAY 475000Lu // 475ms +#define LIANA_BASE_PING 200000Lu // 200ms #define LIANA_PAUSE_DELAY LIANA_BASE_PING #define LIANA_DELAY_IGNORE 0Lu #define LIANA_BUFFER_AHEAD 2 enum { - LIANA_ENTRY_DURATION = 0, - LIANA_ENTRY_ID, - LIANA_ENTRY_SKIP + LIANA_SINK_SET = 0, + LIANA_SINK_UNSET, + LIANA_SINK_BUFFER, + LIANA_SINK_BUFFER_AND_QUEUE, + LIANA_SINK_PAUSE, + LIANA_SINK_SEEK +}; + +enum { + LIANA_LOAD_ENTRY = 0, + LIANA_UNLOAD_ENTRY }; enum { @@ -41,6 +49,7 @@ struct lia_timing { u64 at; u64 seek_pos; u8 pause; + bool ended; }; struct lia_list_entry { @@ -49,6 +58,7 @@ struct lia_list_entry { u64 paused_at; bool held; u64 offset; + bool ended; u64 duration; str name; }; diff --git a/src/liana/server.c b/src/liana/server.c index 1fc96fb..3eb303e 100644 --- a/src/liana/server.c +++ b/src/liana/server.c @@ -164,7 +164,7 @@ static void signal_callback(void *userdata) } if (!conn->errored) { conn->id = al_inc_u16(); - aki_packet_pool_init(&conn->pool, 48, server->loop, packet_pool_callback, conn); + aki_packet_pool_init(&conn->pool, 96, server->loop, packet_pool_callback, conn); al_array_push(node->connections, conn); handle_connection(conn, packet); } else { diff --git a/src/liana/vcr.c b/src/liana/vcr.c index 5d7b851..192b4f1 100644 --- a/src/liana/vcr.c +++ b/src/liana/vcr.c @@ -3,9 +3,9 @@ #include "vcr.h" #include "handler.h" -#define VCR_BUFFER_INIT 64 -#define VCR_BUFFER_BUFFERED (VCR_BUFFER_INIT - 8) -#define VCR_BUFFER_LOW (VCR_BUFFER_INIT - 32) +#define VCR_BUFFER_INIT 256 +#define VCR_BUFFER_BUFFERED (VCR_BUFFER_INIT - 16) +#define VCR_BUFFER_LOW (VCR_BUFFER_INIT - 64) static void signal_callback(void *userdata) { @@ -47,7 +47,7 @@ static aki_thread_result AKI_THREADCALL vcr_track_thread(void *userdata) } 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) { + if (count <= vcr->mark.low && buffered) { aki_signal_send(&vcr->signal); } } @@ -167,9 +167,8 @@ void lia_vcr_cork(struct lia_vcr_track *track) void lia_vcr_uncork(struct lia_vcr_track *track) { if (al_atomic_load(s32)(&track->state, AL_ATOMIC_RELAXED) != LIANA_STREAM_STOPPED) { - // This can be reached during normal operation. Whether or not that should - // be allowed is up for consideration. - return; + // This can be reached during normal operation. Whether or not that makes + // sense is up for consideration. } al_atomic_store(s32)(&track->state, LIANA_STREAM_RUNNING, AL_ATOMIC_RELAXED); aki_mutex_lock(&track->mutex); diff --git a/src/libsink/common.h b/src/libsink/common.h index 0b1a6c9..c29e494 100644 --- a/src/libsink/common.h +++ b/src/libsink/common.h @@ -1,13 +1,9 @@ #pragma once -#define CAMU_SINK_LOCAL 1 +#define CAMU_SINK_LOCAL 0 enum { - CAMU_SINK_CLEAR = 0, - CAMU_SINK_SET, - CAMU_SINK_BUFFER, - CAMU_SINK_BUFFER_AND_QUEUE, + CAMU_SINK_SET = 0, CAMU_SINK_PAUSE, - CAMU_SINK_SEEK, - CAMU_SINK_DURATION + CAMU_SINK_SEEK }; diff --git a/src/libsink/sink.c b/src/libsink/sink.c index b9554c4..a2262bc 100644 --- a/src/libsink/sink.c +++ b/src/libsink/sink.c @@ -22,7 +22,6 @@ enum { BUFFER_CONFIGURED, BUFFER_SET_OR_BUFFERED, BUFFER_ADDED, - BUFFER_REMOVED }; enum { @@ -39,14 +38,14 @@ enum { #define ENTRY_MAX_AGE 7 +#define BUFFER_EMPTY(buf) ((buf)->state == BUFFER_INIT || (buf)->state == BUFFER_QUEUED) + #ifdef CAMU_SINK_NO_VIDEO -#define ENTRY_VIDEO_READY_OR_EMPTY(entry) true +#define VIDEO_READY_OR_EMPTY(entry) true #else -#define ENTRY_VIDEO_READY_OR_EMPTY(entry) \ - ((entry)->video.state == BUFFER_INIT || (entry)->video.state == BUFFER_QUEUED || (entry)->video.state == BUFFER_ADDED) +#define VIDEO_READY_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->video) || (entry)->video.state == BUFFER_ADDED) #endif -#define ENTRY_AUDIO_READY_OR_EMPTY(entry) \ - ((entry)->audio.state == BUFFER_INIT || (entry)->audio.state == BUFFER_QUEUED || (entry)->audio.state == BUFFER_ADDED) +#define AUDIO_READY_OR_EMPTY(entry) (BUFFER_EMPTY(&(entry)->audio) || (entry)->audio.state == BUFFER_ADDED) #if defined CAMU_SCREEN_THREADED && defined CAMU_MIXER_THREADED #define BLOCKING_SLEEP(delay) aki_thread_sleep(delay) @@ -127,6 +126,7 @@ static void set_or_queue_entry(struct camu_sink_entry *entry) } else { add_audio_if_set_and_buffered(entry); } + #ifndef CAMU_SINK_NO_VIDEO if (entry->video.state == BUFFER_INIT) { entry->video.state = BUFFER_QUEUED; @@ -141,7 +141,7 @@ static void sink_local_pause(struct camu_sink *sink, struct camu_sink_entry *ent { if (camu_clock_is_paused(&entry->clock)) { camu_clock_resume(&entry->clock, 0); - bool no_audio = entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED; + bool no_audio = BUFFER_EMPTY(&entry->audio); if (!no_audio && sink->audio.state == SINK_PAUSED) { sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL); sink->audio.state = SINK_PLAYING; @@ -313,22 +313,23 @@ static void maybe_remove_previous(struct camu_sink *sink) void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) { u8 state = entry->audio.state; - if (state == BUFFER_ADDED || state == BUFFER_SET_OR_BUFFERED) { + al_assert(state != BUFFER_ADDED); + if (state == BUFFER_CONFIGURED) { + state = BUFFER_SET_OR_BUFFERED; + } else if (state == BUFFER_SET_OR_BUFFERED) { + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); + // It's possible for this entry's video buffer to have been added and removed by EOF + // before this point. This needs to be a consideration for keeping sync. + // Same but reversed in add_video_if_set_and_buffered(). + if (VIDEO_READY_OR_EMPTY(entry)) { + maybe_remove_previous(entry->sink); + } camu_audio_buffer_unpause(&entry->audio.buf); queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = START, .value.i = CAMU_SINK_AUDIO }); - } - if (state == BUFFER_SET_OR_BUFFERED) { - bool can_resume = ENTRY_VIDEO_READY_OR_EMPTY(entry); - entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf); - if (can_resume) { - maybe_remove_previous(entry->sink); - } state = BUFFER_ADDED; - } else if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; } entry->audio.state = state; } @@ -337,23 +338,20 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry) void add_video_if_set_and_buffered(struct camu_sink_entry *entry) { u8 state = entry->video.state; - if (state == BUFFER_ADDED || state == BUFFER_SET_OR_BUFFERED) { + al_assert(state != BUFFER_ADDED); + if (state == BUFFER_CONFIGURED) { + state = BUFFER_SET_OR_BUFFERED; + } else if (state == BUFFER_SET_OR_BUFFERED) { + entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); + if (AUDIO_READY_OR_EMPTY(entry)) { + maybe_remove_previous(entry->sink); + } bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); queue_cmd(entry->sink, (struct camu_sink_cmd){ .op = single_frame ? STOP : START, .value.i = CAMU_SINK_VIDEO }); - } - if (state == BUFFER_SET_OR_BUFFERED) { - bool can_resume = ENTRY_AUDIO_READY_OR_EMPTY(entry); - entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf); - if (can_resume) { - maybe_remove_previous(entry->sink); - } - state = BUFFER_ADDED; - } else if (state == BUFFER_CONFIGURED) { - state = BUFFER_SET_OR_BUFFERED; } entry->video.state = state; } @@ -366,7 +364,9 @@ static void audio_buffer_callback(void *userdata, u8 op) switch (op) { case CAMU_BUFFER_BUFFERED: aki_mutex_lock(&sink->mutex); - add_audio_if_set_and_buffered(entry); + if (!entry->ended) { + add_audio_if_set_and_buffered(entry); + } aki_mutex_unlock(&sink->mutex); break; case CAMU_BUFFER_CORK: @@ -388,8 +388,10 @@ static void audio_buffer_callback(void *userdata, u8 op) case CAMU_BUFFER_EOF: { lia_vcr_cork(entry->audio.track); aki_mutex_lock(&sink->mutex); - camu_clock_end(&entry->clock); remove_entry_audio_buffer(sink, entry); + camu_clock_end(&entry->clock); + // EOF on the audio buffer unconditionally + // advances the queue. if (sink->queued) { set_or_queue_entry(sink->queued); sink->current = sink->queued; @@ -411,11 +413,15 @@ static void video_buffer_callback(void *userdata, u8 op) struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; struct camu_sink *sink = entry->sink; switch (op) { - case CAMU_BUFFER_BUFFERED: + case CAMU_BUFFER_BUFFERED: { + bool single_frame = camu_video_buffer_is_single_frame(&entry->video.buf); aki_mutex_lock(&sink->mutex); - add_video_if_set_and_buffered(entry); + if (single_frame || !entry->ended) { + add_video_if_set_and_buffered(entry); + } aki_mutex_unlock(&sink->mutex); break; + } case CAMU_BUFFER_CORK: lia_vcr_cork(entry->video.track); break; @@ -430,7 +436,8 @@ static void video_buffer_callback(void *userdata, u8 op) if (!single_frame) { remove_entry_video_buffer(sink, entry); } - bool run_queue = !single_frame && (entry->audio.state == BUFFER_INIT || entry->audio.state == BUFFER_QUEUED); + // If this entry has no audio, advance the queue. + bool run_queue = !single_frame && BUFFER_EMPTY(&entry->audio); if (run_queue) { camu_clock_end(&entry->clock); if (sink->queued) { @@ -463,10 +470,9 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent { // Entry has both audio and video configured. #ifndef CAMU_SINK_NO_VIDEO - if (entry->audio.state != BUFFER_INIT && entry->audio.state != BUFFER_QUEUED && - entry->video.state != BUFFER_INIT && entry->video.state != BUFFER_QUEUED) { + if (!BUFFER_EMPTY(&entry->audio) && !BUFFER_EMPTY(&entry->video)) { camu_video_buffer_set_latency(&entry->video.buf, -camu_mixer_get_latency(sink->audio.mixer)); - } else if (entry->audio.state != BUFFER_INIT && entry->audio.state != BUFFER_QUEUED) { + } else if (!BUFFER_EMPTY(&entry->audio)) { #if !CAMU_SINK_LOCAL camu_audio_buffer_set_latency(&entry->audio.buf, -camu_mixer_get_latency(sink->audio.mixer)); #endif @@ -698,6 +704,21 @@ static void maybe_cleanup_old_entries(struct camu_sink *sink) } } +static void switch_to(struct camu_sink *sink, struct camu_sink_entry *entry) +{ + if (entry->ended || camu_clock_is_ended(&entry->clock)) { + if (sink->current) { + remove_entry_buffers(sink, sink->current); + } + } else { + if (sink->current) { + al_array_push(sink->previous, sink->current); + } + } + set_or_queue_entry(entry); + sink->current = entry; +} + static void clock_callback(void *userdata, u8 op) { struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata; @@ -705,11 +726,7 @@ static void clock_callback(void *userdata, u8 op) if (op == CAMU_CLOCK_PAUSED) { aki_mutex_lock(&sink->mutex); if (sink->target) { - if (sink->current) { - al_array_push(sink->previous, sink->current); - } - set_or_queue_entry(sink->target); - sink->current = sink->target; + switch_to(sink, sink->target); sink->target = NULL; } aki_mutex_unlock(&sink->mutex); @@ -766,7 +783,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn u8 op = aki_packet_read_u8(packet); #if CAMU_SINK_LOCAL - if (op == CAMU_SINK_CLEAR) { + if (op == LIANA_SINK_UNSET) { if (sink->current) { remove_entry_buffers(sink, sink->current); sink->current = NULL; @@ -774,7 +791,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn goto out; } #else - if (op == CAMU_SINK_CLEAR) { al_assert(false); } // Unimplemented. + if (op == LIANA_SINK_UNSET) { al_assert(false); } // Unimplemented. #endif str addr; @@ -785,6 +802,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn u64 at = aki_packet_read_u64(packet); u64 seek_pos = aki_packet_read_u64(packet); u8 pause = aki_packet_read_u8(packet); + bool ended = aki_packet_read_bool(packet); bool created; struct camu_sink_entry *entry = get_entry_from_id(sink, node_id, &created); @@ -794,6 +812,9 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn if (created) { camu_clock_set(&entry->clock, seek_pos / 1000000.0); + if (ended) { + camu_clock_end(&entry->clock); + } struct camu_renderer *renderer = NULL; #ifndef CAMU_SINK_NO_VIDEO renderer = sink->video.renderer; @@ -801,9 +822,11 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn lia_client_connect(&entry->client, sink->loop, sink->type, &addr, port, node_id, seek_pos, renderer); } - if (op == CAMU_SINK_BUFFER) { + entry->ended = ended; + + if (op == LIANA_SINK_BUFFER) { goto out; - } else if (op == CAMU_SINK_BUFFER_AND_QUEUE) { + } else if (op == LIANA_SINK_BUFFER_AND_QUEUE) { sink->queued = entry; goto out; } @@ -811,14 +834,10 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn #if CAMU_SINK_LOCAL (void)at; (void)pause; - if (sink->current) { - if (!camu_clock_is_paused(&sink->current->clock)) { - camu_clock_pause(&sink->current->clock, 0); - } - al_array_push(sink->previous, sink->current); + if (sink->current && !camu_clock_is_paused(&sink->current->clock)) { + camu_clock_pause(&sink->current->clock, 0); } - set_or_queue_entry(entry); - sink->current = entry; + switch_to(sink, entry); // This will resume a user paused stream. camu_clock_resume(&entry->clock, 0); #else @@ -849,20 +868,10 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn sink->current = entry; break; case LIANA_PAUSE_PAUSE: - if (camu_clock_is_ended(&entry->clock)) { - // Server didn't know this entry was ended, but it is. - queue_cmd(sink, (struct camu_sink_cmd){ - .op = END, - .value.i = entry->sequence - }); - camu_clock_pause(&sink->current->clock, at); - sink->current = entry; - } else if (sink->current) { + if (sink->current) { if (camu_clock_is_ended(&sink->current->clock)) { // Server thought we weren't done, be we are. - al_array_push(sink->previous, sink->current); - set_or_queue_entry(entry); - sink->current = entry; + switch_to(sink, entry); } else { sink->target = entry; camu_clock_pause(&sink->current->clock, at); @@ -882,9 +891,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn camu_clock_pause(&prev_target->clock, at); } else if (sink->current) { if (camu_clock_is_ended(&sink->current->clock)) { - al_array_push(sink->previous, sink->current); - set_or_queue_entry(entry); - sink->current = entry; + switch_to(sink, entry); } else { sink->target = entry; camu_clock_pause(&sink->current->clock, at); @@ -897,7 +904,7 @@ static bool set_command_callback(void *userdata, struct aki_rpc_connection *conn out: aki_mutex_unlock(&sink->mutex); - if (op != CAMU_SINK_BUFFER) { + if (op != LIANA_SINK_BUFFER) { maybe_cleanup_old_entries(sink); } @@ -912,15 +919,17 @@ static bool pause_command_callback(void *userdata, struct aki_rpc_connection *co (void)conn; (void)rpacket; + s32 sequence = aki_packet_read_s32(packet); + u64 at = aki_packet_read_u64(packet); + u8 pause = aki_packet_read_u8(packet); + aki_mutex_lock(&sink->mutex); struct camu_sink_entry *current = sink->current; if (!current) goto out; + #if CAMU_SINK_LOCAL sink_local_pause(sink, current); #else - s32 sequence = aki_packet_read_s32(packet); - u64 at = aki_packet_read_u64(packet); - u8 pause = aki_packet_read_u8(packet); if (current->sequence == sequence) { if (pause == LIANA_PAUSE_PAUSE) { camu_clock_pause(&sink->current->clock, at); @@ -957,7 +966,10 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con aki_mutex_lock(&sink->mutex); struct camu_sink_entry *current = sink->current; aki_mutex_unlock(&sink->mutex); + if (!current) goto out; + if (current->sequence == sequence) { + current->ended = false; #if CAMU_SINK_LOCAL (void)at; lia_client_seek(¤t->client, pos, 0); @@ -966,6 +978,7 @@ static bool seek_command_callback(void *userdata, struct aki_rpc_connection *con #endif } +out: aki_packet_free(packet); return false; } diff --git a/src/libsink/sink.h b/src/libsink/sink.h index 6d560f6..b830230 100644 --- a/src/libsink/sink.h +++ b/src/libsink/sink.h @@ -41,6 +41,7 @@ struct camu_sink_entry { u16 lru; struct camu_clock clock; struct lia_client client; + bool ended; struct { u8 state; bool armed; diff --git a/src/mixer/mixer.c b/src/mixer/mixer.c index 987bc8b..e6fd178 100644 --- a/src/mixer/mixer.c +++ b/src/mixer/mixer.c @@ -30,12 +30,20 @@ static s32 data_callback(void *userdata, u8 *data, s32 frame_count, bool *silenc } else { *silence = true; } + struct camu_audio_buffer *buf; + // Read from all current buffers on each callback to + // emulate the behavior of an actual mixer. + al_array_foreach(mixer->buffers, i, buf) { + if (i < mixer->buffers.size - 1) { + camu_audio_buffer_read(buf, data, req); + } + } do { if (al_array_is_empty(mixer->buffers)) { al_memset(data, 0, req); } else { size_t signal; - struct camu_audio_buffer *buf = al_array_last(mixer->buffers); + buf = al_array_last(mixer->buffers); *silence = false; if ((signal = camu_audio_buffer_read(buf, data, req)) < req) { #ifdef CAMU_MIXER_THREADED diff --git a/src/render/renderer_libplacebo.c b/src/render/renderer_libplacebo.c index eeef565..8dde7f5 100644 --- a/src/render/renderer_libplacebo.c +++ b/src/render/renderer_libplacebo.c @@ -204,27 +204,33 @@ static void renderer_lp_render(struct camu_renderer *renderer, struct camu_scree struct camu_screen_video *video; while (scr->videos.size > 0) { - video = &al_array_last(scr->videos); - if (camu_video_buffer_read(video->buf, &mix) && mix.frames) { - target.crop = mix.frames[0]->crop; - target.crop.x1 *= video->view.zoom / video->view.stretch; - target.crop.y1 *= video->view.zoom * video->view.stretch; - target.crop.x0 += video->view.x_offset; - target.crop.y0 += video->view.y_offset; - target.crop.x1 += video->view.x_offset; - target.crop.y1 += video->view.y_offset; - target.rotation = video->view.rotation; - //lr->params.color_adjustment = pl_color_adjustment( - // .saturation = 0.0 - //); - pl_render_image_mix(lr->renderer, &mix, &target, &lr->params); - break; + bool any_eof = false; + al_array_foreach_ptr(scr->videos, i, video) { + if (camu_video_buffer_read(video->buf, &mix) && mix.frames) { + target.crop = mix.frames[0]->crop; + target.crop.x1 *= video->view.zoom / video->view.stretch; + target.crop.y1 *= video->view.zoom * video->view.stretch; + target.crop.x0 += video->view.x_offset; + target.crop.y0 += video->view.y_offset; + target.crop.x1 += video->view.x_offset; + target.crop.y1 += video->view.y_offset; + target.rotation = video->view.rotation; + //lr->params.color_adjustment = pl_color_adjustment( + // .saturation = 0.0 + //); + pl_render_image_mix(lr->renderer, &mix, &target, &lr->params); + } else { + any_eof = true; + } + } + if (any_eof) { #ifdef CAMU_SCREEN_THREADED - } else { if (al_atomic_load(u8)(&scr->queued, AL_ATOMIC_RELAXED)) { camu_screen_run_queue(scr); } #endif + } else { + break; } } diff --git a/src/server/list.c b/src/server/list.c index 8ff123c..e3565d0 100644 --- a/src/server/list.c +++ b/src/server/list.c @@ -1,4 +1,62 @@ +#include "../libsink/common.h" +#include "../server/common.h" + #include "list.h" +#include "server.h" + +void camu_list_callback(void *userdata, u8 op, struct lia_list_entry *entry) +{ + struct camu_server *server = (struct camu_server *)userdata; + (void)server; + (void)op; + (void)entry; +} + +void camu_list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) +{ + struct camu_server_sink *sink = (struct camu_server_sink *)userdata; + switch (op) { + case LIANA_SINK_SET: + case LIANA_SINK_BUFFER: + case LIANA_SINK_BUFFER_AND_QUEUE: { + struct camu_server_resource *resource = (struct camu_server_resource *)entry->opaque; + struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); + aki_packet_write_u8(packet, op); + aki_packet_write_str(packet, &sink->server->addr); + aki_packet_write_u16(packet, CAMU_PORT); + aki_packet_write_u16(packet, resource->node->id); + aki_packet_write_s32(packet, sequence); + aki_packet_write_u64(packet, timing->at); + aki_packet_write_u64(packet, timing->seek_pos); + aki_packet_write_u8(packet, timing->pause); + aki_packet_write_bool(packet, timing->ended); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + case LIANA_SINK_UNSET: { + struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); + aki_packet_write_u8(packet, op); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + case LIANA_SINK_PAUSE: { + struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE); + aki_packet_write_s32(packet, sequence); + aki_packet_write_u64(packet, timing->at); + aki_packet_write_u8(packet, timing->pause); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + case LIANA_SINK_SEEK: { + struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK); + aki_packet_write_s32(packet, sequence); + aki_packet_write_u64(packet, timing->at); + aki_packet_write_u64(packet, timing->seek_pos); + aki_rpc_connection_command(sink->conn, packet, NULL, NULL); + break; + } + } +} void camu_list_init(struct camu_list *list, str *name) { diff --git a/src/server/list.h b/src/server/list.h index f3e52e6..e472024 100644 --- a/src/server/list.h +++ b/src/server/list.h @@ -2,10 +2,15 @@ #include "../liana/list.h" +#include "resource.h" + struct camu_list { str name; struct lia_list impl; }; +void camu_list_callback(void *userdata, u8 op, struct lia_list_entry *entry); +void camu_list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing); + void camu_list_init(struct camu_list *list, str *name); void camu_list_free(struct camu_list *list); diff --git a/src/server/resource.h b/src/server/resource.h new file mode 100644 index 0000000..c1cf619 --- /dev/null +++ b/src/server/resource.h @@ -0,0 +1,13 @@ +#pragma once + +#include "../cache/entry.h" +#include "../liana/server.h" +#include "../portal/src/post.h" + +struct camu_server_resource { + str unique_id; + struct camu_post *post; + struct cch_entry *entry; + struct lia_node *node; + u64 duration; +}; diff --git a/src/server/server.c b/src/server/server.c index 49ea234..d7f8bc9 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -1,8 +1,6 @@ #include #include -#include "../libsink/common.h" - #include "server.h" #include "common.h" #include "list.h" @@ -19,54 +17,31 @@ static struct camu_user *get_user_by_username(struct camu_server *server, str *u return NULL; } -static void list_sink_callback(void *userdata, u8 op, struct lia_list_entry *entry, s32 sequence, struct lia_timing *timing) +static struct camu_server_client *get_client_by_connection(struct camu_server *server, struct aki_rpc_connection *conn) { - struct camu_server_sink *sink = (struct camu_server_sink *)userdata; - switch (op) { - case CAMU_SINK_CLEAR: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); - aki_packet_write_u8(packet, op); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - break; - } - case CAMU_SINK_SET: - case CAMU_SINK_BUFFER: - case CAMU_SINK_BUFFER_AND_QUEUE: { - struct camu_server_resource *resource = (struct camu_server_resource *)entry->opaque; - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SET); - aki_packet_write_u8(packet, op); - aki_packet_write_str(packet, &sink->server->addr); - aki_packet_write_u16(packet, CAMU_PORT); - aki_packet_write_u16(packet, resource->node->id); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); - aki_packet_write_u64(packet, timing->seek_pos); - aki_packet_write_u8(packet, timing->pause); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - break; - } - case CAMU_SINK_PAUSE: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_PAUSE); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); - aki_packet_write_u8(packet, timing->pause); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - break; - } - case CAMU_SINK_SEEK: { - struct aki_packet *packet = aki_rpc_get_packet(sink->conn->rpc, CAMU_SINK_SEEK); - aki_packet_write_s32(packet, sequence); - aki_packet_write_u64(packet, timing->at); - aki_packet_write_u64(packet, timing->seek_pos); - aki_rpc_connection_command(sink->conn, packet, NULL, NULL); - break; + struct camu_server_client *client; + al_array_foreach(server->clients, i, client) { + if (client->conn == conn) return client; } - case CAMU_SINK_DURATION: { - struct camu_server_resource *resource = (struct camu_server_resource *)entry->opaque; - (void)resource; - break; + return NULL; +} + +static struct camu_list *get_list_from_name(struct camu_server *server, str *name) +{ + struct camu_list *list; + al_array_foreach(server->lists, i, list) { + if (al_str_eq(&list->name, name)) return list; } + return NULL; +} + +static struct camu_server_sink *get_sink_from_name(struct camu_server *server, str *name) +{ + struct camu_server_sink *sink; + al_array_foreach(server->sinks, i, sink) { + if (al_str_eq(&sink->name, name)) return sink; } + return NULL; } static bool identify_command_callback(void *userdata, struct aki_rpc_connection *conn, @@ -109,7 +84,7 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection sink->server = server; al_array_push(server->sinks, sink); struct camu_list *list = al_array_at(server->lists, 0); - lia_list_add_sink(&list->impl, list_sink_callback, sink); + lia_list_add_sink(&list->impl, camu_list_sink_callback, sink); al_log_info("server", "New sink."); break; } @@ -119,33 +94,6 @@ static bool identify_command_callback(void *userdata, struct aki_rpc_connection return true; } -static struct camu_server_client *get_client_by_connection(struct camu_server *server, struct aki_rpc_connection *conn) -{ - struct camu_server_client *client; - al_array_foreach(server->clients, i, client) { - if (client->conn == conn) return client; - } - return NULL; -} - -static struct camu_list *get_list_from_name(struct camu_server *server, str *name) -{ - struct camu_list *list; - al_array_foreach(server->lists, i, list) { - if (al_str_eq(&list->name, name)) return list; - } - return NULL; -} - -static struct camu_server_sink *get_sink_from_name(struct camu_server *server, str *name) -{ - struct camu_server_sink *sink; - al_array_foreach(server->sinks, i, sink) { - if (al_str_eq(&sink->name, name)) return sink; - } - return NULL; -} - static bool client_command_command_callback(void *userdata, struct aki_rpc_connection *conn, struct aki_packet *packet, struct aki_packet *rpacket) { @@ -171,7 +119,7 @@ static bool client_command_command_callback(void *userdata, struct aki_rpc_conne struct camu_list *list = get_list_from_name(server, &name); aki_packet_read_str(packet, &name); struct camu_server_sink *sink = get_sink_from_name(server, &name); - lia_list_add_sink(&list->impl, list_sink_callback, sink); + lia_list_add_sink(&list->impl, camu_list_sink_callback, sink); break; } } diff --git a/src/server/server.h b/src/server/server.h index 528f08a..6246f79 100644 --- a/src/server/server.h +++ b/src/server/server.h @@ -7,6 +7,7 @@ #include "../liana/server.h" #include "user.h" +#include "resource.h" #include "local_compat.h" struct camu_server_node { @@ -24,14 +25,6 @@ struct camu_server_sink { struct camu_server *server; }; -struct camu_server_resource { - str unique_id; - struct camu_post *post; - struct cch_entry *entry; - struct lia_node *node; - u64 duration; -}; - struct camu_server { struct aki_event_loop *loop; str addr; -- cgit v1.2.3-101-g0448