summaryrefslogtreecommitdiff
path: root/src/buffer
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-04-09 11:24:01 -0400
committerAndrew Opalach <andrew@akon.city> 2024-04-09 11:24:01 -0400
commit02f3d3565602146bbbfce85b2719246f24036cb9 (patch)
treec6588ffe297b777e36260effa4fa42b958ca6ba3 /src/buffer
parentbbf3314165182e402ff25acccddc004a87f81ef0 (diff)
downloadcamu-02f3d3565602146bbbfce85b2719246f24036cb9.tar.gz
camu-02f3d3565602146bbbfce85b2719246f24036cb9.tar.bz2
camu-02f3d3565602146bbbfce85b2719246f24036cb9.zip
Massive restructure and many changes
- The server-side list concept is still a wip Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/buffer')
-rw-r--r--src/buffer/audio.c104
-rw-r--r--src/buffer/audio.h11
-rw-r--r--src/buffer/clock.c85
-rw-r--r--src/buffer/clock.h5
-rw-r--r--src/buffer/common_internal.h10
-rw-r--r--src/buffer/frame_queue.h4
-rw-r--r--src/buffer/video.c59
-rw-r--r--src/buffer/video.h14
8 files changed, 178 insertions, 114 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index e32ef1f..bfefed7 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -2,6 +2,7 @@
#include "audio.h"
#include "common.h"
+#include "common_internal.h"
#define BUFFER_USEC (12 * 1000000L)
#define BUFFER_WATERMARK_LOW (4 * 1000000L) // Must be a most half of the buffer size.
@@ -11,18 +12,11 @@
#ifdef CAMU_AUDIO_BUFFER_FADE
#define FADE 1.0
-#define FADE_LENGTH 18000
-#define FADE_STEP(rate) (FADE / FADE_LENGTH)
-//#define FADE_STEP(rate) ((FADE / FADE_LENGTH) / (FADE_LENGTH - 1.0))
+#define FADE_LENGTH(rate) (rate * 0.4)
+#define FADE_STEP(rate) (FADE / FADE_LENGTH(rate))
#endif
enum {
- FLOWING = 0,
- FLUSHED,
- SIGNALED
-};
-
-enum {
PAUSE_PRE = 0,
PAUSE_PAUSED,
#ifdef CAMU_AUDIO_BUFFER_FADE
@@ -31,8 +25,8 @@ enum {
PAUSE_PLAYING
};
-static bool setup_optimal_resampler(struct camu_audio_buffer *buf, struct camu_lav_resampler *resamp,
- struct camu_lav_resample_fmt *fmt)
+static bool setup_optimal_resampler(struct camu_audio_buffer *buf, struct camu_ff_resampler *resamp,
+ struct camu_ff_resample_fmt *fmt)
{
/*
switch ((enum AVSampleFormat)fmt->in_format) {
@@ -63,9 +57,9 @@ static bool setup_optimal_resampler(struct camu_audio_buffer *buf, struct camu_l
*/
buf->fade_rate = buf->fmt.req_sample_rate / 10;
- buf->bytes_per_sample = (s32)camu_lav_resample_fmt_bytes_per_sample(fmt);
+ buf->bytes_per_sample = (s32)camu_ff_resample_fmt_bytes_per_sample(fmt);
- return camu_lav_resampler_init(resamp, fmt);
+ return camu_ff_resampler_init(resamp, fmt);
}
bool camu_audio_buffer_init(struct camu_audio_buffer *buf, struct camu_clock *clock,
@@ -73,15 +67,16 @@ bool camu_audio_buffer_init(struct camu_audio_buffer *buf, struct camu_clock *cl
{
buf->mixer = mixer;
buf->clock = clock;
- buf->pause = PAUSE_PRE;
buf->ignore_desync = false;
+ al_atomic_size_t_store(&buf->continue_mark, 0, AL_ATOMIC_RELAXED);
buf->buffered = false;
+ al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
+ buf->pts = -1.0;
+ buf->pause = PAUSE_PRE;
#ifdef CAMU_AUDIO_BUFFER_FADE
buf->fade_period = 0;
buf->fade_offset = 0;
#endif
- al_atomic_size_t_store(&buf->continue_mark, 0, AL_ATOMIC_RELAXED);
- al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
#ifdef CAMU_MIXER_THREADED
al_atomic_bool_store(&buf->ref, false, AL_ATOMIC_RELAXED);
#endif
@@ -106,12 +101,12 @@ bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_stre
in_format_name, buf->fmt.in_channel_layout.nb_channels, buf->fmt.in_sample_rate,
req_format_name, buf->fmt.req_channel_layout.nb_channels, buf->fmt.req_sample_rate);
- buf->size = camu_lav_resample_fmt_usec_to_bytes(&buf->fmt, BUFFER_USEC);
+ buf->size = camu_ff_resample_fmt_usec_to_bytes(&buf->fmt, BUFFER_USEC);
buf->data = (u8 *)al_malloc(buf->size);
al_ring_buffer_init(&buf->rb, buf->data, buf->size);
- buf->watermark.low = camu_lav_resample_fmt_usec_to_bytes(&buf->fmt, BUFFER_WATERMARK_LOW);
- buf->watermark.high = camu_lav_resample_fmt_usec_to_bytes(&buf->fmt, BUFFER_WATERMARK_HIGH);
+ buf->watermark.low = camu_ff_resample_fmt_usec_to_bytes(&buf->fmt, BUFFER_WATERMARK_LOW);
+ buf->watermark.high = camu_ff_resample_fmt_usec_to_bytes(&buf->fmt, BUFFER_WATERMARK_HIGH);
camu_peak_buffer_init(&buf->peak);
buf->stream = stream;
@@ -119,35 +114,43 @@ bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_stre
return true;
}
+#ifdef _DEBUG_
+#define BUFFER_OCCUPIED_DEBUG(buf) \
+ camu_ff_resample_fmt_bytes_to_sec(&buf->fmt, buf->size - al_ring_buffer_space(&buf->rb))
+#else
+#define BUFFER_OCCUPIED_DEBUG(buf) 0
+#endif
+
static bool push_internal(struct camu_audio_buffer *buf, u8 **data, s32 sample_count, f64 pts)
{
if (buf->fmt.resampler_needed) {
- sample_count = camu_lav_resampler_convert(&buf->resamp, (const u8 **)data, sample_count);
- data = camu_lav_resampler_get_data(&buf->resamp);
+ sample_count = camu_ff_resampler_convert(&buf->resamp, (const u8 **)data, sample_count);
+ data = camu_ff_resampler_get_data(&buf->resamp);
}
+ if (sample_count <= 0) return false;
- f64 duration = camu_lav_resample_fmt_samples_to_sec(&buf->fmt, sample_count);
+ f64 duration = camu_ff_resample_fmt_samples_to_sec(&buf->fmt, sample_count);
if (pts + duration < camu_clock_get_base_pts(buf->clock)) {
return false;
}
+ if (buf->pts == -1.0) buf->pts = pts;
- if (sample_count <= 0) return false;
-
- size_t size, have = camu_lav_resample_fmt_samples_to_bytes(&buf->fmt, (size_t)sample_count);
+ size_t size, have = camu_ff_resample_fmt_samples_to_bytes(&buf->fmt, (size_t)sample_count);
size_t space = al_ring_buffer_space(&buf->rb);
if (!buf->buffered && buf->size - space > buf->watermark.high) {
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
buf->buffered = true;
+ al_log_debug("audio_buffer", "Buffered (watermark: %.2fs).", BUFFER_OCCUPIED_DEBUG(buf));
}
size_t peak = buf->peak.buf.size;
if (space < (have + peak)) {
camu_peak_buffer_push(&buf->peak, data[0], have);
if (peak >= buf->watermark.low) {
- // Discard the peak buffer. If this happens the user of this buffer is
- // taking _way_ too long to stop.
+ // Discard the peak buffer. If this happens the writer of this buffer
+ // is taking _way_ too long to stop.
camu_peak_buffer_flush(&buf->peak, &size);
}
- // Continuing based on an outdated continue_mark value is safe as long as the
+ // Continuing based on an outdated continue mark value is safe as long as the
// peak buffer is smaller than the low watermark and the low watermark is
// less than or equal to half the buffer size.
al_atomic_size_t_store(&buf->continue_mark, buf->watermark.low + peak, AL_ATOMIC_RELAXED);
@@ -165,7 +168,7 @@ static bool push_internal(struct camu_audio_buffer *buf, u8 **data, s32 sample_c
return true;
}
-#ifdef HAVE_FFMPEG
+#ifdef CAMU_HAVE_FFMPEG
static void push_av_frame_internal(struct camu_audio_buffer *buf, AVFrame *frame)
{
// It doesn't matter if there's a race here with flush(), as long
@@ -185,7 +188,7 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_frame *fr
switch (frame->type) {
case CAMU_NORMAL:
break;
-#ifdef HAVE_FFMPEG
+#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
push_av_frame_internal(buf, frame->av.frame);
break;
@@ -195,23 +198,32 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_frame *fr
al_free(frame);
}
+void camu_audio_buffer_unpause(struct camu_audio_buffer *buf)
+{
+ buf->pause = PAUSE_PRE;
+}
+
+// Not thread-safe, must be called while the buffer is not be read or written to.
void camu_audio_buffer_reset(struct camu_audio_buffer *buf)
{
al_atomic_size_t_store(&buf->continue_mark, 0, AL_ATOMIC_RELAXED);
+ buf->buffered = false;
al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
- buf->pts = camu_clock_get_base_pts(buf->clock) + camu_mixer_get_latency(buf->mixer);
+ buf->pts = -1.0;
buf->pause = PAUSE_PRE;
- buf->buffered = false;
al_ring_buffer_reset(&buf->rb);
}
+// flush() always comes from the same thread as push().
void camu_audio_buffer_flush(struct camu_audio_buffer *buf)
{
- al_atomic_u8_store(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
if (!buf->buffered) {
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
buf->buffered = true;
+ al_log_debug("audio_buffer", "Buffered (watermark: %.2fs).", BUFFER_OCCUPIED_DEBUG(buf));
}
+ al_atomic_u8_store(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
+ al_log_debug("audio_buffer", "Flush requested.");
}
#ifdef CAMU_AUDIO_BUFFER_FADE
@@ -263,13 +275,12 @@ static void handle_fade(struct camu_audio_buffer *buf, u8 *data, size_t size, bo
// The sole purpose of PAUSE_PRE is to avoid a fade in at the start of the stream.
size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t req)
{
- f64 pts = camu_clock_get_pts(buf->clock, camu_mixer_get_latency(buf->mixer));
- if (pts < 0.0) {
+ if (camu_clock_is_paused(buf->clock)) {
#ifdef CAMU_AUDIO_BUFFER_FADE
if (buf->pause == PAUSE_PLAYING) {
// Fade out. Signified by >0 fade_period and pause = PAUSE_FADING.
buf->volume = FADE;
- buf->fade_period = FADE_LENGTH;
+ buf->fade_period = FADE_LENGTH(buf->fade_rate);
buf->fade_offset = 0;
buf->pause = PAUSE_FADING;
} else if (buf->fade_period == 0) { // Fade out done.
@@ -283,7 +294,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
} else if (buf->pause == PAUSE_FADING || buf->pause == PAUSE_PAUSED) {
// Fade in. Signified by >0 fade_period and pause = PAUSE_PLAYING.
buf->volume = 1.0 - FADE;
- buf->fade_period = FADE_LENGTH;
+ buf->fade_period = FADE_LENGTH(buf->fade_rate);
buf->fade_offset = 0;
// We don't want to consider sync if reversing an active
// fade but we do if pause = PAUSE_PAUSED.
@@ -297,6 +308,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
return req;
#endif
}
+ f64 pts = camu_clock_get_pts(buf->clock, camu_mixer_get_latency(buf->mixer));
size_t ret, signal = req;
size_t size = al_ring_buffer_occupied(&buf->rb);
// Cut off fade if it's reaching too far.
@@ -309,20 +321,20 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
// Possibly 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_PRE || buf->pause == PAUSE_PAUSED))) {
- if (fabs(pts) >= 0.322) {
+ if (UNLIKELY(fabs(pts) >= 0.322 && buf->pause != PAUSE_PRE)) {
al_log_warn("audio_buffer", "Abnormally large audio desync of %.5fs", pts);
}
if (pts > 0.0) {
- ret = camu_lav_resample_fmt_sec_to_bytes(&buf->fmt, pts);
+ ret = camu_ff_resample_fmt_sec_to_bytes(&buf->fmt, pts);
ret = AL_MIN(ret, size);
al_log_debug("audio_buffer", "Skipping %.5fs of audio (%zu bytes).", pts, ret);
ret = al_ring_buffer_discard(&buf->rb, ret);
size -= ret;
- buf->pts += camu_lav_resample_fmt_bytes_to_sec(&buf->fmt, ret);
+ buf->pts += camu_ff_resample_fmt_bytes_to_sec(&buf->fmt, ret);
// Could go on to underrun.
} else if (pts < 0.0) {
pts = -pts;
- ret = camu_lav_resample_fmt_sec_to_bytes(&buf->fmt, pts);
+ ret = camu_ff_resample_fmt_sec_to_bytes(&buf->fmt, pts);
ret = AL_MIN(ret, req);
al_log_debug("audio_buffer", "Delaying audio by %.5fs (%zu bytes).", pts, ret);
al_memset(data, 0, ret);
@@ -341,9 +353,9 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
// If flow = FLUSHED, we know we've recieved all the data
// this stream wants to send. Signal that we are done.
signal = size;
+ buf->callback(buf->userdata, CAMU_BUFFER_EOF);
al_atomic_u8_store(&buf->flow, SIGNALED, AL_ATOMIC_RELAXED);
al_log_debug("audio_buffer", "Flushed (signal: %zu).", signal);
- buf->callback(buf->userdata, CAMU_BUFFER_EOF);
} else {
// Put silence into the remainder of the request buffer.
al_memset(data + size, 0, req - size);
@@ -368,7 +380,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
#endif
// Read as much as we can from the ring buffer.
ret = al_ring_buffer_read(&buf->rb, data, req);
- buf->pts += camu_lav_resample_fmt_bytes_to_sec(&buf->fmt, ret);
+ buf->pts += camu_ff_resample_fmt_bytes_to_sec(&buf->fmt, ret);
if (flow == FLOWING) {
ret = al_atomic_size_t_load(&buf->continue_mark, AL_ATOMIC_RELAXED);
// Request to uncork if we are below the mark.
@@ -378,9 +390,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
}
#ifdef CAMU_AUDIO_BUFFER_FADE
}
- if (buf->fade_period > 0) {
- handle_fade(buf, data, req, buf->pause == PAUSE_FADING);
- }
+ if (buf->fade_period > 0) handle_fade(buf, data, req, buf->pause == PAUSE_FADING);
#endif
}
// Never return less than req, unless we hit the end of the stream.
@@ -390,7 +400,7 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
void camu_audio_buffer_free(struct camu_audio_buffer *buf)
{
if (buf->fmt.resampler_needed) {
- camu_lav_resampler_close(&buf->resamp);
+ camu_ff_resampler_close(&buf->resamp);
}
if (buf->data) al_free(buf->data);
camu_peak_buffer_free(&buf->peak);
diff --git a/src/buffer/audio.h b/src/buffer/audio.h
index 1a8baf0..c4ba142 100644
--- a/src/buffer/audio.h
+++ b/src/buffer/audio.h
@@ -4,7 +4,7 @@
#include <al/ring_buffer.h>
#include "../codec/codec.h"
-#include "../codec/libav/resampler.h"
+#include "../codec/ffmpeg/resampler.h"
#include "../mixer/mixer.h"
#include "clock.h"
@@ -27,10 +27,8 @@ struct camu_audio_buffer {
s32 fade_rate;
s32 bytes_per_sample;
- struct camu_lav_resampler resamp;
- struct camu_lav_resample_fmt fmt;
-
- atomic_u8 flow;
+ struct camu_ff_resampler resamp;
+ struct camu_ff_resample_fmt fmt;
u8 *data;
size_t size;
@@ -41,6 +39,8 @@ struct camu_audio_buffer {
struct camu_peak_buffer peak;
atomic_size_t continue_mark;
+ atomic_u8 flow;
+
#ifdef CAMU_MIXER_THREADED
atomic_bool ref;
#endif
@@ -52,6 +52,7 @@ struct camu_audio_buffer {
bool camu_audio_buffer_init(struct camu_audio_buffer *buf, struct camu_clock *clock, struct camu_mixer *mixer);
bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_stream *stream);
void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_frame *frame);
+void camu_audio_buffer_unpause(struct camu_audio_buffer *buf);
void camu_audio_buffer_reset(struct camu_audio_buffer *buf);
void camu_audio_buffer_flush(struct camu_audio_buffer *buf);
size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t req);
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
index f8284ec..032e44a 100644
--- a/src/buffer/clock.c
+++ b/src/buffer/clock.c
@@ -4,57 +4,72 @@
void camu_clock_set(struct camu_clock *clock, u64 base, u64 ts, u64 delay)
{
- al_atomic_f64_store(&clock->base, base / 1000000.f, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->base, base / 1000000.0, AL_ATOMIC_RELAXED);
al_atomic_f64_store(&clock->tick, 0.0, AL_ATOMIC_RELAXED);
- al_atomic_u64_store(&clock->start, ts, AL_ATOMIC_RELAXED);
+ al_atomic_u64_store(&clock->start, ts + delay, AL_ATOMIC_RELAXED);
clock->delay = delay;
+ clock->pause_at = -1.0;
}
-bool camu_clock_calc_tick(struct camu_clock *clock)
+f64 calc_tick_internal(struct camu_clock *clock)
{
- u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED);
- if (start == 0L) return false;
- f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED);
- u64 ts = (tick == 0.0 && clock->delay == 0) ? start : aki_get_timestamp();
- al_assert(start <= ts);
- u64 diff = ts - start;
- if (tick == 0.0) {
- if (!(diff >= clock->delay && clock->delay >= (diff - clock->delay))) {
- // We are past the requested start time.
- return false;
- }
- diff = clock->delay - (diff - clock->delay);
- tick = aki_get_tick() - al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
- }
- al_atomic_f64_store(&clock->tick, tick + (diff / 1000000.f), AL_ATOMIC_RELAXED);
- return true;
+ s64 diff = clock->delay == 0 ? 0 : ((al_atomic_u64_load(&clock->start, AL_ATOMIC_ACQUIRE)) - aki_get_timestamp());
+ f64 tick = aki_get_tick() - al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
+ tick += diff / 1000000.0;
+ al_atomic_f64_store(&clock->tick, tick, AL_ATOMIC_RELAXED);
+ al_atomic_u64_store(&clock->start, 0, AL_ATOMIC_RELEASE);
+ return tick;
+}
+
+bool camu_clock_is_late(struct camu_clock *clock)
+{
+ s64 diff = (s64)(al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED)) - aki_get_timestamp();
+ return clock->delay > 0 && diff < 0;
}
void camu_clock_resume(struct camu_clock *clock)
{
- u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_ACQUIRE);
- if (start == 0L) return;
- al_atomic_u64_store(&clock->start, 0L, AL_ATOMIC_RELEASE);
+ f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_ACQUIRE);
+ if (tick != 0.0) return;
+ u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED);
+ if (start != 0L) {
+ // No longer report paused, calculate tick on first call to get_pts().
+ al_atomic_f64_store(&clock->tick, -1.0, AL_ATOMIC_RELEASE);
+ } else {
+ tick = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->tick, aki_get_tick() - tick, AL_ATOMIC_RELEASE);
+ }
}
-static f64 get_pts_internal(struct camu_clock *clock)
+void camu_clock_arm_resume(struct camu_clock *clock, u64 ts)
{
- f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED);
- return aki_get_tick() - tick;
+ al_atomic_u64_store(&clock->start, ts, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->tick, -1.0, AL_ATOMIC_RELAXED);
}
void camu_clock_pause(struct camu_clock *clock)
{
- u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_ACQUIRE);
- if (start != 0L) return;
- al_atomic_f64_store(&clock->base, get_pts_internal(clock), AL_ATOMIC_RELAXED);
- al_atomic_u64_store(&clock->start, aki_get_timestamp(), AL_ATOMIC_RELEASE);
+ f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_ACQUIRE);
+ if (tick == 0.0) return;
+ al_atomic_f64_store(&clock->base, aki_get_tick() - tick, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->tick, 0.0, AL_ATOMIC_RELEASE);
+}
+
+void camu_clock_arm_pause(struct camu_clock *clock, f64 pts)
+{
+ clock->pause_at = pts;
}
bool camu_clock_is_paused(struct camu_clock *clock)
{
- u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED);
- return start != 0L;
+ f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED);
+ if (clock->pause_at >= 0.0 && (aki_get_tick() - tick) >= clock->pause_at) {
+ al_atomic_f64_store(&clock->base, clock->pause_at, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->tick, 0.0, AL_ATOMIC_RELAXED);
+ clock->pause_at = -1.0;
+ return true;
+ }
+ return tick == 0.0;
}
f64 camu_clock_get_base_pts(struct camu_clock *clock)
@@ -64,10 +79,8 @@ f64 camu_clock_get_base_pts(struct camu_clock *clock)
f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset)
{
- u64 start = al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED);
- if (start != 0L) return -1.0;
- f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
- f64 pts = get_pts_internal(clock) + offset;
- if (pts < base) return -1.0;
+ f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED);
+ if (tick == -1.0) tick = calc_tick_internal(clock);
+ f64 pts = (aki_get_tick() - tick) + offset;
return pts;
}
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
index a7399b5..bdd50d8 100644
--- a/src/buffer/clock.h
+++ b/src/buffer/clock.h
@@ -7,12 +7,15 @@ struct camu_clock {
atomic_f64 base, tick;
atomic_u64 start;
u64 delay;
+ f64 pause_at;
};
void camu_clock_set(struct camu_clock *clock, u64 base, u64 ts, u64 delay);
-bool camu_clock_calc_tick(struct camu_clock *clock);
+bool camu_clock_is_late(struct camu_clock *clock);
void camu_clock_resume(struct camu_clock *clock);
+void camu_clock_arm_resume(struct camu_clock *clock, u64 ts);
void camu_clock_pause(struct camu_clock *clock);
+void camu_clock_arm_pause(struct camu_clock *clock, f64 pts);
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);
diff --git a/src/buffer/common_internal.h b/src/buffer/common_internal.h
new file mode 100644
index 0000000..0dc1e5f
--- /dev/null
+++ b/src/buffer/common_internal.h
@@ -0,0 +1,10 @@
+#pragma once
+
+enum {
+ // Flowing.
+ FLOWING,
+ // Got request to flush.
+ FLUSHED,
+ // Realized flush.
+ SIGNALED
+};
diff --git a/src/buffer/frame_queue.h b/src/buffer/frame_queue.h
index df58497..512b7f2 100644
--- a/src/buffer/frame_queue.h
+++ b/src/buffer/frame_queue.h
@@ -1,6 +1,6 @@
#pragma once
-#ifdef HAVE_FFMPEG
+#ifdef CAMU_HAVE_FFMPEG
#include <libavutil/frame.h>
#endif
@@ -15,7 +15,7 @@ enum {
struct camu_frame_queue {
void (*push)(struct camu_frame_queue *, struct camu_frame *, f64);
-#ifdef HAVE_FFMPEG
+#ifdef CAMU_HAVE_FFMPEG
void (*push_av_frame)(struct camu_frame_queue *, AVFrame *, f64);
#endif
void (*flush)(struct camu_frame_queue *);
diff --git a/src/buffer/video.c b/src/buffer/video.c
index 15af78e..ca8baab 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -2,6 +2,7 @@
#include "video.h"
#include "common.h"
+#include "common_internal.h"
//#define CAMU_VIDEO_BUFFER_FORCE_SCALER
@@ -10,12 +11,14 @@
#define BUFFER_WATERMARK_HIGH 14
#define BUFFER_WATERMARK_RESET BUFFER_WATERMARK_HIGH + 10.
-bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock,
+bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock, f64 offset,
struct camu_renderer *renderer)
{
buf->clock = clock;
+ buf->offset = offset;
buf->queue = renderer->create_queue(renderer);
buf->buffered = false;
+ al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
#ifdef CAMU_SCREEN_THREADED
al_atomic_bool_store(&buf->ref, false, AL_ATOMIC_RELAXED);
#endif
@@ -31,7 +34,7 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stre
buf->avg_frame_duration = 0.0;
break;
}
-#ifdef HAVE_FFMPEG
+#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
s32 width = stream->av.stream->codecpar->width;
s32 height = stream->av.stream->codecpar->height;
@@ -50,7 +53,7 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stre
buf->fmt.req_width = width;
buf->fmt.req_height = height;
buf->fmt.req_format = AV_PIX_FMT_RGBA;
- if (camu_lav_scaler_init(&buf->scale, &buf->fmt) && buf->fmt.scaler_needed) {
+ if (camu_ff_scaler_init(&buf->scale, &buf->fmt) && buf->fmt.scaler_needed) {
format_name = av_get_pix_fmt_name(buf->fmt.req_format);
al_log_info("video_buffer", "Scaling stream to: %s (%dx%d)", format_name,
buf->fmt.req_width, buf->fmt.req_height);
@@ -65,19 +68,27 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stre
return true;
}
+bool camu_video_buffer_is_single_frame(struct camu_video_buffer *buf)
+{
+ return buf->single_frame;
+}
+
static void after_push_internal(struct camu_video_buffer *buf)
{
s32 count = buf->queue->count(buf->queue);
if (!buf->buffered && (buf->single_frame || count >= BUFFER_WATERMARK_BUFFERED)) {
+ // Preserve order of, flush -> callback -> set flow, for single frames.
+ if (buf->single_frame) buf->queue->flush(buf->queue);
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
buf->buffered = true;
+ al_log_debug("video_buffer", "Buffered (watermark: %d frames).", count);
+ if (buf->single_frame) al_atomic_u8_store(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
}
if (count >= BUFFER_WATERMARK_RESET) buf->queue->reset(buf->queue);
else if (count >= BUFFER_WATERMARK_HIGH) buf->callback(buf->userdata, CAMU_BUFFER_CORK);
- if (buf->single_frame) camu_video_buffer_flush(buf);
}
-#ifdef HAVE_FFMPEG
+#ifdef CAMU_HAVE_FFMPEG
static void push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame)
{
f64 pts = 0.0;
@@ -91,7 +102,7 @@ static void push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame
}
#ifdef CAMU_VIDEO_BUFFER_FORCE_SCALER
if (buf->fmt.scaler_needed) {
- if (!camu_lav_scaler_scale(&buf->scale, (const u8 **)frame->data, frame->linesize)) {
+ if (!camu_ff_scaler_scale(&buf->scale, (const u8 **)frame->data, frame->linesize)) {
return;
}
av_frame_free(&frame);
@@ -105,11 +116,16 @@ static void push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame
void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_frame *frame)
{
+ // A single frame will be sent again after a seek, discard it here (for now).
+ if (buf->single_frame && buf->buffered) {
+ camu_frame_discard(frame);
+ return;
+ }
switch (frame->type) {
case CAMU_NORMAL:
buf->queue->push(buf->queue, frame, 0.0);
break;
-#ifdef HAVE_FFMPEG
+#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT:
push_av_frame_internal(buf, frame->av.frame);
al_free(frame);
@@ -119,34 +135,41 @@ void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_frame *fr
after_push_internal(buf);
}
+// Not thread-safe, must be called while the buffer is not be read or written to.
void camu_video_buffer_reset(struct camu_video_buffer *buf)
{
buf->queue->reset(buf->queue);
buf->buffered = false;
+ al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
}
-bool camu_video_buffer_is_single_frame(struct camu_video_buffer *buf)
-{
- return buf->single_frame;
-}
-
+// flush() always comes from the same thread as push().
void camu_video_buffer_flush(struct camu_video_buffer *buf)
{
buf->queue->flush(buf->queue);
if (!buf->buffered) {
buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
buf->buffered = true;
+ al_log_debug("video_buffer", "Buffered (watermark: %d frames).", buf->queue->count(buf->queue));
}
+ al_atomic_u8_store(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
}
bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out)
{
- f64 pts = buf->single_frame ? 0.0 : camu_clock_get_pts(buf->clock, 0);
- if (pts < 0.0) pts = camu_clock_get_base_pts(buf->clock);
- u8 ret = buf->queue->read(buf->queue, pts, out);
- if (ret == CAMU_QUEUE_EOF || (ret == CAMU_QUEUE_OK && buf->single_frame)) {
+ bool paused = buf->single_frame || camu_clock_is_paused(buf->clock);
+ f64 base = camu_clock_get_base_pts(buf->clock);
+ if (!paused) {
+ f64 pts = camu_clock_get_pts(buf->clock, buf->offset);
+ if (pts > base) base = pts;
+ }
+ u8 ret = buf->queue->read(buf->queue, base, out);
+ u8 flow = al_atomic_u8_load(&buf->flow, AL_ATOMIC_RELAXED);
+ if (flow == FLUSHED && (ret == CAMU_QUEUE_EOF || (buf->single_frame && ret == CAMU_QUEUE_OK))) {
buf->callback(buf->userdata, CAMU_BUFFER_EOF);
- } else if (buf->queue->count(buf->queue) <= BUFFER_WATERMARK_LOW) {
+ al_atomic_u8_store(&buf->flow, SIGNALED, AL_ATOMIC_RELAXED);
+ al_log_debug("video_buffer", "Flushed.");
+ } else if (flow == FLOWING && buf->queue->count(buf->queue) <= BUFFER_WATERMARK_LOW) {
buf->callback(buf->userdata, CAMU_BUFFER_UNCORK);
}
return ret == CAMU_QUEUE_OK || ret == CAMU_QUEUE_MORE;
@@ -155,7 +178,7 @@ bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out)
void camu_video_buffer_free(struct camu_video_buffer *buf)
{
#ifdef CAMU_VIDEO_BUFFER_FORCE_SCALER
- if (buf->fmt.scaler_needed) camu_lav_scaler_close(&buf->scale);
+ if (buf->fmt.scaler_needed) camu_ff_scaler_close(&buf->scale);
#endif
buf->queue->free(&buf->queue);
}
diff --git a/src/buffer/video.h b/src/buffer/video.h
index a6349f8..5b92dc6 100644
--- a/src/buffer/video.h
+++ b/src/buffer/video.h
@@ -3,7 +3,7 @@
#include <al/types.h>
#include "../codec/codec.h"
-#include "../codec/libav/scaler.h"
+#include "../codec/ffmpeg/scaler.h"
#include "../render/renderer.h"
#include "../screen/screen.h"
@@ -16,15 +16,18 @@ struct camu_video_buffer {
struct camu_clock *clock;
f32 start_time;
f32 avg_frame_duration;
+ f64 offset;
bool single_frame;
- struct camu_lav_scaler scale;
- struct camu_lav_scale_fmt fmt;
+ struct camu_ff_scaler scale;
+ struct camu_ff_scale_fmt fmt;
struct camu_frame_queue *queue;
bool buffered;
+ atomic_u8 flow;
+
#ifdef CAMU_SCREEN_THREADED
atomic_bool ref;
#endif
@@ -33,11 +36,12 @@ struct camu_video_buffer {
void *userdata;
};
-bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock, struct camu_renderer *renderer);
+bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock, f64 offset,
+ struct camu_renderer *renderer);
bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stream *stream);
+bool camu_video_buffer_is_single_frame(struct camu_video_buffer *buf);
void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_frame *frame);
void camu_video_buffer_reset(struct camu_video_buffer *buf);
-bool camu_video_buffer_is_single_frame(struct camu_video_buffer *buf);
void camu_video_buffer_flush(struct camu_video_buffer *buf);
bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out);
void camu_video_buffer_free(struct camu_video_buffer *buf);