summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/buffer/audio.c25
-rw-r--r--src/buffer/clock.c10
-rw-r--r--src/buffer/video.c10
-rw-r--r--src/cache/handle.c4
-rw-r--r--src/codec/ffmpeg/decoder.c5
-rw-r--r--src/liana/list.c115
-rw-r--r--src/liana/list.h20
-rw-r--r--src/liana/server.c2
-rw-r--r--src/liana/vcr.c13
-rw-r--r--src/libsink/common.h10
-rw-r--r--src/libsink/sink.c155
-rw-r--r--src/libsink/sink.h1
-rw-r--r--src/mixer/mixer.c10
-rw-r--r--src/render/renderer_libplacebo.c38
-rw-r--r--src/server/list.c58
-rw-r--r--src/server/list.h5
-rw-r--r--src/server/resource.h13
-rw-r--r--src/server/server.c98
-rw-r--r--src/server/server.h9
19 files changed, 324 insertions, 277 deletions
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(&current->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 <al/log.h>
#include <al/lib.h>
-#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;