summaryrefslogtreecommitdiff
path: root/src
diff options
context:
space:
mode:
Diffstat (limited to 'src')
-rw-r--r--src/buffer/audio.c158
-rw-r--r--src/buffer/audio.h10
-rw-r--r--src/buffer/clock.c8
-rw-r--r--src/buffer/clock.h2
-rw-r--r--src/buffer/video.c4
-rw-r--r--src/buffer/video.h2
-rw-r--r--src/cache/handle.c2
-rw-r--r--src/codec/ffmpeg/decoder.c2
-rw-r--r--src/codec/ffmpeg/demuxer.c4
-rw-r--r--src/fruits/cmsrv/cmsrv.c51
-rw-r--r--src/liana/client.c4
-rw-r--r--src/liana/list.c2
-rw-r--r--src/liana/server.c5
-rw-r--r--src/liana/vcr.c6
-rw-r--r--src/libsink/sink.c285
-rw-r--r--src/libsink/sink.h2
-rw-r--r--src/mixer/mixer.c18
-rw-r--r--src/mixer/mixer.h8
-rw-r--r--src/render/renderer_libplacebo.c2
-rw-r--r--src/server/server.c2
-rw-r--r--src/sink/common.c81
-rw-r--r--src/sink/common.h84
-rw-r--r--src/sink/desktop.c2
-rw-r--r--src/sink/meson.build2
-rw-r--r--src/util/queue.h12
25 files changed, 477 insertions, 281 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index d438d25..707219f 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -20,7 +20,7 @@
enum {
PAUSE_PAUSED = 0,
#ifdef CAMU_AUDIO_BUFFER_FADE
- PAUSE_FADING,
+ PAUSE_FADING_OUT,
#endif
PAUSE_PLAYING
};
@@ -33,7 +33,7 @@ static void reset_buffer_state(struct camu_audio_buffer *buf)
al_atomic_store(s32)(&buf->volume.set, 0, AL_ATOMIC_RELAXED);
#ifdef CAMU_AUDIO_BUFFER_FADE
buf->fade.offset = 0;
- buf->fade.volume = 1.f;
+ buf->fade.volume = -1.f;
#endif
buf->buffered = false;
al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
@@ -43,9 +43,10 @@ static void reset_buffer_state(struct camu_audio_buffer *buf)
bool camu_audio_buffer_init(struct camu_audio_buffer *buf, struct camu_clock *clock)
{
buf->clock = clock;
- reset_buffer_state(buf);
- buf->ignore_desync = false;
buf->latency = 0.0;
+ buf->ignore_desync = false;
+ al_atomic_store(bool)(&buf->no_video, false, AL_ATOMIC_RELAXED);
+ reset_buffer_state(buf);
#ifdef CAMU_MIXER_THREADED
al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
#endif
@@ -82,7 +83,7 @@ bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_code
buf->mark.min = camu_audio_format_sec_to_bytes(&buf->fmt.req, BUFFER_MARK_MIN);
buf->mark.buffered = camu_audio_format_sec_to_bytes(&buf->fmt.req, BUFFER_MARK_BUFFERED);
- camu_peak_buffer_init(&buf->peak, 1024 * 16);
+ camu_peak_buffer_init(&buf->peak, KB(16));
buf->stream = stream;
@@ -100,6 +101,16 @@ void camu_audio_buffer_set_latency(struct camu_audio_buffer *buf, f64 latency)
buf->latency = latency;
}
+void camu_audio_buffer_set_ignore_desync(struct camu_audio_buffer *buf, bool ignore_desync)
+{
+ buf->ignore_desync = ignore_desync;
+}
+
+void camu_audio_buffer_set_no_video(struct camu_audio_buffer *buf, bool no_video)
+{
+ al_atomic_store(bool)(&buf->no_video, no_video, AL_ATOMIC_RELAXED);
+}
+
#ifdef _DEBUG_
#define BUFFERED_SECONDS_DEBUG(buf) \
camu_audio_format_bytes_to_sec(&buf->fmt.req, al_ring_buffer_occupied(&buf->rb))
@@ -113,26 +124,37 @@ static inline bool frame_is_late(struct camu_clock *clock, f64 base, f64 pts, f6
return pts + duration < base;
}
-// return value of false indicates we pushed to the peak buffer.
-static bool push_internal(struct camu_audio_buffer *buf, f64 pts, u8 **data, s32 sample_count, bool flush)
+// A return value of false signals that we pushed to the peak buffer.
+static bool push_internal(struct camu_audio_buffer *buf, f64 pts, u8 **data, s32 sample_count)
{
- if (!flush) {
+ f64 base_pts = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
+
+ if (sample_count > 0) {
+ al_assert(data);
f64 duration = camu_audio_format_samples_to_sec(&buf->fmt.in, sample_count);
- f64 base = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
- if (frame_is_late(buf->clock, base, pts, duration)) return true;
+ if (frame_is_late(buf->clock, base_pts, pts, duration)) {
+ return true;
+ }
if (buf->fmt.resampler_needed) {
sample_count = buf->resampler->convert(buf->resampler, (const u8 **)data, sample_count);
data = buf->resampler->get_data(buf->resampler);
}
- if (sample_count == 0) return true;
- if (base == -1.0) al_atomic_store(f64)(&buf->pts, pts, AL_ATOMIC_RELEASE);
} else {
- al_assert(!data && sample_count == 0);
+ al_assert(!data);
if (buf->fmt.resampler_needed) {
sample_count = buf->resampler->flush(buf->resampler);
data = buf->resampler->get_data(buf->resampler);
}
- if (sample_count == 0) return true;
+ }
+
+ if (base_pts == -1.0) {
+ // We want this push() to set the pts even if sample_count = 0.
+ al_atomic_store(f64)(&buf->pts, pts, AL_ATOMIC_RELEASE);
+ }
+
+ if (sample_count == 0) {
+ // The resampler is holding all the data.
+ return true;
}
size_t space = al_ring_buffer_space(&buf->rb);
@@ -177,7 +199,7 @@ static void push_av_frame_internal(struct camu_audio_buffer *buf, AVFrame *frame
s32 sample_count = frame->nb_samples;
AVStream *stream = buf->stream->av.stream;
f64 pts = frame->best_effort_timestamp * av_q2d(stream->time_base);
- push_internal(buf, pts, frame->data, sample_count, false);
+ push_internal(buf, pts, frame->data, sample_count);
av_frame_free(&frame);
}
#endif
@@ -186,13 +208,13 @@ static void push_av_frame_internal(struct camu_audio_buffer *buf, AVFrame *frame
void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_frame *frame)
{
al_assert(al_atomic_load(u8)(&buf->flow, AL_ATOMIC_RELAXED) == FLOWING);
+
switch (frame->mode) {
case CAMU_NORMAL: {
s32 sample_count = frame->audio.sample_count;
- f64 pts = frame->pts;
u8 *planes[CAMU_PLANAR_DATA_POINTERS] = { 0 };
planes[0] = frame->data;
- push_internal(buf, pts, planes, sample_count, false);
+ push_internal(buf, frame->pts, planes, sample_count);
break;
}
#ifdef CAMU_HAVE_FFMPEG
@@ -202,6 +224,7 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_fra
}
#endif
}
+
al_free(frame);
}
@@ -209,14 +232,17 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_fra
void camu_audio_buffer_flush(struct camu_audio_buffer *buf)
{
al_log_debug("audio_buffer", "Flush requested.");
- if (!push_internal(buf, 0.0, NULL, 0, true)) {
+
+ if (!push_internal(buf, 0.0, NULL, 0)) {
al_log_debug("audio_buffer", "Buffer filled by flush.");
}
+
if (!buf->buffered) {
al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", BUFFERED_SECONDS_DEBUG(buf));
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
buf->buffered = true;
}
+
al_atomic_store(u8)(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
}
@@ -236,14 +262,6 @@ void camu_audio_buffer_unpause(struct camu_audio_buffer *buf)
al_atomic_add(s32)(&buf->unpause, 1, AL_ATOMIC_RELAXED);
}
-static void increment_pts(struct camu_audio_buffer *buf, f64 amount, f64 *base)
-{
- *base += amount;
- al_atomic_store(f64)(&buf->pts, *base, AL_ATOMIC_RELEASE);
-}
-
-#define NOT_PAUSED(pause) (pause != PAUSE_PAUSED)
-
// PAUSE_PAUSED signifies that the last read was silence. Meaning we can skip
// around in the buffer without worrying about pops.
size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t req)
@@ -258,25 +276,31 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
al_atomic_sub(s32)(&buf->volume.set, 1, AL_ATOMIC_RELEASE);
}
- if (camu_clock_is_paused(buf->clock)) {
+ bool set_clock = al_atomic_load(bool)(&buf->no_video, AL_ATOMIC_RELAXED);
+ f64 pts = camu_clock_get_pts(buf->clock, buf->latency, set_clock);
+ if (pts == -1.0) {
#ifdef CAMU_AUDIO_BUFFER_FADE
if (buf->pause == PAUSE_PLAYING) {
- buf->pause = PAUSE_FADING;
- buf->fade.offset = 0;
- } else if (buf->fade.volume == 0.f || buf->pause == PAUSE_PAUSED) {
+ buf->pause = PAUSE_FADING_OUT;
buf->fade.offset = 0;
+ } else if (buf->pause == PAUSE_PAUSED || buf->fade.volume == 0.f) {
al_memset(data, 0, req);
- if (NOT_PAUSED(buf->pause)) {
+ if (buf->pause != PAUSE_PAUSED) {
+ // Fade out finished.
buf->pause = PAUSE_PAUSED;
buf->callback(buf->userdata, CAMU_BUFFER_PAUSED);
}
+ buf->fade.offset = 0;
return req;
}
- } else if (buf->pause == PAUSE_FADING || buf->pause == PAUSE_PAUSED) {
- if (buf->pause == PAUSE_FADING) buf->pause = PAUSE_PLAYING;
-#else
+ } else if (buf->pause == PAUSE_FADING_OUT || buf->pause == PAUSE_PAUSED) {
+ if (buf->pause == PAUSE_FADING_OUT) {
+ // If unpaused during a fade out, resume immediately to not break the stream.
+ buf->pause = PAUSE_PLAYING;
+ }
+#else // No fade and clock paused.
al_memset(data, 0, req);
- if (NOT_PAUSED(buf->pause)) {
+ if (buf->pause != PAUSE_PAUSED) {
buf->pause = PAUSE_PAUSED;
buf->callback(buf->userdata, CAMU_BUFFER_PAUSED);
}
@@ -285,36 +309,36 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
}
if (al_atomic_load(s32)(&buf->unpause, AL_ATOMIC_ACQUIRE) > 0) {
- // If pause != PLAYING here, this will cause unexpected behavior.
+ // If pause != PLAYING here this will likely cause unexpected behavior.
buf->pause = PAUSE_PAUSED;
al_atomic_sub(s32)(&buf->unpause, 1, AL_ATOMIC_RELEASE);
}
- f64 base = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
- f64 pts = camu_clock_get_pts(buf->clock, buf->latency, false);
+ struct camu_audio_format *fmt = &buf->fmt.req;
size_t ret, signal = req;
size_t have = al_ring_buffer_occupied(&buf->rb);
#ifdef CAMU_AUDIO_BUFFER_FADE
// Cut off fade if it's reaching too far.
if (have < buf->fade.offset) have = 0;
else have -= buf->fade.offset;
- if (buf->pause != PAUSE_FADING && buf->pause != PAUSE_PLAYING) {
#endif
- pts -= base;
- // 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)) {
+
+ f64 base_pts = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
+ // Unpause and possibly attempt syncing to the clock.
+ if (UNLIKELY(buf->pause == PAUSE_PAUSED)) {
+ if (!buf->ignore_desync) {
+ pts -= base_pts;
if (pts > 0.0) { // Skip.
- ret = camu_audio_format_sec_to_bytes(&buf->fmt.req, pts);
+ ret = camu_audio_format_sec_to_bytes(fmt, pts);
ret = MIN(ret, have);
al_log_info("audio_buffer", "Skipping %fs of audio (%zu bytes).", pts, ret);
ret = al_ring_buffer_discard(&buf->rb, ret);
have -= ret;
- increment_pts(buf, camu_audio_format_bytes_to_sec(&buf->fmt.req, ret), &base);
+ base_pts += camu_audio_format_bytes_to_sec(fmt, ret);
// Could go on to underrun.
} else if (pts < 0.0) { // Delay.
pts = -pts;
- ret = camu_audio_format_sec_to_bytes(&buf->fmt.req, pts);
+ ret = camu_audio_format_sec_to_bytes(fmt, pts);
ret = MIN(ret, req);
al_log_info("audio_buffer", "Delaying audio by %fs.", pts);
al_memset(data, 0, ret);
@@ -324,13 +348,11 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
if (req == 0) return signal;
}
}
- if (buf->pause != PAUSE_PLAYING) buf->pause = PAUSE_PLAYING;
-#ifdef CAMU_AUDIO_BUFFER_FADE
+ buf->pause = PAUSE_PLAYING;
}
-#endif
u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE);
- if (have < req) { // We don't have enough data to fulfill our request.
+ if (UNLIKELY(have < req)) { // We don't have enough data to fulfill our request.
if (flow == FLUSHED) { // Stream is flushed.
// Check peak buffer for any remaining data.
ret = camu_peak_buffer_get_size(&buf->peak);
@@ -340,17 +362,16 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
al_ring_buffer_write(&buf->rb, camu_peak_buffer_flush(&buf->peak), ret);
have += ret;
}
- // Only signal EOF if what, if any, data from the peak buffer
- // still isn't enough for this request.
+ // Only signal EOF if the data from the peak buffer, if any, wasn't
+ // enough for this request.
if (have < req) {
signal = have;
al_log_debug("audio_buffer", "Flushed (signal: %zu).", signal);
buf->callback(buf->userdata, CAMU_BUFFER_EOF);
- flow = SIGNALED;
al_atomic_store(u8)(&buf->flow, SIGNALED, AL_ATOMIC_RELEASE);
}
} else {
- // Silence remainder of the request.
+ // Silence the remainder of the request.
al_memset(data + have, 0, req - have);
if (flow == FLOWING) {
// We hit an underrun because data wasn't coming in fast enough.
@@ -368,11 +389,11 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
req = MIN(req, have);
}
- if (req > 0) {
+ if (LIKELY(req > 0)) {
#ifdef CAMU_AUDIO_BUFFER_FADE
- if (buf->pause == PAUSE_FADING) {
- // Peek into the ring buffer. This won't consume the data, so when
- // we resume it will fade in from the same position the fade out started.
+ if (buf->pause == PAUSE_FADING_OUT) {
+ // Peeking into the ring buffer won't consume data. So, when resumed we
+ // will fade in from the same position that the fade out started.
ret = al_ring_buffer_peek(&buf->rb, data, buf->fade.offset, req);
buf->fade.offset += ret;
} else {
@@ -380,32 +401,37 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
// Read as much as we determined we can.
ret = al_ring_buffer_read(&buf->rb, data, req);
al_assert(ret == req);
- increment_pts(buf, camu_audio_format_bytes_to_sec(&buf->fmt.req, ret), &base);
+ base_pts += camu_audio_format_bytes_to_sec(fmt, ret);
#ifdef CAMU_AUDIO_BUFFER_FADE
}
+
f32 step = 0.f;
- if (buf->pause == PAUSE_FADING || buf->fade.volume > buf->volume.user) {
- step = -FADE_STEP(&buf->fmt.req);
+ if (buf->pause == PAUSE_FADING_OUT || buf->fade.volume > buf->volume.user) {
+ step = -FADE_STEP(fmt);
} else if (buf->fade.volume < buf->volume.user) {
- step = FADE_STEP(&buf->fmt.req);
+ step = FADE_STEP(fmt);
}
+
if (step != 0.f || buf->fade.volume != 1.f) {
if (buf->volume.user > 1.f) {
- // Only adjust step if it's an increase.
+ // Only adjust the step if it's an increase to avoid
+ // drawn-out fades if the volume is low.
step *= buf->volume.user;
}
- buf->fade.volume = apply_volume(data, req, &buf->fmt.req, buf->fade.volume, buf->volume.user, step);
+ buf->fade.volume = apply_volume(data, req, fmt, buf->fade.volume, buf->volume.user, step);
}
#else
if (buf->volume.user != 1.f) {
- apply_volume(data, req, &buf->fmt.req, buf->volume.user, 0.f, 0.f);
+ apply_volume(data, req, fmt, buf->volume.user, 0.f, 0.f);
}
#endif
}
+ al_atomic_store(f64)(&buf->pts, base_pts, AL_ATOMIC_RELEASE);
+
// Check if we should request to uncork.
#ifdef CAMU_AUDIO_BUFFER_FADE
- if (flow == FLOWING && buf->pause != PAUSE_FADING) {
+ if (flow == FLOWING && buf->pause != PAUSE_FADING_OUT) {
#else
if (flow == FLOWING) {
#endif
diff --git a/src/buffer/audio.h b/src/buffer/audio.h
index 6fc3cca..ca2d405 100644
--- a/src/buffer/audio.h
+++ b/src/buffer/audio.h
@@ -14,12 +14,14 @@
struct camu_audio_buffer {
struct camu_codec_stream *stream;
+ struct camu_clock *clock;
atomic(f64) pts;
+ f64 latency;
+ bool ignore_desync;
+ atomic(bool) no_video;
+
u8 pause;
atomic(s32) unpause;
- struct camu_clock *clock;
- bool ignore_desync;
- f64 latency;
struct camu_resampler_format fmt;
struct camu_resampler *resampler;
@@ -61,6 +63,8 @@ bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_code
struct camu_mixer *mixer);
void camu_audio_buffer_set_volume(struct camu_audio_buffer *buf, f32 volume);
void camu_audio_buffer_set_latency(struct camu_audio_buffer *buf, f64 latency);
+void camu_audio_buffer_set_ignore_desync(struct camu_audio_buffer *buf, bool ignore_desync);
+void camu_audio_buffer_set_no_video(struct camu_audio_buffer *buf, bool no_video);
void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_frame *frame);
void camu_audio_buffer_flush(struct camu_audio_buffer *buf);
void camu_audio_buffer_reset(struct camu_audio_buffer *buf);
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
index b711a27..3e3efc0 100644
--- a/src/buffer/clock.c
+++ b/src/buffer/clock.c
@@ -108,7 +108,7 @@ f64 camu_clock_get_base_pts(struct camu_clock *clock)
return clock->base;
}
-f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool skip_set)
+f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool allow_set)
{
f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE);
if (pause == PAUSED) return -1.0;
@@ -117,11 +117,11 @@ f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool skip_set)
f64 tick = al_atomic_load(f64)(&clock->tick, AL_ATOMIC_RELAXED);
if (tick == -1.0) {
- if (skip_set) {
- return -1.0;
- } else {
+ if (allow_set) {
tick = al_atomic_compare_and_swap(f64)(&clock->tick, -1.0, current);
if (tick == -1.0) tick = current;
+ } else {
+ return -1.0;
}
}
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
index edad287..a2762a6 100644
--- a/src/buffer/clock.h
+++ b/src/buffer/clock.h
@@ -48,4 +48,4 @@ bool camu_clock_is_armed(struct camu_clock *clock);
bool camu_clock_is_paused(struct camu_clock *clock);
f64 camu_clock_get_base_pts(struct camu_clock *clock);
-f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool skip_set);
+f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset, bool allow_set);
diff --git a/src/buffer/video.c b/src/buffer/video.c
index e2b6f8e..47c1ea9 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -28,7 +28,7 @@ bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *cl
buf->avg_frame_duration = 0.0;
buf->queue = NULL;
buf->buffered = false;
- buf->weighted_read = false;
+ buf->weighted_read = true;
al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
#ifdef CAMU_SCREEN_THREADED
al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
@@ -242,7 +242,7 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weig
{
f64 base = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
if (!buf->single_frame) {
- f64 pts = camu_clock_get_pts(buf->clock, buf->latency, buf->weighted_read);
+ f64 pts = camu_clock_get_pts(buf->clock, buf->latency, !buf->weighted_read);
if (pts > base) {
base = pts;
al_atomic_store(f64)(&buf->pts, base, AL_ATOMIC_RELEASE);
diff --git a/src/buffer/video.h b/src/buffer/video.h
index 5313f75..632a28f 100644
--- a/src/buffer/video.h
+++ b/src/buffer/video.h
@@ -15,8 +15,8 @@ struct camu_video_buffer {
struct camu_codec_stream *stream;
struct camu_clock *clock;
- f64 latency;
atomic(f64) pts;
+ f64 latency;
f64 last_pts;
bool single_frame;
diff --git a/src/cache/handle.c b/src/cache/handle.c
index dddf172..98bc2fa 100644
--- a/src/cache/handle.c
+++ b/src/cache/handle.c
@@ -38,6 +38,7 @@ s32 cch_handle_read(struct cch_handle *handle, u8 *buf, s32 size)
return CAMU_ERR_EOF;
}
struct cch_backing *backing = handle->entry->backing;
+ al_assert(size >= 0);
size_t available = (size_t)size;
if (backing->mode == CACHE_BACKING_MAPPED) {
u8 *ptr = backing->get_ptr(backing, handle->pointer, &available);
@@ -49,6 +50,7 @@ s32 cch_handle_read(struct cch_handle *handle, u8 *buf, s32 size)
backing->read(backing, buf, handle->pointer, &available);
}
handle->pointer += available;
+ al_assert(available <= INT32_MAX);
return available > 0 ? (s32)available : CAMU_ERR_EOF;
}
diff --git a/src/codec/ffmpeg/decoder.c b/src/codec/ffmpeg/decoder.c
index ce974f3..c2d95ce 100644
--- a/src/codec/ffmpeg/decoder.c
+++ b/src/codec/ffmpeg/decoder.c
@@ -58,9 +58,11 @@ static bool ff_decoder_init(struct camu_decoder *dec, struct camu_renderer *rend
AVCodecParameters *codecpar = stream->av.stream->codecpar;
const AVCodec *codec = avcodec_find_decoder(codecpar->codec_id);
+#ifdef CAMU_FF_DECODER_HWACCEL
if (strcmp(codec->name, "libdav1d") == 0) {
codec = avcodec_find_decoder_by_name("av1");
}
+#endif
if (!codec) {
al_log_error("ff_decoder", "Failed to find decoder.");
diff --git a/src/codec/ffmpeg/demuxer.c b/src/codec/ffmpeg/demuxer.c
index 286be73..1e92186 100644
--- a/src/codec/ffmpeg/demuxer.c
+++ b/src/codec/ffmpeg/demuxer.c
@@ -3,7 +3,7 @@
#include "demuxer.h"
#include "avio.h"
-#define DEMUX_BUFFER_SIZE (4096 * 4)
+#define DEMUX_BUFFER_SIZE 4096
static void close_internal(struct camu_ff_demuxer *av)
{
@@ -164,9 +164,9 @@ static bool ff_demuxer_seek(struct camu_demuxer *demux, u64 pos)
al_assert(pos <= INT64_MAX);
if (pos > av->duration) pos = av->duration;
if (av->eof) {
- avformat_flush(av->format_context);
av->eof = false;
}
+ avformat_flush(av->format_context);
s64 ts = (s64)pos;
if (avformat_seek_file(av->format_context, -1, INT64_MIN, ts, ts, 0) < 0) {
return false;
diff --git a/src/fruits/cmsrv/cmsrv.c b/src/fruits/cmsrv/cmsrv.c
index 80d7c03..6f00a05 100644
--- a/src/fruits/cmsrv/cmsrv.c
+++ b/src/fruits/cmsrv/cmsrv.c
@@ -1,8 +1,9 @@
-#define CAMU_LOCAL_SOCKET
+#define CMSRV_LOCAL_SOCKET
+#define CMSRV_USE_UI
#include <al/log.h>
#include <nnwt/common.h>
-#ifdef CAMU_LOCAL_SOCKET
+#ifdef CMSRV_LOCAL_SOCKET
#include <nnwt/line_processor.h>
#endif
#include <nnwt/timer.h>
@@ -20,18 +21,20 @@ struct cmsrv {
struct nn_event_loop loop;
struct nn_signal quit_signal;
struct camu_server server;
-#ifdef CAMU_LOCAL_SOCKET
+#ifdef CMSRV_LOCAL_SOCKET
struct {
struct nn_socket sock;
struct nn_line_processor cli;
} local;
#endif
+#ifdef CMSRV_USE_UI
struct cmsrv_ui ui;
struct nn_poll input_poll;
struct nn_timer render_timer;
+#endif
};
-#ifdef CAMU_LOCAL_SOCKET
+#ifdef CMSRV_LOCAL_SOCKET
static u8 server_line_callback(void *userdata, str *line)
{
struct cmsrv *s = (struct cmsrv *)userdata;
@@ -73,6 +76,7 @@ static u8 server_line_callback(void *userdata, str *line)
}
#endif
+#ifdef CMSRV_USE_UI
static void render_timer_callback(void *userdata, struct nn_timer *timer)
{
struct cmsrv *s = (struct cmsrv *)userdata;
@@ -107,12 +111,6 @@ static void input_poll_callback(void *userdata, s32 revents)
}
}
-static void quit_signal_callback(void *userdata)
-{
- struct cmsrv *s = (struct cmsrv *)userdata;
- nn_event_loop_break_one(&s->loop);
-}
-
static s32 log_callback(void *userdata, u8 level, char *message)
{
struct cmsrv *s = (struct cmsrv *)userdata;
@@ -120,6 +118,23 @@ static s32 log_callback(void *userdata, u8 level, char *message)
cmsrv_ui_push_message(&s->ui, message);
return al_strlen(message);
}
+#endif
+
+static void server_meta_callback(void *userdata, u8 op, struct lia_list *list, struct lia_list_entry *entry)
+{
+ struct cmsrv *s = (struct cmsrv *)userdata;
+ (void)s;
+ (void)list;
+ if (op == LIANA_META_PLAYING) {
+ al_log_info("cmsrv", "now playing: %ls.", entry->name.data);
+ }
+}
+
+static void quit_signal_callback(void *userdata)
+{
+ struct cmsrv *s = (struct cmsrv *)userdata;
+ nn_event_loop_break_one(&s->loop);
+}
static struct cmsrv s = { 0 };
@@ -137,12 +152,14 @@ s32 wmain(s32 argc, wchar_t **argv)
{
(void)argc;
(void)argv;
- if (!nn_common_init() || !cmsrv_ui_init(&s.ui, &s.server)) return EXIT_FAILURE;
+ if (!nn_common_init()) return EXIT_FAILURE;
+#ifdef CMSRV_USE_UI
+ if (!cmsrv_ui_init(&s.ui, &s.server)) return EXIT_FAILURE;
+ al_set_print(log_callback, &s);
+#endif
signal(SIGINT, sigint_handler);
- al_set_print(log_callback, &s);
-
#ifdef CAMU_HAVE_FFMPEG
camu_ff_set_default_log_callback();
#endif
@@ -153,11 +170,13 @@ s32 wmain(s32 argc, wchar_t **argv)
nn_signal_start(&s.quit_signal);
camu_server_init(&s.server, &s.loop);
+ s.server.meta_callback = server_meta_callback;
+ s.server.userdata = &s;
if (!camu_server_listen(&s.server, CAMU_TEST_TYPE, CAMU_TEST_ADDR, CAMU_PORT)) {
return EXIT_FAILURE;
}
-#ifdef CAMU_LOCAL_SOCKET
+#ifdef CMSRV_LOCAL_SOCKET
s.local.sock.type = NNWT_SOCKET_UNIX;
nn_socket_init(&s.local.sock, NNWT_SOCKET_NONBLOCKING);
s.local.cli.callback = server_line_callback;
@@ -169,6 +188,7 @@ s32 wmain(s32 argc, wchar_t **argv)
}
#endif
+#ifdef CMSRV_USE_UI
nn_poll_init(&s.input_poll, input_poll_callback, &s);
nn_poll_set(&s.input_poll, cmsrv_ui_get_input_fd(&s.ui), NNWT_POLL_READ);
nn_poll_start(&s.input_poll, &s.loop);
@@ -176,10 +196,13 @@ s32 wmain(s32 argc, wchar_t **argv)
nn_timer_init(&s.render_timer, &s.loop, render_timer_callback, &s);
nn_timer_set_repeat(&s.render_timer, NNWT_TS_FROM_USEC(100000));
nn_timer_again(&s.render_timer);
+#endif
nn_event_loop_run(&s.loop);
+#ifdef CMSRV_USE_UI
cmsrv_ui_close(&s.ui);
+#endif
#ifdef CAMU_HAVE_FFMPEG
camu_ff_free_default_log_callback();
diff --git a/src/liana/client.c b/src/liana/client.c
index f1d2a61..daf32b0 100644
--- a/src/liana/client.c
+++ b/src/liana/client.c
@@ -158,6 +158,7 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
};
client->callback(client->userdata, LIANA_CLIENT_RESUME_AT, NULL, &time);
if (client->mask == 0) {
+ al_assert(client->connection_id == 0);
al_log_warn("liana", "Handling reconnect on unconfigured client.");
} else {
client->callback(client->userdata, LIANA_CLIENT_RECONNECTED, NULL, NULL);
@@ -167,7 +168,8 @@ static bool connection_callback(void *userdata, struct nn_packet_stream *stream)
stream->packet_sent_callback = packet_sent_callback;
struct nn_packet *packet = nn_packet_create();
nn_packet_write_u32(packet, client->id);
- nn_packet_write_u32(packet, client->connection_id);
+ //nn_packet_write_u32(packet, client->connection_id);
+ nn_packet_write_u32(packet, 0);
nn_packet_write_s32(packet, client->mask);
nn_packet_write_u64(packet, client->pos);
if (client->mask == 0) {
diff --git a/src/liana/list.c b/src/liana/list.c
index 218c51e..3c8a9c0 100644
--- a/src/liana/list.c
+++ b/src/liana/list.c
@@ -633,7 +633,7 @@ static void handle_clear(struct lia_list *list)
static void run_queue(struct lia_list *list)
{
if (!list->cmd) {
- if (list->queue.size == 0) return;
+ if (al_array_is_empty(list->queue)) return;
al_array_pop_at(list->queue, 0, list->cmd);
}
struct lia_list_cmd *cmd = list->cmd;
diff --git a/src/liana/server.c b/src/liana/server.c
index 3b7bc2d..87bcf57 100644
--- a/src/liana/server.c
+++ b/src/liana/server.c
@@ -4,7 +4,6 @@
#include "handler.h"
#include "handlers.h"
#include "list.h"
-#include "vcr.h"
static inline u32 get_incremental_id(struct lia_server *server)
{
@@ -167,9 +166,7 @@ static void handle_connection(struct lia_node_connection *conn, struct nn_packet
stream->packet_sent_callback = data_packet_sent_callback;
stream->packets_sent_callback = data_packets_sent_callback;
stream->connection_closed_callback = data_connection_closed_callback;
-#ifndef VCR_BUFFER_WHOLE_FILE
nn_thread_create(&conn->thread, handler_thread, conn);
-#endif
}
}
@@ -202,7 +199,7 @@ static void signal_callback(void *userdata)
struct nn_packet_stream *stream = conn->stream;
if (!conn->errored) {
conn->id = get_incremental_id(server);
- nn_packet_pool_init(&conn->pool, 512, server->loop, packet_pool_callback, conn);
+ nn_packet_pool_init(&conn->pool, 1024, 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 c12c422..06f0fe4 100644
--- a/src/liana/vcr.c
+++ b/src/liana/vcr.c
@@ -103,6 +103,8 @@ static nn_thread_result NNWT_THREADCALL vcr_track_thread(void *userdata)
bool success = track->client->handle_packet(track->client, packet);
if (!success) {
al_log_error("liana", "Error handling packet, exiting track thread.");
+ return_entire_cache(track);
+ track->cache.disabled = true;
nn_packet_cache_unlock(&track->cache);
nn_mutex_unlock(&track->mutex);
return 0;
@@ -186,7 +188,7 @@ void lia_vcr_add_track(struct lia_vcr *vcr, struct lia_vcr_track *track)
bool lia_vcr_is_empty(struct lia_vcr *vcr)
{
- return vcr->tracks.size == 0;
+ return al_array_is_empty(vcr->tracks);
}
static struct lia_vcr_track *get_track_from_index(struct lia_vcr *vcr, s32 index)
@@ -347,9 +349,7 @@ void lia_vcr_flush(struct lia_vcr *vcr)
al_array_foreach(vcr->tracks, i, track) {
if (VCR_TRACK_THREADED(track)) {
vcr_track_close_internal(track);
-#ifndef VCR_BUFFER_WHOLE_FILE
return_entire_cache(track);
-#endif
al_atomic_store(u8)(&track->buffered, 0, AL_ATOMIC_RELAXED);
track->client->flush(track->client);
nn_packet_cache_enable(&track->cache);
diff --git a/src/libsink/sink.c b/src/libsink/sink.c
index 1569286..f588cfa 100644
--- a/src/libsink/sink.c
+++ b/src/libsink/sink.c
@@ -38,51 +38,64 @@ enum {
// Command queue commands.
enum {
- START,
+ START = 0,
STOP,
+ ADD_BUFFER,
+ REMOVE_BUFFER,
+ CLEAR,
+ CLOSE,
TOGGLE_PAUSE,
SEEK,
SKIP,
SHUFFLE,
RESEEK,
UNSET,
- END,
- CLOSE
+ END
};
// Number of entries to keep buffered at one time.
#define ENTRY_MAX_AGE 7
-// If a buffer is still INIT or QUEUED after an entry is configured, it's "empty".
+// If a buffer is still INIT or QUEUED after the entry is configured, it's "empty".
#define BUFFER_EMPTY(buf) ((buf)->state == BUFFER_INIT || (buf)->state == BUFFER_QUEUED)
#define BUFFER_NOT_EMPTY(buf) (!BUFFER_EMPTY(buf))
-#ifdef CAMU_SINK_NO_VIDEO
-#define VIDEO_ADDED_OR_EMPTY(entry) true
-#else
-#define VIDEO_ADDED_OR_EMPTY(entry) ((entry)->video.state == BUFFER_ADDED || BUFFER_EMPTY(&(entry)->video))
+#define AUDIO_EMPTY(entry) BUFFER_EMPTY(&(entry)->audio)
+#ifndef CAMU_SINK_NO_VIDEO
+#define VIDEO_EMPTY(entry) BUFFER_EMPTY(&(entry)->video)
#endif
-#define AUDIO_ADDED_OR_EMPTY(entry) ((entry)->audio.state == BUFFER_ADDED || BUFFER_EMPTY(&(entry)->audio))
-#ifdef CAMU_SINK_NO_VIDEO
-#define VIDEO_ENDED_OR_EMPTY(entry) true
+#define AUDIO_NOT_EMPTY(entry) BUFFER_NOT_EMPTY(&(entry)->audio)
+#ifndef CAMU_SINK_NO_VIDEO
+#define VIDEO_NOT_EMPTY(entry) BUFFER_NOT_EMPTY(&(entry)->video)
+#endif
+
+#define AUDIO_ADDED_OR_EMPTY(entry) ((entry)->audio.state == BUFFER_ADDED || BUFFER_EMPTY(&(entry)->audio))
+#ifndef CAMU_SINK_NO_VIDEO
+#define VIDEO_ADDED_OR_EMPTY(entry) ((entry)->video.state == BUFFER_ADDED || BUFFER_EMPTY(&(entry)->video))
#else
-#define VIDEO_ENDED_OR_EMPTY(entry) ((entry)->video.state == BUFFER_ENDED || BUFFER_EMPTY(&(entry)->video))
+#define VIDEO_ADDED_OR_EMPTY(entry) true
#endif
+
#define AUDIO_ENDED_OR_EMPTY(entry) ((entry)->audio.state == BUFFER_ENDED || BUFFER_EMPTY(&(entry)->audio))
+#ifndef CAMU_SINK_NO_VIDEO
+#define VIDEO_ENDED_OR_EMPTY(entry) ((entry)->video.state == BUFFER_ENDED || BUFFER_EMPTY(&(entry)->video))
+#else
+#define VIDEO_ENDED_OR_EMPTY(entry) true
+#endif
-#ifdef CAMU_SINK_NO_VIDEO
-#define VIDEO_REMOVED_OR_EMPTY(entry) true
+#define AUDIO_NOT_ADDED(entry) ((entry)->audio.state != BUFFER_ADDED)
+#ifndef CAMU_SINK_NO_VIDEO
+#define VIDEO_NOT_ADDED(entry) ((entry)->video.state != BUFFER_ADDED)
#else
-#define VIDEO_REMOVED_OR_EMPTY(entry) ((entry)->video.state != BUFFER_ADDED)
+#define VIDEO_NOT_ADDED(entry) true
#endif
-#define AUDIO_REMOVED_OR_EMPTY(entry) ((entry)->audio.state != BUFFER_ADDED)
#ifndef CAMU_SINK_NO_VIDEO
-#define ENTRY_IS_SINGLE_FRAME(entry) camu_video_buffer_is_single_frame(&(entry)->video.buf)
+#define VIDEO_IS_SINGLE_FRAME(entry) camu_video_buffer_is_single_frame(&(entry)->video.buf)
#endif
-#if defined CAMU_SCREEN_THREADED && defined CAMU_MIXER_THREADED
+#if defined CAMU_SCREEN_THREADED && defined CAMU_MIXER_THREADED_START_STOP
#define BLOCKING_SLEEP(delay) nn_thread_sleep(delay)
#else
#define BLOCKING_SLEEP(delay) nn_event_loop_sleep(sink->loop, delay)
@@ -110,12 +123,48 @@ static inline bool entry_video_buffer_held(struct camu_sink_entry *entry)
}
#endif
+static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
+{
+ camu_queue_push(sink->queue, cmd);
+ nn_signal_send(&sink->queue_signal);
+}
+
+static inline void add_entry_audio_buffer(struct camu_sink_entry *entry)
+{
+ struct camu_sink *sink = entry->sink;
+#ifdef CAMU_MIXER_THREADED_START_STOP
+ sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+#else
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = ADD_BUFFER,
+ .value.i = CAMU_SINK_AUDIO,
+ .opaque = entry
+ });
+#endif
+}
+
+#ifndef CAMU_SINK_NO_VIDEO
+static inline void add_entry_video_buffer(struct camu_sink_entry *entry)
+{
+ struct camu_sink *sink = entry->sink;
+ sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+}
+#endif
+
static void remove_entry_audio_buffer(struct camu_sink *sink, struct camu_sink_entry *entry)
{
al_assert(!entry->ended && entry->audio.state != BUFFER_ENDED);
al_assert(entry->audio.state != BUFFER_INIT);
if (entry->audio.state == BUFFER_ADDED) {
+#ifdef CAMU_MIXER_THREADED_START_STOP
sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+#else
+ queue_cmd(sink, (struct camu_sink_cmd){
+ .op = REMOVE_BUFFER,
+ .value.i = CAMU_SINK_AUDIO,
+ .opaque = entry
+ });
+#endif
entry->audio.state = BUFFER_SET_OR_BUFFERED;
} else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) {
entry->audio.state = BUFFER_CONFIGURED;
@@ -187,12 +236,12 @@ 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);
- if (!BUFFER_EMPTY(&entry->audio) && sink->audio.state == SINK_PAUSED) {
+ if (AUDIO_NOT_EMPTY(entry) && sink->audio.state == SINK_PAUSED) {
sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_AUDIO, NULL);
sink->audio.state = SINK_PLAYING;
}
#ifndef CAMU_SINK_NO_VIDEO
- if (!ENTRY_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PAUSED) {
+ if (!VIDEO_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PAUSED) {
sink->callback(sink->userdata, CAMU_SINK_START, CAMU_SINK_VIDEO, NULL);
sink->video.state = SINK_PLAYING;
}
@@ -200,7 +249,7 @@ static void sink_local_pause(struct camu_sink *sink, struct camu_sink_entry *ent
} else {
camu_clock_pause(&entry->clock, 0);
#ifndef CAMU_SINK_NO_VIDEO
- if (!ENTRY_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PLAYING) {
+ if (!VIDEO_IS_SINGLE_FRAME(entry) && sink->video.state == SINK_PLAYING) {
sink->callback(sink->userdata, CAMU_SINK_STOP, CAMU_SINK_VIDEO, NULL);
sink->video.state = SINK_PAUSED;
}
@@ -224,7 +273,6 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
{
switch (cmd->op) {
case START: {
- nn_mutex_lock(&sink->mutex);
switch (cmd->value.i) {
case CAMU_SINK_AUDIO:
if (sink->audio.state == SINK_PAUSED) {
@@ -241,11 +289,9 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
#endif
}
- nn_mutex_unlock(&sink->mutex);
break;
}
case STOP: {
- nn_mutex_lock(&sink->mutex);
switch (cmd->value.i) {
case CAMU_SINK_AUDIO:
if (sink->audio.state == SINK_PLAYING) {
@@ -262,9 +308,54 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
break;
#endif
}
- nn_mutex_unlock(&sink->mutex);
break;
}
+ case ADD_BUFFER: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ sink->callback(sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ break;
+#endif
+ }
+ break;
+ }
+ case REMOVE_BUFFER: {
+ struct camu_sink_entry *entry = (struct camu_sink_entry *)cmd->opaque;
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ sink->callback(sink->userdata, CAMU_SINK_REMOVE_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ break;
+#endif
+ }
+ break;
+ }
+ case CLEAR: {
+ switch (cmd->value.i) {
+ case CAMU_SINK_AUDIO:
+ sink->callback(sink->userdata, CAMU_SINK_CLEAR, CAMU_SINK_AUDIO, NULL);
+ break;
+#ifndef CAMU_SINK_NO_VIDEO
+ case CAMU_SINK_VIDEO:
+ sink->callback(sink->userdata, CAMU_SINK_CLEAR, CAMU_SINK_VIDEO, NULL);
+ break;
+#endif
+ }
+ break;
+ }
+ case CLOSE: {
+ nn_signal_stop(&sink->queue_signal);
+ sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
+ return;
+ }
case SKIP: {
if (!sink->conn) return;
struct nn_packet *packet = nn_rpc_get_packet(&sink->client, CAMU_SERVER_LIST_ACTION);
@@ -338,11 +429,6 @@ static void handle_sink_cmd(struct camu_sink *sink, struct camu_sink_cmd *cmd)
nn_rpc_connection_command(sink->conn, packet, NULL, NULL);
break;
}
- case CLOSE: {
- nn_signal_stop(&sink->queue_signal);
- sink->callback(sink->userdata, CAMU_SINK_EXIT, 0, NULL);
- return;
- }
}
}
@@ -358,17 +444,11 @@ static void queue_signal_callback(void *userdata)
}
}
-static void queue_cmd(struct camu_sink *sink, struct camu_sink_cmd cmd)
-{
- camu_queue_push(sink->queue, cmd);
- nn_signal_send(&sink->queue_signal);
-}
-
static void mixer_callback(void *userdata, u8 op)
{
struct camu_sink *sink = (struct camu_sink *)userdata;
if (op == CAMU_MIXER_EMPTY) {
- al_log_debug("sink", "Mixer empty.");
+ al_log_info("sink", "Mixer empty.");
queue_cmd(sink, (struct camu_sink_cmd){
.op = STOP,
.value.i = CAMU_SINK_AUDIO
@@ -510,31 +590,34 @@ void add_audio_if_set_and_buffered(struct camu_sink_entry *entry)
entry->audio.state = BUFFER_SET_OR_BUFFERED;
} else if (entry->audio.state == BUFFER_SET_OR_BUFFERED) {
entry->audio.state = BUFFER_ADDED;
-#ifndef CAMU_SINK_LOCAL
+ bool ignore_video = true;
+#ifndef CAMU_SINK_NO_VIDEO
+ ignore_video = VIDEO_EMPTY(entry) || VIDEO_IS_SINGLE_FRAME(entry);
+#endif
+ camu_audio_buffer_set_no_video(&entry->audio.buf, ignore_video);
+#ifdef CAMU_SINK_LOCAL
+ // If we're local we don't have to worry about syncing audio-only entries.
+ camu_audio_buffer_set_ignore_desync(&entry->audio.buf, ignore_video);
+#else
+ // Force resync.
camu_audio_buffer_unpause(&entry->audio.buf);
#endif
if (VIDEO_ADDED_OR_EMPTY(entry)) {
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ add_entry_audio_buffer(entry);
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = entry->audio.buffer_paused ? STOP : START,
.value.i = CAMU_SINK_AUDIO
});
#ifndef CAMU_SINK_NO_VIDEO
- bool stop_video = false;
- // Single frames are unconditionally added in add_video_if_set_and_buffered().
- if (!ENTRY_IS_SINGLE_FRAME(entry)) {
- if (entry->video.state == BUFFER_ADDED) {
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
- } else {
- // Entry has no video.
- stop_video = true;
- }
+ if (!ignore_video) {
+ // Single frames are unconditionally added in add_video_if_set_and_buffered().
+ add_entry_video_buffer(entry);
}
#endif
maybe_remove_previous(entry->sink);
#ifndef CAMU_SINK_NO_VIDEO
- // This must come after maybe_remove_previous().
- if (stop_video) {
+ if (VIDEO_EMPTY(entry)) {
+ // This must come after maybe_remove_previous().
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = STOP,
.value.i = CAMU_SINK_VIDEO
@@ -558,19 +641,19 @@ void add_video_if_set_and_buffered(struct camu_sink_entry *entry)
entry->video.state = BUFFER_SET_OR_BUFFERED;
} else if (entry->video.state == BUFFER_SET_OR_BUFFERED) {
entry->video.state = BUFFER_ADDED;
- bool single_frame = ENTRY_IS_SINGLE_FRAME(entry);
+ bool single_frame = VIDEO_IS_SINGLE_FRAME(entry);
if (AUDIO_ADDED_OR_EMPTY(entry)) {
if (entry->audio.state == BUFFER_ADDED) {
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_AUDIO, &entry->audio.buf);
+ add_entry_audio_buffer(entry);
queue_cmd(entry->sink, (struct camu_sink_cmd){
.op = entry->audio.buffer_paused ? STOP : START,
.value.i = CAMU_SINK_AUDIO
});
}
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ add_entry_video_buffer(entry);
maybe_remove_previous(entry->sink);
} else if (single_frame) {
- entry->sink->callback(entry->sink->userdata, CAMU_SINK_ADD_BUFFER, CAMU_SINK_VIDEO, &entry->video.buf);
+ add_entry_video_buffer(entry);
}
// Stopping video here is needed if skipping from a video to an image.
queue_cmd(entry->sink, (struct camu_sink_cmd){
@@ -591,11 +674,11 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
maybe_add_to_previous(sink, current, target);
} else {
#ifndef CAMU_SINK_NO_VIDEO
- if (ENTRY_IS_SINGLE_FRAME(current)) {
+ if (VIDEO_IS_SINGLE_FRAME(current)) {
remove_entry_video_buffer(sink, current);
}
#endif
- al_assert(AUDIO_REMOVED_OR_EMPTY(current) && VIDEO_REMOVED_OR_EMPTY(current));
+ al_assert(AUDIO_NOT_ADDED(current) && VIDEO_NOT_ADDED(current));
}
if (sink->reconnecting) {
@@ -607,7 +690,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
if (!target->ended) {
add_or_queue_entry(target);
#ifndef CAMU_SINK_NO_VIDEO
- } else if (ENTRY_IS_SINGLE_FRAME(target)) {
+ } else if (VIDEO_IS_SINGLE_FRAME(target)) {
add_video_if_set_and_buffered(target);
#endif
}
@@ -620,6 +703,7 @@ static void switch_to(struct camu_sink *sink, struct camu_sink_entry *target)
static void pause_and_swap_to(struct camu_sink *sink, struct camu_sink_entry *target, u64 at)
{
struct camu_sink_entry *current = sink->current;
+ al_assert(target != current);
if (!current || current->ended) {
switch_to(sink, target);
} else {
@@ -693,11 +777,11 @@ static void audio_buffer_callback(void *userdata, u8 op)
// This assert likely doesn't matter due to the handling of the ENDED state.
al_assert(entry->audio.state == BUFFER_SET_OR_BUFFERED);
entry->audio.state = BUFFER_ENDED;
- bool run_queue = VIDEO_ENDED_OR_EMPTY(entry);
+ bool end_entry = VIDEO_ENDED_OR_EMPTY(entry);
#ifndef CAMU_SINK_NO_VIDEO
- run_queue = run_queue || ENTRY_IS_SINGLE_FRAME(entry);
+ end_entry = end_entry || VIDEO_IS_SINGLE_FRAME(entry);
#endif
- if (run_queue) {
+ if (end_entry) {
end_entry_and_advance_queue(sink, entry);
}
nn_mutex_unlock(&sink->mutex);
@@ -731,7 +815,10 @@ static void video_buffer_callback(void *userdata, u8 op)
bool swapped = false;
nn_mutex_lock(&sink->mutex);
al_log_info("sink", "Video EOF.");
- if (!ENTRY_IS_SINGLE_FRAME(entry)) {
+ if (AUDIO_NOT_EMPTY(entry)) {
+ camu_audio_buffer_set_no_video(&entry->audio.buf, true);
+ }
+ if (!VIDEO_IS_SINGLE_FRAME(entry)) {
if (entry->video.state == BUFFER_ADDED) {
remove_entry_video_buffer(sink, entry);
}
@@ -785,7 +872,7 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent
{
#ifdef CAMU_SINK_LOCAL
#ifndef CAMU_SINK_NO_VIDEO
- if (BUFFER_NOT_EMPTY(&entry->audio) && BUFFER_NOT_EMPTY(&entry->video)) {
+ if (AUDIO_NOT_EMPTY(entry) && VIDEO_NOT_EMPTY(entry)) {
f64 audio = camu_mixer_get_latency(sink->audio.mixer);
s32 frames = audio / entry->video.buf.avg_frame_duration;
frames -= sink->video.renderer->get_latency(sink->video.renderer);
@@ -800,7 +887,7 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent
// latency directly into the audio buffer.
f64 audio = camu_mixer_get_latency(sink->audio.mixer);
#ifndef CAMU_SINK_NO_VIDEO
- if (BUFFER_NOT_EMPTY(&entry->video)) {
+ if (VIDEO_NOT_EMPTY(entry)) {
s32 frames = audio / entry->video.buf.avg_frame_duration;
frames += sink->video.renderer->get_latency(sink->video.renderer);
camu_video_buffer_set_latency(&entry->video.buf, frames);
@@ -810,6 +897,19 @@ static void evaluate_latency(struct camu_sink *sink, struct camu_sink_entry *ent
#endif
}
+static void run_queue_by_opaque(struct camu_sink *sink, void *opaque)
+{
+ camu_queue_lock(sink->queue);
+ struct camu_sink_cmd *cmd;
+ al_array_foreach_ptr(sink->queue.a, i, cmd) {
+ if (cmd->opaque == opaque) {
+ handle_sink_cmd(sink, cmd);
+ al_array_remove_at_iter(sink->queue.a, i);
+ }
+ }
+ camu_queue_unlock(sink->queue);
+}
+
static void client_callback(void *userdata, u8 op, struct camu_codec_stream *stream, void *opaque)
{
struct camu_sink_entry *entry = (struct camu_sink_entry *)userdata;
@@ -833,7 +933,9 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
} else {
assert(false);
}
- if (VIDEO_ADDED_OR_EMPTY(entry)) evaluate_latency(sink, entry);
+ if (VIDEO_ADDED_OR_EMPTY(entry)) {
+ evaluate_latency(sink, entry);
+ }
nn_mutex_unlock(&sink->mutex);
break;
#ifndef CAMU_SINK_NO_VIDEO
@@ -852,7 +954,9 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
} else {
assert(false);
}
- if (AUDIO_ADDED_OR_EMPTY(entry)) evaluate_latency(sink, entry);
+ if (AUDIO_ADDED_OR_EMPTY(entry)) {
+ evaluate_latency(sink, entry);
+ }
nn_mutex_unlock(&sink->mutex);
break;
case CAMU_STREAM_SUBTITLE:
@@ -879,14 +983,14 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
struct camu_codec_frame *frame = (struct camu_codec_frame *)opaque;
switch (stream->type) {
case CAMU_STREAM_AUDIO:
- if (BUFFER_NOT_EMPTY(&entry->audio)) {
+ if (AUDIO_NOT_EMPTY(entry)) {
camu_audio_buffer_push(&entry->audio.buf, frame);
return;
}
break;
#ifndef CAMU_SINK_NO_VIDEO
case CAMU_STREAM_VIDEO:
- if (BUFFER_NOT_EMPTY(&entry->video)) {
+ if (VIDEO_NOT_EMPTY(entry)) {
camu_video_buffer_push(&entry->video.buf, frame);
return;
}
@@ -917,11 +1021,6 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
nn_mutex_lock(&sink->mutex);
- // We need to call this if entry was added to previous then,
- // - it's being cleaned up after (ENTRY_MAX_AGE - 1) entries were added but none buffered.
- // - it was seeked.
- maybe_remove_from_previous(sink, entry);
-
if (reconnect) {
entry->ended = false;
if (entry == sink->current) {
@@ -933,6 +1032,11 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
}
}
+ // We need to call this if entry was added to previous then,
+ // - it's being cleaned up after ENTRY_MAX_AGE - 1 entries were added but none buffered.
+ // - it was seeked.
+ maybe_remove_from_previous(sink, entry);
+
bool skip_audio = sink->audio.state == SINK_PAUSED;
if (entry->audio.state == BUFFER_ADDED) {
remove_entry_audio_buffer(sink, entry);
@@ -941,11 +1045,9 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
if (entry->audio.state > BUFFER_CONFIGURED) {
entry->audio.state = BUFFER_CONFIGURED;
}
+
#ifndef CAMU_SINK_NO_VIDEO
- bool ignore_video = BUFFER_EMPTY(&entry->video);
- if (!ignore_video) {
- ignore_video = reconnect && ENTRY_IS_SINGLE_FRAME(entry);
- }
+ bool ignore_video = VIDEO_EMPTY(entry) || (reconnect && VIDEO_IS_SINGLE_FRAME(entry));
if (!ignore_video) {
if (entry->video.state == BUFFER_ADDED) {
remove_entry_video_buffer(sink, entry);
@@ -956,6 +1058,9 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
}
#endif
+ // Resolve any queued REMOVE_BUFFER requests before waiting.
+ run_queue_by_opaque(sink, entry);
+
nn_mutex_unlock(&sink->mutex);
while ( // Block until buffers are removed.
@@ -967,11 +1072,11 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
) { BLOCKING_SLEEP(NNWT_TS_FROM_USEC(2000)); }
if (reconnect) {
- if (BUFFER_NOT_EMPTY(&entry->audio)) {
+ if (AUDIO_NOT_EMPTY(entry)) {
camu_audio_buffer_reset(&entry->audio.buf);
}
#ifndef CAMU_SINK_NO_VIDEO
- if (BUFFER_NOT_EMPTY(&entry->video) && !ignore_video) {
+ if (VIDEO_NOT_EMPTY(entry) && !ignore_video) {
camu_video_buffer_reset(&entry->video.buf);
}
#endif
@@ -998,14 +1103,16 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
case LIANA_CLIENT_RECONNECTED: {
nn_mutex_lock(&sink->mutex);
if (entry == sink->reconnecting) {
+ al_assert(entry == sink->current);
if (entry->audio.state > BUFFER_QUEUED) {
add_audio_if_set_and_buffered(entry);
}
#ifndef CAMU_SINK_NO_VIDEO
- if (entry->video.state > BUFFER_QUEUED && !ENTRY_IS_SINGLE_FRAME(entry)) {
+ if (entry->video.state > BUFFER_QUEUED && !VIDEO_IS_SINGLE_FRAME(entry)) {
add_video_if_set_and_buffered(entry);
}
#endif
+ sink->reconnecting = NULL;
}
nn_mutex_unlock(&sink->mutex);
break;
@@ -1013,7 +1120,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
case LIANA_CLIENT_EOF: {
switch (stream->type) {
case CAMU_STREAM_AUDIO: {
- if (BUFFER_NOT_EMPTY(&entry->audio)) {
+ if (AUDIO_NOT_EMPTY(entry)) {
camu_audio_buffer_flush(&entry->audio.buf);
}
break;
@@ -1021,7 +1128,7 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
#ifndef CAMU_SINK_NO_VIDEO
case CAMU_STREAM_VIDEO: {
// Single frames are immediately flushed inside the buffer.
- if (BUFFER_NOT_EMPTY(&entry->video) && !ENTRY_IS_SINGLE_FRAME(entry)) {
+ if (VIDEO_NOT_EMPTY(entry) && !VIDEO_IS_SINGLE_FRAME(entry)) {
camu_video_buffer_flush(&entry->video.buf);
}
break;
@@ -1053,16 +1160,20 @@ static void client_callback(void *userdata, u8 op, struct camu_codec_stream *str
sink->queued = NULL;
}
+ bool removed;
+ al_array_check_remove(sink->entries, entry, removed);
+
+ nn_mutex_unlock(&sink->mutex);
+
lia_client_free(&entry->client);
camu_audio_buffer_free(&entry->audio.buf);
#ifndef CAMU_SINK_NO_VIDEO
camu_video_buffer_free(&entry->video.buf);
#endif
- bool removed;
- al_array_check_remove(sink->entries, entry, removed);
- nn_mutex_unlock(&sink->mutex);
al_free(entry);
+
al_log_info("sink", "Entry closed by %s.", removed ? "disconnect" : "cleanup");
+
break;
}
}
@@ -1440,7 +1551,7 @@ void camu_sink_toggle_pause(struct camu_sink *sink)
if (current) {
queue_cmd(sink, (struct camu_sink_cmd){
.op = TOGGLE_PAUSE,
- .value.f = camu_clock_get_pts(&current->clock, 0.0, true),
+ .value.f = camu_clock_get_pts(&current->clock, 0.0, false),
.opaque = current
});
}
@@ -1488,6 +1599,10 @@ void camu_sink_stop(struct camu_sink *sink)
.value.i = CAMU_SINK_AUDIO
});
queue_cmd(sink, (struct camu_sink_cmd){
+ .op = CLEAR,
+ .value.i = CAMU_SINK_AUDIO
+ });
+ queue_cmd(sink, (struct camu_sink_cmd){
.op = CLOSE
});
}
diff --git a/src/libsink/sink.h b/src/libsink/sink.h
index c08f763..10fc182 100644
--- a/src/libsink/sink.h
+++ b/src/libsink/sink.h
@@ -26,9 +26,9 @@ enum {
enum {
CAMU_SINK_ADD_BUFFER = 0,
CAMU_SINK_REMOVE_BUFFER,
- CAMU_SINK_SWAP_BUFFER,
CAMU_SINK_START,
CAMU_SINK_STOP,
+ CAMU_SINK_CLEAR,
CAMU_SINK_EXIT
};
diff --git a/src/mixer/mixer.c b/src/mixer/mixer.c
index 5a21b19..1360934 100644
--- a/src/mixer/mixer.c
+++ b/src/mixer/mixer.c
@@ -13,10 +13,12 @@ static s32 data_callback(void *userdata, u8 *data, s32 frame_count, bool *silenc
{
struct camu_mixer *mixer = (struct camu_mixer *)userdata;
size_t req = camu_audio_format_samples_to_bytes(&mixer->fmt.req, (size_t)frame_count);
+#ifdef CAMU_MIXER_THREADED_START_STOP
if (UNLIKELY(mixer->paused)) {
al_memset(data, 0, req);
*silence = MIXER_WANT_INITIAL_SILENCE;
} else {
+#endif
#ifdef CAMU_MIXER_THREADED
if (al_atomic_load(u8)(&mixer->queued, AL_ATOMIC_RELAXED)) {
camu_mixer_run_queue(mixer);
@@ -58,7 +60,9 @@ static s32 data_callback(void *userdata, u8 *data, s32 frame_count, bool *silenc
}
break;
}
+#ifdef CAMU_MIXER_THREADED_START_STOP
}
+#endif
return frame_count;
}
@@ -260,7 +264,7 @@ void camu_mixer_clear(struct camu_mixer *mixer)
void camu_mixer_pause(struct camu_mixer *mixer)
{
-#ifdef CAMU_MIXER_THREADED
+#ifdef CAMU_MIXER_THREADED_START_STOP
nn_mutex_lock(&mixer->mutex);
#endif
if (!mixer->paused) {
@@ -268,25 +272,29 @@ void camu_mixer_pause(struct camu_mixer *mixer)
mixer->paused = true;
}
#ifdef CAMU_MIXER_THREADED
+ // This assumes stop() blocks until the output actually stops.
+ // That may not be the behavior we want to require.
run_queue_internal(mixer);
+#endif
+#ifdef CAMU_MIXER_THREADED_START_STOP
nn_mutex_unlock(&mixer->mutex);
#endif
}
void camu_mixer_resume(struct camu_mixer *mixer)
{
-#ifdef CAMU_MIXER_THREADED
+#ifdef CAMU_MIXER_THREADED_START_STOP
nn_mutex_lock(&mixer->mutex);
#endif
if (mixer->paused) {
// start() can internally call data_callback once before returning.
- // In that call mixer->paused will still be true. So, we have a special case to
- // immediately return silence and avoid any possible locking.
+ // In that call mixer->paused will still be true. So, if MIXER_THREADED_START_STOP
+ // is defined, we have a special case to immediately return silence to avoid a deadlock.
// Outputs can treat that silence as part of the stream with MIXER_WANT_INITIAL_SILENCE.
mixer->audio->start(mixer->audio);
mixer->paused = false;
}
-#ifdef CAMU_MIXER_THREADED
+#ifdef CAMU_MIXER_THREADED_START_STOP
nn_mutex_unlock(&mixer->mutex);
#endif
}
diff --git a/src/mixer/mixer.h b/src/mixer/mixer.h
index 16ec924..802d1b5 100644
--- a/src/mixer/mixer.h
+++ b/src/mixer/mixer.h
@@ -9,6 +9,14 @@
#include "../codec/codec.h"
+#ifdef CAMU_MIXER_THREADED
+// Do we need to lock in order to synchronize the mixers paused state.
+// Disabling this is a very specific optimization to allow the audio device to
+// buffer data during start(). It requires pause(), resume(), and remove_buffer()
+// to all come from the same thread.
+//#define CAMU_MIXER_THREADED_START_STOP
+#endif
+
enum {
CAMU_MIXER_EMPTY = 0
};
diff --git a/src/render/renderer_libplacebo.c b/src/render/renderer_libplacebo.c
index ad7e7fc..3de0fdc 100644
--- a/src/render/renderer_libplacebo.c
+++ b/src/render/renderer_libplacebo.c
@@ -269,7 +269,7 @@ static void renderer_lp_render(struct camu_renderer *renderer, struct camu_scree
}
}
- if (scr->videos.size == 0 && !force) {
+ if (al_array_is_empty(scr->videos) && !force) {
return;
}
diff --git a/src/server/server.c b/src/server/server.c
index c2b30a0..e949ee1 100644
--- a/src/server/server.c
+++ b/src/server/server.c
@@ -371,7 +371,7 @@ static void simple_search_portal_callback(void *userdata0, void *userdata1, stru
return;
}
case CAMU_CLIENT_GET_PAGE: {
- if (!result->page || result->page->list.size == 0) break;
+ if (!result->page || al_array_is_empty(result->page->list)) break;
struct camu_post *post = camu_post_cache_get(&server->cache, &al_array_at(result->page->list, 0));
if (!post) break;
struct cch_entry *entry = entry_from_post(server, post, 0);
diff --git a/src/sink/common.c b/src/sink/common.c
new file mode 100644
index 0000000..f720a3c
--- /dev/null
+++ b/src/sink/common.c
@@ -0,0 +1,81 @@
+#include <al/log.h>
+
+#include "../libsink/sink.h"
+
+#include "common.h"
+
+bool camu_default_sink_callback(struct camu_screen *scr, struct camu_mixer *mixer, u8 op, u8 type, void *opaque)
+{
+ switch (op) {
+ case CAMU_SINK_ADD_BUFFER:
+ switch (type) {
+ case CAMU_SINK_AUDIO: {
+ struct camu_audio_buffer *buf = (struct camu_audio_buffer *)opaque;
+ camu_mixer_add_buffer(mixer, buf);
+ al_log_info("camu_desktop", "Audio buffer added.");
+ break;
+ }
+ case CAMU_SINK_VIDEO: {
+ struct camu_video_buffer *buf = (struct camu_video_buffer *)opaque;
+ camu_screen_add_buffer(scr, buf);
+ al_log_info("camu_desktop", "Video buffer added.");
+ break;
+ }
+ }
+ break;
+ case CAMU_SINK_REMOVE_BUFFER:
+ switch (type) {
+ case CAMU_SINK_AUDIO: {
+ struct camu_audio_buffer *buf = (struct camu_audio_buffer *)opaque;
+ camu_mixer_remove_buffer(mixer, buf);
+ al_log_info("camu_desktop", "Audio buffer removed.");
+ break;
+ }
+ case CAMU_SINK_VIDEO: {
+ struct camu_video_buffer *buf = (struct camu_video_buffer *)opaque;
+ camu_screen_remove_buffer(scr, buf);
+ al_log_info("camu_desktop", "Video buffer removed.");
+ break;
+ }
+ }
+ break;
+ case CAMU_SINK_START:
+ switch (type) {
+ case CAMU_SINK_AUDIO:
+ camu_mixer_resume(mixer);
+ al_log_info("camu_desktop", "Audio started.");
+ break;
+ case CAMU_SINK_VIDEO:
+ camu_screen_set_state(scr, CAMU_SCREEN_PLAYING);
+ camu_screen_wake(scr);
+ al_log_info("camu_desktop", "Video started.");
+ break;
+ }
+ break;
+ case CAMU_SINK_STOP:
+ switch (type) {
+ case CAMU_SINK_AUDIO:
+ camu_mixer_pause(mixer);
+ al_log_info("camu_desktop", "Audio stopped.");
+ break;
+ case CAMU_SINK_VIDEO:
+ camu_screen_set_state(scr, CAMU_SCREEN_PAUSED);
+ al_log_info("camu_desktop", "Video stopped.");
+ break;
+ }
+ break;
+ case CAMU_SINK_CLEAR:
+ switch (type) {
+ case CAMU_SINK_AUDIO:
+ camu_mixer_clear(mixer);
+ break;
+ case CAMU_SINK_VIDEO:
+ camu_screen_clear(scr);
+ break;
+ }
+ break;
+ case CAMU_SINK_EXIT:
+ return false;
+ }
+ return true;
+}
diff --git a/src/sink/common.h b/src/sink/common.h
index dffaea1..5d02ed1 100644
--- a/src/sink/common.h
+++ b/src/sink/common.h
@@ -1,86 +1,4 @@
-#include <al/log.h>
-
#include "../screen/screen.h"
#include "../mixer/mixer.h"
-#include "../libsink/sink.h"
-
-AL_UNUSED_FUNCTION_PUSH
-
-static bool camu_default_sink_callback(struct camu_screen *scr, struct camu_mixer *mixer, u8 op, u8 type, void *opaque)
-{
- switch (op) {
- case CAMU_SINK_ADD_BUFFER:
- switch (type) {
- case CAMU_SINK_AUDIO: {
- struct camu_audio_buffer *buf = (struct camu_audio_buffer *)opaque;
- camu_mixer_add_buffer(mixer, buf);
- al_log_info("camu_desktop", "Audio buffer added.");
- break;
- }
- case CAMU_SINK_VIDEO: {
- struct camu_video_buffer *buf = (struct camu_video_buffer *)opaque;
- camu_screen_add_buffer(scr, buf);
- al_log_info("camu_desktop", "Video buffer added.");
- break;
- }
- }
- break;
- case CAMU_SINK_REMOVE_BUFFER:
- switch (type) {
- case CAMU_SINK_AUDIO: {
- struct camu_audio_buffer *buf = (struct camu_audio_buffer *)opaque;
- camu_mixer_remove_buffer(mixer, buf);
- al_log_info("camu_desktop", "Audio buffer removed.");
- break;
- }
- case CAMU_SINK_VIDEO: {
- struct camu_video_buffer *buf = (struct camu_video_buffer *)opaque;
- camu_screen_remove_buffer(scr, buf);
- al_log_info("camu_desktop", "Video buffer removed.");
- break;
- }
- }
- break;
- case CAMU_SINK_SWAP_BUFFER:
- switch (type) {
- case CAMU_SINK_AUDIO: {
- break;
- }
- }
- break;
- case CAMU_SINK_START:
- switch (type) {
- case CAMU_SINK_AUDIO: {
- camu_mixer_resume(mixer);
- al_log_info("camu_desktop", "Audio started.");
- break;
- }
- case CAMU_SINK_VIDEO: {
- camu_screen_set_state(scr, CAMU_SCREEN_PLAYING);
- camu_screen_wake(scr);
- al_log_info("camu_desktop", "Video started.");
- break;
- }
- }
- break;
- case CAMU_SINK_STOP:
- switch (type) {
- case CAMU_SINK_AUDIO: {
- camu_mixer_pause(mixer);
- al_log_info("camu_desktop", "Audio stopped.");
- break;
- }
- case CAMU_SINK_VIDEO: {
- camu_screen_set_state(scr, CAMU_SCREEN_PAUSED);
- al_log_info("camu_desktop", "Video stopped.");
- break;
- }
- }
- break;
- case CAMU_SINK_EXIT:
- return false;
- }
- return true;
-}
-AL_UNUSED_FUNCTION_POP
+bool camu_default_sink_callback(struct camu_screen *scr, struct camu_mixer *mixer, u8 op, u8 type, void *opaque);
diff --git a/src/sink/desktop.c b/src/sink/desktop.c
index 5529a3b..cb5e7c6 100644
--- a/src/sink/desktop.c
+++ b/src/sink/desktop.c
@@ -127,8 +127,6 @@ void camu_desktop_stop(struct camu_desktop *c)
camu_input_simulator_stop();
}
#endif
- camu_mixer_pause(&c->mixer);
- camu_mixer_clear(&c->mixer);
// We have already stopped ticking the screen at this point.
camu_screen_clear(&c->scr);
camu_sink_stop(&c->sink);
diff --git a/src/sink/meson.build b/src/sink/meson.build
index 30a5f39..124c7a5 100644
--- a/src/sink/meson.build
+++ b/src/sink/meson.build
@@ -1,4 +1,4 @@
-desktop_src = ['desktop.c', 'input_simulator.c']
+desktop_src = ['desktop.c', 'common.c', 'input_simulator.c']
desktop_deps = [libsink]
desktop_args = ['-DCAMU_MIXER_THREADED', '-DCAMU_SCREEN_THREADED']
desktop = declare_dependency(sources: desktop_src, dependencies: desktop_deps,
diff --git a/src/util/queue.h b/src/util/queue.h
index 1d788cb..bd91342 100644
--- a/src/util/queue.h
+++ b/src/util/queue.h
@@ -9,6 +9,18 @@
struct nn_mutex mutex; \
}
+#define camu_queue_lock(q) \
+AL_MACRO_WRAP \
+({ \
+ nn_mutex_lock(&(q).mutex); \
+})
+
+#define camu_queue_unlock(q) \
+AL_MACRO_WRAP \
+({ \
+ nn_mutex_unlock(&(q).mutex); \
+})
+
#define camu_queue_size(q, r) \
AL_MACRO_WRAP \
({ \