summaryrefslogtreecommitdiff
path: root/src/buffer
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2025-10-30 15:24:58 -0400
committerAndrew Opalach <andrew@akon.city> 2025-10-30 15:24:58 -0400
commitf7e23d3c5e47ec0c105bf506e58e23f635faa20e (patch)
tree99dfc4f34da05d16bdfac6a437bb86a65d9b32f8 /src/buffer
parent90da3b27d939b3b7af1cf7fed10dfaaa7e271622 (diff)
downloadcamu-f7e23d3c5e47ec0c105bf506e58e23f635faa20e.tar.gz
camu-f7e23d3c5e47ec0c105bf506e58e23f635faa20e.tar.bz2
camu-f7e23d3c5e47ec0c105bf506e58e23f635faa20e.zip
Wip sink changes around errored/ended buffers
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/buffer')
-rw-r--r--src/buffer/audio.c38
-rw-r--r--src/buffer/audio.h2
-rw-r--r--src/buffer/common.h10
-rw-r--r--src/buffer/common_internal.h2
-rw-r--r--src/buffer/video.c66
-rw-r--r--src/buffer/video.h8
-rw-r--r--src/buffer/video_null.h6
7 files changed, 88 insertions, 44 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index 699b05e..377877e 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -10,7 +10,7 @@
#include "common.h"
#include "common_internal.h"
-#define BUFFER_SIZE 8.0
+#define BUFFER_SIZE 6.0
#define BUFFER_MARK_MIN 3.25 // Must be a most half of the buffer size.
#define BUFFER_MARK_BUFFERED 0.35
@@ -151,9 +151,9 @@ static bool push_internal(struct camu_audio_buffer *buf, f64 pts, u8 **data, s32
// The maximum space is buf->size - 1.
ptrdiff_t space = al_ring_buffer_space(&buf->rb);
if (!buf->buffered && (buf->size - 1) - space >= buf->mark.buffered) {
+ buf->buffered = true;
log_debug("Buffered (mark: %.1fKB).", buf->mark.buffered / 1024.0);
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
- buf->buffered = true;
}
ptrdiff_t have = (ptrdiff_t)camu_audio_format_samples_to_bytes(&buf->fmt.req, sample_count);
@@ -198,7 +198,15 @@ static void push_av_frame_internal(struct camu_audio_buffer *buf, AVFrame *frame
// A reset() must finish before any data is pushed.
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);
+ u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_RELAXED);
+ // flow could be ERRORED here.
+ if (flow != FLOWING) {
+ // Assert that push() is never called after flush().
+ al_assert(flow != FLUSHED);
+ camu_codec_frame_discard(frame);
+ return;
+ }
+
switch (frame->mode) {
case CAMU_NORMAL: {
s32 sample_count = frame->audio.sample_count;
@@ -214,22 +222,28 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_fra
}
#endif
}
+
al_free(frame);
+
+#ifdef CAMU_BUFFER_SPORADIC_ERRORS
+ ROLL_FOR_BUFFER_ERROR(buf);
+#endif
}
// flush() always comes from the same thread as push().
-void camu_audio_buffer_flush(struct camu_audio_buffer *buf)
+void camu_audio_buffer_flush(struct camu_audio_buffer *buf, bool error)
{
log_debug("Flush requested.");
+ u8 flow = error ? FLUSHED_ERROR : FLUSHED;
+ al_atomic_store(u8)(&buf->flow, flow, AL_ATOMIC_RELAXED);
if (!push_internal(buf, 0.0, NULL, 0)) {
log_debug("Buffer filled by flush.");
}
if (!buf->buffered) {
+ buf->buffered = true;
log_debug("Buffered (flush).");
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
- buf->buffered = true;
}
- al_atomic_store(u8)(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
}
// Not thread-safe, must be called while the buffer is not being read from or written to.
@@ -250,8 +264,19 @@ void camu_audio_buffer_resync(struct camu_audio_buffer *buf)
ptrdiff_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, ptrdiff_t req)
{
+ // Assert this buffer isn't being read before we signaled BUFFER_BUFFERED.
al_assert(buf->buffered);
+ u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE);
+ if (UNLIKELY(flow == ERRORED || flow == FLUSHED_ERROR)) {
+ if (flow == FLUSHED_ERROR) {
+ buf->callback(buf->userdata, CAMU_BUFFER_ERRORED);
+ al_atomic_store(u8)(&buf->flow, ERRORED, AL_ATOMIC_RELEASE);
+ }
+ al_memset(data, 0, req);
+ return req;
+ }
+
struct camu_audio_format *fmt = &buf->fmt.req;
ptrdiff_t ret, signal = req;
@@ -354,7 +379,6 @@ ptrdiff_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, ptrdif
buf->pause = PAUSE_PLAYING;
}
- u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE);
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.
diff --git a/src/buffer/audio.h b/src/buffer/audio.h
index 78a12dc..a740c6c 100644
--- a/src/buffer/audio.h
+++ b/src/buffer/audio.h
@@ -68,7 +68,7 @@ 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_flush(struct camu_audio_buffer *buf, bool error);
void camu_audio_buffer_reset(struct camu_audio_buffer *buf);
void camu_audio_buffer_resync(struct camu_audio_buffer *buf);
ptrdiff_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, ptrdiff_t req);
diff --git a/src/buffer/common.h b/src/buffer/common.h
index 24964ba..bd9ca45 100644
--- a/src/buffer/common.h
+++ b/src/buffer/common.h
@@ -1,5 +1,15 @@
#pragma once
+//#define CAMU_BUFFER_SPORADIC_ERRORS
+#ifdef CAMU_BUFFER_SPORADIC_ERRORS
+#include <al/random.h>
+#define ROLL_FOR_BUFFER_ERROR(buf) do { \
+ if (al_random_int(0, 254) == 72) { \
+ al_atomic_store(u8)(&(buf)->flow, FLUSHED_ERROR, AL_ATOMIC_RELAXED); \
+ } \
+} while (0)
+#endif
+
enum {
CAMU_BUFFER_BUFFERED = 0,
CAMU_BUFFER_CORK,
diff --git a/src/buffer/common_internal.h b/src/buffer/common_internal.h
index 700e278..a9e6ece 100644
--- a/src/buffer/common_internal.h
+++ b/src/buffer/common_internal.h
@@ -9,6 +9,8 @@ enum {
FLOWING,
// Got request to flush.
FLUSHED,
+ // Flushed because of an error.
+ FLUSHED_ERROR,
// Realized flush.
SIGNALED,
// Errored.
diff --git a/src/buffer/video.c b/src/buffer/video.c
index 0e9ce28..d16a106 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -26,7 +26,7 @@ bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *cl
buf->clock = clock;
buf->latency = 0.0;
al_atomic_store(f64)(&buf->pts, -1.0, AL_ATOMIC_RELAXED);
- buf->reset_pts = -1.0;
+ buf->seek_pts = -1.0;
buf->single_frame = true;
buf->queue = NULL;
buf->buffered = false;
@@ -117,27 +117,6 @@ void camu_video_buffer_set_latency(struct camu_video_buffer *buf, f64 latency)
buf->latency = latency;
}
-static void after_push_internal(struct camu_video_buffer *buf)
-{
- // Note that queue->reset() must only be called from the read() thread.
- // If a reset where to happen from the this thread, the latest read()
- // frame could be freed before it was used.
- s32 count = buf->queue->count(buf->queue);
- if (!buf->buffered && (buf->single_frame || count >= BUFFER_MARK_BUFFERED)) {
- // Preserve order of: set flow -> flush -> callback, for single frames.
- if (buf->single_frame) {
- al_atomic_store(u8)(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
- buf->queue->flush(buf->queue);
- }
- buf->buffered = true;
- buf->buffered_with_one_frame = count == 1;
- log_debug("Buffered (mark: %.2fs).", count * buf->avg_frame_duration);
- buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
- } else if (count >= BUFFER_MARK_HIGH) {
- buf->callback(buf->userdata, CAMU_BUFFER_CORK);
- }
-}
-
#ifdef CAMU_HAVE_FFMPEG
static bool push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame)
{
@@ -163,11 +142,16 @@ static bool push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame
void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_frame *frame)
{
- if (buf->single_frame && buf->buffered) {
- log_warn("Unexpected duplicate frame received.");
+ u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE);
+ if (flow != FLOWING) {
+ // A single frame will be FLUSHED after any push().
+ if (buf->single_frame) {
+ log_warn("Unexpected duplicate frame received.");
+ }
camu_codec_frame_discard(frame);
return;
}
+
switch (frame->mode) {
case CAMU_NORMAL: {
f64 base_pts = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
@@ -188,7 +172,25 @@ void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_fra
}
#endif
}
- after_push_internal(buf);
+
+ s32 count = buf->queue->count(buf->queue);
+ if (!buf->buffered && (buf->single_frame || count >= BUFFER_MARK_BUFFERED)) {
+ // Preserve order of: set flow -> flush -> callback, for single frames.
+ if (buf->single_frame) {
+ al_atomic_store(u8)(&buf->flow, FLUSHED, AL_ATOMIC_RELEASE);
+ buf->queue->flush(buf->queue);
+ }
+ buf->buffered = true;
+ buf->buffered_with_one_frame = count == 1;
+ log_debug("Buffered (mark: %.2fs).", count * buf->avg_frame_duration);
+ buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
+ } else if (count >= BUFFER_MARK_HIGH) {
+ buf->callback(buf->userdata, CAMU_BUFFER_CORK);
+ }
+
+#ifdef CAMU_BUFFER_SPORADIC_ERRORS
+ ROLL_FOR_BUFFER_ERROR(buf);
+#endif
}
void camu_video_buffer_push_subtitle(struct camu_video_buffer *buf, struct camu_codec_packet *packet)
@@ -197,10 +199,11 @@ void camu_video_buffer_push_subtitle(struct camu_video_buffer *buf, struct camu_
}
// flush() always comes from the same thread as push().
-void camu_video_buffer_flush(struct camu_video_buffer *buf)
+void camu_video_buffer_flush(struct camu_video_buffer *buf, bool error)
{
log_debug("Flush requested.");
- al_atomic_store(u8)(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
+ u8 flow = error ? FLUSHED_ERROR : FLUSHED;
+ al_atomic_store(u8)(&buf->flow, flow, AL_ATOMIC_RELAXED);
buf->queue->flush(buf->queue);
if (!buf->buffered) {
s32 count = buf->queue->count(buf->queue);
@@ -212,9 +215,11 @@ void camu_video_buffer_flush(struct camu_video_buffer *buf)
}
// Not thread-safe, must be called while the buffer is not being read from or written to.
+// If a queue->reset() happened from the push() thread while the buffer was active,
+// the latest read() frame could be freed before it was used.
void camu_video_buffer_reset(struct camu_video_buffer *buf, f64 pts)
{
- buf->reset_pts = pts;
+ buf->seek_pts = pts;
al_atomic_store(f64)(&buf->pts, -1.0, AL_ATOMIC_RELAXED);
if (buf->queue) buf->queue->reset(buf->queue);
buf->buffered = false;
@@ -223,6 +228,7 @@ void camu_video_buffer_reset(struct camu_video_buffer *buf, f64 pts)
bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weighted)
{
+ // Assert this buffer isn't being read before we signaled BUFFER_BUFFERED.
al_assert(buf->buffered);
f64 base_pts = al_atomic_load(f64)(&buf->pts, AL_ATOMIC_ACQUIRE);
@@ -232,9 +238,9 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weig
}
u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE);
- if (flow == ERRORED) return false;
+ if (flow == ERRORED) return false;
u8 ret = buf->queue->read(buf->queue, base_pts, out);
- if (ret == CAMU_QUEUE_ERR) {
+ if (ret == CAMU_QUEUE_ERR || flow == FLUSHED_ERROR) {
buf->callback(buf->userdata, CAMU_BUFFER_ERRORED);
al_atomic_store(u8)(&buf->flow, ERRORED, AL_ATOMIC_RELEASE);
return false;
diff --git a/src/buffer/video.h b/src/buffer/video.h
index 380943e..fdd7d0b 100644
--- a/src/buffer/video.h
+++ b/src/buffer/video.h
@@ -17,7 +17,7 @@ struct camu_video_buffer {
struct camu_clock *clock;
f64 latency;
atomic(f64) pts;
- f64 reset_pts;
+ f64 seek_pts;
bool single_frame;
f64 avg_frame_duration;
@@ -27,7 +27,7 @@ struct camu_video_buffer {
struct camu_frame_queue *queue;
bool buffered;
- // We need to manually trigger EOF in read().
+ // Need to manually trigger EOF in read().
bool buffered_with_one_frame;
// Flush the renderer on this read and don't allow it to set the clock.
bool weighted_read;
@@ -38,7 +38,7 @@ struct camu_video_buffer {
atomic(u8) ref;
#endif
- // Previous view, set from screen.
+ // Previous view, set from screen::add_buffer_internal().
struct camu_view view;
void (*callback)(void *, u8);
@@ -52,7 +52,7 @@ bool camu_video_buffer_configure_subtitles(struct camu_video_buffer *buf, struct
void camu_video_buffer_set_latency(struct camu_video_buffer *buf, f64 latency);
void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_frame *frame);
void camu_video_buffer_push_subtitle(struct camu_video_buffer *buf, struct camu_codec_packet *packet);
-void camu_video_buffer_flush(struct camu_video_buffer *buf);
+void camu_video_buffer_flush(struct camu_video_buffer *buf, bool error);
void camu_video_buffer_reset(struct camu_video_buffer *buf, f64 pts);
bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out, bool *weighted);
void camu_video_buffer_free(struct camu_video_buffer *buf);
diff --git a/src/buffer/video_null.h b/src/buffer/video_null.h
index 06b284c..f82a4d3 100644
--- a/src/buffer/video_null.h
+++ b/src/buffer/video_null.h
@@ -8,6 +8,7 @@
struct camu_video_buffer {
struct camu_codec_stream *stream;
+ f64 seek_pts;
bool single_frame;
f64 avg_frame_duration;
#ifdef CAMU_SCREEN_THREADED
@@ -47,7 +48,7 @@ static bool camu_video_buffer_configure_subtitles(struct camu_video_buffer *buf,
{
(void)buf;
(void)stream;
- return true;
+ return false;
}
static void camu_video_buffer_set_latency(struct camu_video_buffer *buf, s32 frames)
@@ -69,9 +70,10 @@ static void camu_video_buffer_push_subtitle(struct camu_video_buffer *buf, struc
(void)packet;
}
-static void camu_video_buffer_flush(struct camu_video_buffer *buf)
+static void camu_video_buffer_flush(struct camu_video_buffer *buf, bool error)
{
(void)buf;
+ (void)error;
}
static void camu_video_buffer_reset(struct camu_video_buffer *buf, f64 pts)