summaryrefslogtreecommitdiff
path: root/src/buffer
diff options
context:
space:
mode:
authorAndrew Opalach <andrew@akon.city> 2024-10-21 19:22:50 -0400
committerAndrew Opalach <andrew@akon.city> 2024-10-21 19:22:50 -0400
commit60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2 (patch)
tree08ff2ce7975f523112e7ad2fe4f797b4fc7db5de /src/buffer
parent2f9a0945bfeee3296cec3d38d094e4c49f9cb65f (diff)
downloadcamu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.tar.gz
camu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.tar.bz2
camu-60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2.zip
Everything before initial synced list
Signed-off-by: Andrew Opalach <andrew@akon.city>
Diffstat (limited to 'src/buffer')
-rw-r--r--src/buffer/audio.c427
-rw-r--r--src/buffer/audio.h36
-rw-r--r--src/buffer/clock.c151
-rw-r--r--src/buffer/clock.h56
-rw-r--r--src/buffer/frame_queue.h2
-rw-r--r--src/buffer/peak_buffer.c12
-rw-r--r--src/buffer/peak_buffer.h5
-rw-r--r--src/buffer/video.c74
-rw-r--r--src/buffer/video.h29
-rw-r--r--src/buffer/volume.h36
10 files changed, 444 insertions, 384 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
index 644be42..1bf42bf 100644
--- a/src/buffer/audio.c
+++ b/src/buffer/audio.c
@@ -1,169 +1,139 @@
#include <al/log.h>
+#ifdef CAMU_HAVE_FFMPEG
+#include "../codec/ffmpeg/resampler.h"
+#endif
+
#include "audio.h"
#include "common.h"
#include "common_internal.h"
+#include "volume.h"
-#define BUFFER_USEC (12 * 1000000L)
-#define BUFFER_WATERMARK_LOW (4 * 1000000L) // Must be a most half of the buffer size.
-#define BUFFER_WATERMARK_HIGH (3 * 1000000L)
+#define BUFFER_USEC (3 * 1000000L)
+#define BUFFER_MARK_MIN (1.1 * 1000000L) // Must be a most half of the buffer size.
+#define BUFFER_MARK_BUFFERED (1.0 * 1000000L)
-#define DESYNC_PTS 0.022
+#define LARGE_DESYNC_PTS 0.322
#ifdef CAMU_AUDIO_BUFFER_FADE
-#define FADE 1.0
-#define FADE_LENGTH(rate) (rate * 0.4)
-#define FADE_STEP(rate) (FADE / FADE_LENGTH(rate))
+#define FADE_STEP(fmt) (1.f / (fmt)->sample_rate)
#endif
enum {
- PAUSE_PRE = 0,
- PAUSE_PAUSED,
+ PAUSE_PAUSED = 0,
#ifdef CAMU_AUDIO_BUFFER_FADE
PAUSE_FADING,
#endif
PAUSE_PLAYING
};
-static bool setup_optimal_resampler(struct camu_audio_buffer *buf, struct camu_ff_resampler *resamp,
- struct camu_ff_resample_fmt *fmt)
+static void reset_buffer_state(struct camu_audio_buffer *buf)
{
- /*
- switch ((enum AVSampleFormat)fmt->in_format) {
- case AV_SAMPLE_FMT_S16P:
- fmt->req_format = AV_SAMPLE_FMT_S16;
- break;
- case AV_SAMPLE_FMT_S32P:
- fmt->req_format = AV_SAMPLE_FMT_S32;
- break;
- case AV_SAMPLE_FMT_DBL:
- case AV_SAMPLE_FMT_DBLP:
- case AV_SAMPLE_FMT_FLTP:
- fmt->req_format = AV_SAMPLE_FMT_FLT;
- break;
- case AV_SAMPLE_FMT_S16:
- case AV_SAMPLE_FMT_S32:
- case AV_SAMPLE_FMT_FLT:
- fmt->req_format = fmt->in_format;
- break;
- default:
- al_log_warn("audio_buffer", "Unhandled sample format.");
- break;
- }
-
- fmt->req_format = fmt->in_format;
- fmt->req_channel_count = fmt->in_channel_count;
- fmt->req_sample_rate = fmt->in_sample_rate;
- */
-
- buf->fade_rate = buf->fmt.req_sample_rate / 10;
- buf->bytes_per_sample = (s32)camu_ff_resample_fmt_bytes_per_sample(fmt);
-
- return camu_ff_resampler_init(resamp, fmt);
+ buf->pts = -1.0;
+ buf->pause = PAUSE_PAUSED;
+ buf->volume = 1.f;
+#ifdef CAMU_AUDIO_BUFFER_FADE
+ buf->fade_offset = 0;
+#endif
+ buf->buffered = false;
+ al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
+ al_atomic_store(size_t)(&buf->uncork_at, 0, AL_ATOMIC_RELAXED);
}
bool camu_audio_buffer_init(struct camu_audio_buffer *buf, struct camu_clock *clock,
struct camu_mixer *mixer)
{
- buf->mixer = mixer;
buf->clock = clock;
+ buf->mixer = mixer;
+ reset_buffer_state(buf);
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
+ buf->latency = 0.0;
#ifdef CAMU_MIXER_THREADED
- al_atomic_bool_store(&buf->ref, false, AL_ATOMIC_RELAXED);
+ al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
#endif
return true;
}
-bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_stream *stream)
+bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_codec_stream *stream)
{
- buf->fmt.in_format = stream->av.stream->codecpar->format;
- av_channel_layout_copy(&buf->fmt.in_channel_layout, &stream->av.stream->codecpar->ch_layout);
- buf->fmt.in_channel_count = buf->fmt.in_channel_layout.nb_channels;
- buf->fmt.in_sample_rate = stream->av.stream->codecpar->sample_rate;
-
+ camu_audio_format_copy(&buf->fmt.in, &stream->audio.fmt);
camu_mixer_pick_format(buf->mixer, &buf->fmt);
- if (!setup_optimal_resampler(buf, &buf->resamp, &buf->fmt)) {
+ buf->fmt.resampler_needed = !camu_resampler_format_matches(&buf->fmt);
+ if (buf->fmt.resampler_needed) {
+#ifdef CAMU_HAVE_FFMPEG
+ buf->resamp = camu_ff_resampler_create();
+ if (!buf->resamp->init(buf->resamp, &buf->fmt)) {
+ return false;
+ }
+#else
return false;
+#endif
}
- const char *in_format_name = av_get_sample_fmt_name(buf->fmt.in_format);
- const char *req_format_name = av_get_sample_fmt_name(buf->fmt.req_format);
+ const char *in_format_name = camu_audio_format_name(buf->fmt.in.format);
+ const char *req_format_name = camu_audio_format_name(buf->fmt.req.format);
al_log_info("audio_buffer", "Stream: %s (%dch) %dHz -> %s (%dch) %dHz.",
- 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);
+ in_format_name, buf->fmt.in.channel_count, buf->fmt.in.sample_rate,
+ req_format_name, buf->fmt.req.channel_count, buf->fmt.req.sample_rate);
- buf->size = camu_ff_resample_fmt_usec_to_bytes(&buf->fmt, BUFFER_USEC);
+ buf->size = camu_audio_format_usec_to_bytes(&buf->fmt.req, BUFFER_USEC);
buf->data = (u8 *)al_malloc(buf->size);
al_ring_buffer_init(&buf->rb, buf->data, buf->size);
- 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->mark.min = camu_audio_format_usec_to_bytes(&buf->fmt.req, BUFFER_MARK_MIN);
+ buf->mark.buffered = camu_audio_format_usec_to_bytes(&buf->fmt.req, BUFFER_MARK_BUFFERED);
+ camu_peak_buffer_init(&buf->peak, 1024 * 16);
buf->stream = stream;
return true;
}
+void camu_audio_buffer_set_latency(struct camu_audio_buffer *buf, f64 latency)
+{
+ buf->latency = latency;
+}
+
#ifdef _DEBUG_
-#define BUFFER_OCCUPIED_DEBUG(buf) \
- camu_ff_resample_fmt_bytes_to_sec(&buf->fmt, al_ring_buffer_occupied(&buf->rb))
+#define OCCUPIED_SECONDS_DEBUG(buf) \
+ camu_audio_format_bytes_to_sec(&buf->fmt.req, al_ring_buffer_occupied(&buf->rb))
#else
-#define BUFFER_OCCUPIED_DEBUG(buf) 0
+#define OCCUPIED_SECONDS_DEBUG(buf) 0
#endif
-static bool push_internal(struct camu_audio_buffer *buf, u8 **data, s32 sample_count, f64 pts)
+static bool push_internal(struct camu_audio_buffer *buf, u8 *data, s32 sample_count)
{
- if (buf->fmt.resampler_needed) {
- 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_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;
-
- 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) {
+ if (!buf->buffered && buf->size - space > buf->mark.buffered) {
+ al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", OCCUPIED_SECONDS_DEBUG(buf));
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 writer of this buffer
- // is taking _way_ too long to stop.
- camu_peak_buffer_flush(&buf->peak, &size);
+
+ size_t have = camu_audio_format_samples_to_bytes(&buf->fmt.req, (size_t)sample_count);
+ size_t peak = camu_peak_buffer_size(&buf->peak);
+ if (space < have + peak) {
+ al_assert(have < buf->mark.min);
+ camu_peak_buffer_push(&buf->peak, data, have);
+ peak += have;
+ if (peak >= buf->mark.min) {
+ // If this happens the writer of this buffer is taking way too long to stop.
+ al_log_warn("audio_buffer", "Overrun likely, discarding peak buffer.");
+ camu_peak_buffer_flush(&buf->peak);
+ peak = 0;
}
- // 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);
+ al_assert(peak <= buf->mark.min);
+ // We adjust the min mark by the peak buffer size just for consistency.
+ al_atomic_store(size_t)(&buf->uncork_at, buf->mark.min - peak, AL_ATOMIC_RELAXED);
buf->callback(buf->userdata, CAMU_BUFFER_CORK);
return false;
}
if (peak > 0) {
- u8 *ptr = camu_peak_buffer_flush(&buf->peak, &size);
- al_ring_buffer_write(&buf->rb, ptr, size);
+ al_ring_buffer_write(&buf->rb, camu_peak_buffer_flush(&buf->peak), peak);
}
- al_ring_buffer_write(&buf->rb, data[0], have);
+ al_ring_buffer_write(&buf->rb, data, have);
return true;
}
@@ -171,23 +141,40 @@ static bool push_internal(struct camu_audio_buffer *buf, u8 **data, s32 sample_c
#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
- // as any pushed frame is from the same stream.
- // A reset() must finish before any data is pushed.
- if (al_atomic_u8_load(&buf->flow, AL_ATOMIC_RELAXED) == FLOWING) {
- AVStream *stream = buf->stream->av.stream;
- f64 pts = (frame->best_effort_timestamp - stream->start_time) * av_q2d(stream->time_base);
- push_internal(buf, frame->data, frame->nb_samples, pts);
+ s32 sample_count = frame->nb_samples;
+ AVStream *stream = buf->stream->av.stream;
+ f64 pts = frame->best_effort_timestamp * av_q2d(stream->time_base);
+ f64 duration = camu_audio_format_samples_to_sec(&buf->fmt.in, frame->nb_samples);
+ if (pts + duration >= camu_clock_get_base_pts(buf->clock)) {
+ if (buf->pts == -1.0) buf->pts = pts;
+ u8 **data = frame->data;
+ if (buf->fmt.resampler_needed) {
+ sample_count = buf->resamp->convert(buf->resamp, (const u8 **)data, sample_count);
+ data = buf->resamp->get_data(buf->resamp);
+ }
+ if (sample_count > 0) push_internal(buf, data[0], sample_count);
}
av_frame_free(&frame);
}
#endif
-void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_frame *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)
{
- switch (frame->type) {
- case CAMU_NORMAL:
+ al_assert(al_atomic_load(u8)(&buf->flow, AL_ATOMIC_RELAXED) == FLOWING);
+ switch (frame->mode) {
+ case CAMU_NORMAL: {
+ u8 *store[AV_NUM_DATA_POINTERS] = { 0 };
+ store[0] = frame->data;
+ u8 **data = store;
+ s32 sample_count = frame->audio.sample_count;
+ if (buf->fmt.resampler_needed) {
+ sample_count = buf->resamp->convert(buf->resamp, (const u8 **)store, sample_count);
+ data = buf->resamp->get_data(buf->resamp);
+ }
+ push_internal(buf, data[0], sample_count);
break;
+ }
#ifdef CAMU_HAVE_FFMPEG
case CAMU_FFMPEG_COMPAT: {
push_av_frame_internal(buf, frame->av.frame);
@@ -198,92 +185,56 @@ void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_frame *fr
al_free(frame);
}
+// Not thread-safe, must be called while the buffer is not being read from or written to.
void camu_audio_buffer_unpause(struct camu_audio_buffer *buf)
{
- buf->pause = PAUSE_PRE;
+ buf->pause = PAUSE_PAUSED;
}
-// Not thread-safe, must be called while the buffer is not being read from or written to.
+// Not thread-safe.
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 = -1.0;
- buf->pause = PAUSE_PRE;
+ reset_buffer_state(buf);
al_ring_buffer_reset(&buf->rb);
+ camu_peak_buffer_flush(&buf->peak);
+ if (buf->fmt.resampler_needed) {
+ buf->resamp->flush(buf->resamp);
+ }
}
// flush() always comes from the same thread as push().
void camu_audio_buffer_flush(struct camu_audio_buffer *buf)
{
+ if (buf->fmt.resampler_needed) {
+ s32 sample_count = buf->resamp->flush(buf->resamp);
+ if (sample_count > 0) {
+ u8 **data = buf->resamp->get_data(buf->resamp);
+ if (!push_internal(buf, data[0], sample_count)) {
+ al_log_debug("audio_buffer", "Buffer filled by resampler flush.");
+ }
+ }
+ }
if (!buf->buffered) {
+ al_log_debug("audio_buffer", "Buffered (mark: %.2fs).", OCCUPIED_SECONDS_DEBUG(buf));
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_atomic_store(u8)(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
al_log_debug("audio_buffer", "Flush requested.");
}
-#ifdef CAMU_AUDIO_BUFFER_FADE
-#define MAX_BYTES 8
-#define TYPED_CAST(f, type) (f64)(*((type *)f))
-#define TYPED_CLAMP(f, type, min, max, ret) *((type *)ret) = (type)AL_CLAMP(f, min, max)
+#define NOT_PAUSED(pause) (pause != PAUSE_PAUSED)
-// This is obviously bad for optimization, should just make a function for each type.
-static void handle_fade(struct camu_audio_buffer *buf, u8 *data, size_t size, bool out)
-{
- u8 bytes[MAX_BYTES];
- u8 *ptr = bytes + (MAX_BYTES - buf->bytes_per_sample);
- f64 value;
- u32 sample_count = size / buf->bytes_per_sample;
- u32 channel_count = 2;
- for (u32 i = 0; i < sample_count; i++) {
- al_memcpy(ptr, data, buf->bytes_per_sample);
- switch (buf->bytes_per_sample) {
- case 2:
- value = TYPED_CAST(ptr, s16);
- break;
- case 4:
- value = TYPED_CAST(ptr, s32);
- break;
- }
- value *= (buf->volume);
- switch (buf->bytes_per_sample) {
- case 2:
- TYPED_CLAMP(value, s16, INT16_MIN, INT16_MAX, ptr);
- break;
- case 4:
- TYPED_CLAMP(value, s32, INT32_MIN, INT32_MAX, ptr);
- break;
- }
- al_memcpy(data, ptr, buf->bytes_per_sample);
- if ((i + 1) % channel_count == 0) {
- if (out) buf->volume = AL_CLAMP((buf->volume - FADE_STEP(buf->fade_rate)), 0.0, 1.0);
- else buf->volume = AL_CLAMP((buf->volume + FADE_STEP(buf->fade_rate)), 0.0, 1.0);
- if (buf->fade_period > 0) buf->fade_period--;
- }
- data += buf->bytes_per_sample;
- }
-}
-#endif
-
-#define NOT_PAUSED(pause) (pause != PAUSE_PAUSED && pause != PAUSE_PRE)
-
-// PAUSE_PAUSED means the last read was silence, meaning it's safe to skip.
-// The sole purpose of PAUSE_PRE is to avoid a fade in at the start of the stream.
+// 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)
{
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_rate);
- buf->fade_offset = 0;
buf->pause = PAUSE_FADING;
- } else if (buf->fade_period == 0) { // Fade out done.
+ buf->fade_offset = 0;
+ } else if (buf->volume == 0.f) { // Fade out done.
al_memset(data, 0, req);
if (NOT_PAUSED(buf->pause)) {
buf->pause = PAUSE_PAUSED;
@@ -292,12 +243,6 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
return req;
}
} 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_rate);
- buf->fade_offset = 0;
- // We don't want to consider sync if reversing an active
- // fade but we do if pause = PAUSE_PAUSED.
if (buf->pause == PAUSE_FADING) buf->pause = PAUSE_PLAYING;
#else
al_memset(data, 0, req);
@@ -308,66 +253,87 @@ 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));
+ u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_ACQUIRE);
+ if (flow == SIGNALED) {
+ buf->callback(buf->userdata, CAMU_BUFFER_EOF);
+ return 0;
+ }
+ f64 pts = camu_clock_get_pts(buf->clock, buf->latency);
size_t ret, signal = req;
- size_t size = al_ring_buffer_occupied(&buf->rb);
- // Cut off fade if it's reaching too far.
- if (size < buf->fade_offset) size = 0;
- else size -= buf->fade_offset;
+ 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 -= buf->pts;
- // 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 (UNLIKELY(fabs(pts) >= 0.322 && buf->pause != PAUSE_PRE)) {
+ // Attempt syncing to the clock.
+ // For this to work the mixer must report a reasonably accurate value for latency.
+ if (UNLIKELY(!buf->ignore_desync && buf->pause == PAUSE_PAUSED)) {
+ if (UNLIKELY(fabs(pts) >= LARGE_DESYNC_PTS)) {
al_log_warn("audio_buffer", "Abnormally large audio desync of %.5fs", pts);
}
- if (pts > 0.0) {
- 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);
+ if (pts > 0.0) { // Skip.
+ ret = camu_audio_format_sec_to_bytes(&buf->fmt.req, pts);
+ ret = AL_MIN(ret, have);
+ al_log_debug("audio_buffer", "Skipping %fs of audio (%zu bytes).", pts, ret);
ret = al_ring_buffer_discard(&buf->rb, ret);
- size -= ret;
- buf->pts += camu_ff_resample_fmt_bytes_to_sec(&buf->fmt, ret);
+ have -= ret;
+ buf->pts += camu_audio_format_bytes_to_sec(&buf->fmt.req, ret);
// Could go on to underrun.
- } else if (pts < 0.0) {
+ } else if (pts < 0.0) { // Delay.
pts = -pts;
- ret = camu_ff_resample_fmt_sec_to_bytes(&buf->fmt, pts);
+ ret = camu_audio_format_sec_to_bytes(&buf->fmt.req, pts);
ret = AL_MIN(ret, req);
- al_log_debug("audio_buffer", "Delaying audio by %.5fs (%zu bytes).", pts, ret);
+ al_log_debug("audio_buffer", "Delaying audio by %fs (%zu bytes).", pts, ret);
al_memset(data, 0, ret);
data += ret;
+ req -= ret;
// We can continue to delay.
- if ((req -= ret) == 0) return signal;
+ if (req == 0) return signal;
}
}
if (buf->pause != PAUSE_PLAYING) buf->pause = PAUSE_PLAYING;
#ifdef CAMU_AUDIO_BUFFER_FADE
}
#endif
- u8 flow = al_atomic_u8_load(&buf->flow, AL_ATOMIC_RELAXED);
- if (size < req) { // We don't have enough data to fulfill our request.
- if (flow == FLUSHED) {
- // 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);
+ if (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_size(&buf->peak);
+ if (ret > 0) {
+ // This doesn't break single reader, single writer because
+ // pushing data after a flush is not allowed.
+ 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.
+ if (have < req) {
+ buf->callback(buf->userdata, CAMU_BUFFER_EOF);
+ flow = SIGNALED;
+ al_atomic_store(u8)(&buf->flow, SIGNALED, AL_ATOMIC_RELEASE);
+ signal = have;
+ al_log_debug("audio_buffer", "Flushed (signal: %zu).", signal);
+ }
} else {
- // Put silence into the remainder of the request buffer.
- al_memset(data + size, 0, req - size);
+ // Silence remainder of the request.
+ al_memset(data + have, 0, req - have);
if (flow == FLOWING) {
- // Our stream hit an underrun because data wasn't coming in fast
- // enough. It's also possible to underrun if this function takes too
- // long which isn't checked here. That is unlikely though.
- al_log_warn("audio_buffer", "Underrun (req: %zu, have: %zu).", req, size);
+ // We hit an underrun because data wasn't coming in fast enough.
+ // An underrun can also happen in the audio output if read() (this function)
+ // takes too long. That isn't checked here.
+ al_log_warn("audio_buffer", "Underrun (req: %zu, have: %zu).", req, have);
+ // If we have data by the next read(), try to skip ahead to maintain sync.
+ // This might exacerbate the underrun issue but an underrun is already
+ // unexpected behavior, trying to stay in sync comes first.
+ buf->pause = PAUSE_PAUSED;
}
- // If flow = SIGNALED, we are just waiting to be removed or reset.
+ // If flow = SIGNALED, we are waiting to be removed or reset.
}
- req = size;
+ // The amount of available data could've increased in the flow = FLUSHED case.
+ req = AL_MIN(req, have);
}
if (req > 0) {
#ifdef CAMU_AUDIO_BUFFER_FADE
@@ -378,29 +344,38 @@ size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t re
buf->fade_offset += ret;
} else {
#endif
- // Read as much as we can from the ring buffer.
+ // Read as much as we determined we can.
ret = al_ring_buffer_read(&buf->rb, data, req);
- 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.
- if (ret && ((buf->size - size) - req) >= ret) {
- buf->callback(buf->userdata, CAMU_BUFFER_UNCORK);
- }
- }
+ al_assert(ret == req);
+ buf->pts += camu_audio_format_bytes_to_sec(&buf->fmt.req, ret);
#ifdef CAMU_AUDIO_BUFFER_FADE
}
- if (buf->fade_period > 0) handle_fade(buf, data, req, buf->pause == PAUSE_FADING);
+ f32 step = 0.f;
+ if (buf->pause == PAUSE_FADING) {
+ step = -FADE_STEP(&buf->fmt.req);
+ } else if (buf->volume < 1.f) {
+ step = FADE_STEP(&buf->fmt.req);
+ }
+ if (step != 0.f || buf->volume != 1.f) {
+ buf->volume = apply_volume(data, req, &buf->fmt.req, buf->volume, step);
+ }
#endif
}
- // Never return less than req, unless we hit the end of the stream.
+ // Check if we should request to uncork.
+ if (flow == FLOWING) {
+ ret = al_atomic_load(size_t)(&buf->uncork_at, AL_ATOMIC_RELAXED);
+ if (ret && have - req <= ret) {
+ buf->callback(buf->userdata, CAMU_BUFFER_UNCORK);
+ }
+ }
+ // To signal EOF, return less then req.
return signal;
}
void camu_audio_buffer_free(struct camu_audio_buffer *buf)
{
if (buf->fmt.resampler_needed) {
- camu_ff_resampler_close(&buf->resamp);
+ buf->resamp->free(&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 c4ba142..3adccc7 100644
--- a/src/buffer/audio.h
+++ b/src/buffer/audio.h
@@ -4,45 +4,44 @@
#include <al/ring_buffer.h>
#include "../codec/codec.h"
-#include "../codec/ffmpeg/resampler.h"
#include "../mixer/mixer.h"
#include "clock.h"
#include "peak_buffer.h"
-//#define CAMU_AUDIO_BUFFER_FADE
+#define CAMU_AUDIO_BUFFER_FADE
struct camu_audio_buffer {
- struct camu_stream *stream;
+ struct camu_codec_stream *stream;
struct camu_mixer *mixer;
f64 pts;
u8 pause;
- bool ignore_desync;
struct camu_clock *clock;
+ bool ignore_desync;
+ f64 latency;
- f64 volume;
- s32 fade_period;
- size_t fade_offset;
- s32 fade_rate;
-
- s32 bytes_per_sample;
- struct camu_ff_resampler resamp;
- struct camu_ff_resample_fmt fmt;
+ struct camu_resampler_format fmt;
+ struct camu_resampler *resamp;
u8 *data;
size_t size;
struct al_ring_buffer rb;
- struct { size_t low, high; } watermark;
+ struct { size_t min, buffered; } mark;
bool buffered;
+ atomic(u8) flow;
+
struct camu_peak_buffer peak;
- atomic_size_t continue_mark;
+ atomic(size_t) uncork_at;
- atomic_u8 flow;
+ f32 volume;
+#ifdef CAMU_AUDIO_BUFFER_FADE
+ size_t fade_offset;
+#endif
#ifdef CAMU_MIXER_THREADED
- atomic_bool ref;
+ atomic(u8) ref;
#endif
void (*callback)(void *, u8);
@@ -50,8 +49,9 @@ 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);
+bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_codec_stream *stream);
+void camu_audio_buffer_set_latency(struct camu_audio_buffer *buf, f64 latency);
+void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_codec_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);
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
index d26e4a8..c4a07a6 100644
--- a/src/buffer/clock.c
+++ b/src/buffer/clock.c
@@ -1,116 +1,117 @@
#include <al/lib.h>
+#include <aki/thread.h>
#include "clock.h"
-void camu_clock_set(struct camu_clock *clock, struct shrb_seek_req *req)
-{
- al_atomic_f64_store(&clock->tick, 0.0, AL_ATOMIC_RELAXED);
- if (req->paused) {
- clock->pause_at = req->paused_at / 1000000.0;
- al_atomic_f64_store(&clock->base, clock->pause_at, AL_ATOMIC_RELAXED);
- al_atomic_u64_store(&clock->start, 0Lu, AL_ATOMIC_RELAXED);
- } else {
- clock->pause_at = -1.0;
- al_atomic_f64_store(&clock->base, req->base / 1000000.0, AL_ATOMIC_RELAXED);
- al_atomic_u64_store(&clock->start, req->ts + req->delay, AL_ATOMIC_RELAXED);
- clock->ts = req->ts;
- }
- clock->offset = 0L;
- clock->delay = req->delay;
-}
+#define RUNNING 0.0
+#define PAUSED -DBL_MAX
+#define ENDED DBL_MAX
-f64 calc_tick_internal(struct camu_clock *clock)
+void camu_clock_init(struct camu_clock *clock, void (*callback)(void *, u8), void *userdata)
{
- s64 diff = clock->delay == 0Lu ? 0L
- : ((al_atomic_u64_load(&clock->start, AL_ATOMIC_ACQUIRE) + clock->offset)
- - aki_get_timestamp());
- f64 tick = aki_get_tick() - al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
- //al_printf("DIFFERENCE: %f\n", diff / 1000000.0);
- tick += diff / 1000000.0;
- al_atomic_f64_store(&clock->tick, tick, AL_ATOMIC_RELAXED);
- al_atomic_u64_store(&clock->start, 0Lu, AL_ATOMIC_RELEASE);
- return tick;
+ clock->callback = callback;
+ clock->userdata = userdata;
}
-bool camu_clock_is_late(struct camu_clock *clock)
+static f64 calc_tick_offset(f64 tick, u64 now, u64 target)
{
- s64 diff = (s64)(al_atomic_u64_load(&clock->start, AL_ATOMIC_RELAXED)) - aki_get_timestamp();
- return clock->delay > 0L && diff < 0L;
+ if (now > target) return tick - (now - target) / 1000000.0;
+ else return tick + (target - now) / 1000000.0;
}
-void camu_clock_ready(struct camu_clock *clock)
+void camu_clock_set(struct camu_clock *clock, f64 base)
{
- if (clock->pause_at != -1.0) {
- clock->pause_at = -1.0;
- return;
- }
- camu_clock_resume(clock);
+ clock->base = base;
+ clock->offset = 0.0;
+ al_atomic_store(f64)(&clock->tick, -1.0, AL_ATOMIC_RELAXED);
+ al_atomic_store(f64)(&clock->pause, PAUSED, AL_ATOMIC_RELAXED);
+ clock->paused_at = 0.0;
}
-void camu_clock_offset(struct camu_clock *clock, s64 offset)
+void camu_clock_offset(struct camu_clock *clock, f64 amount)
{
- clock->offset += offset;
- //al_printf("OFFSET: %f\n", clock->offset / 1000000.0);
+ clock->offset += amount;
}
-void camu_clock_resume(struct camu_clock *clock)
+void camu_clock_pause(struct camu_clock *clock, u64 target)
{
- 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 != 0Lu) {
- // No longer report paused, calculate tick on first call to get_pts().
- al_atomic_f64_store(&clock->tick, -1.0, AL_ATOMIC_RELEASE);
+ f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE);
+ if (pause == ENDED) return;
+ al_assert(clock->paused_at == -1.0);
+ f64 tick = aki_get_tick();
+ if (target > 0) {
+ tick = calc_tick_offset(tick, aki_get_timestamp(), target);
+ al_atomic_store(f64)(&clock->pause, tick, AL_ATOMIC_RELEASE);
} else {
- tick = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
- al_atomic_f64_store(&clock->tick, aki_get_tick() - tick, AL_ATOMIC_RELEASE);
+ al_atomic_store(f64)(&clock->pause, PAUSED, AL_ATOMIC_RELEASE);
}
- clock->pause_at = -1.0;
+ clock->paused_at = tick;
}
-void camu_clock_arm_resume(struct camu_clock *clock, u64 ts)
+void camu_clock_resume(struct camu_clock *clock, u64 target)
{
- al_atomic_u64_store(&clock->start, ts, AL_ATOMIC_RELAXED);
- al_atomic_f64_store(&clock->tick, -1.0, AL_ATOMIC_RELAXED);
- clock->pause_at = -1.0;
+ f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE);
+ if (pause == ENDED) return;
+ al_assert(clock->paused_at != -1.0);
+ f64 tick = aki_get_tick();
+ //if (pause != PAUSED && tick < pause) {
+ // tick = pause;
+ //} else if (target > 0) {
+ if (target > 0) {
+ tick = calc_tick_offset(tick, aki_get_timestamp(), target);
+ al_printf("%f\n", tick);
+ }
+ if (clock->paused_at == 0.0) {
+ if (target == 0) {
+ // target = 0 can never be synced.
+ al_atomic_store(f64)(&clock->tick, -1.0, AL_ATOMIC_RELAXED);
+ } else {
+ al_atomic_store(f64)(&clock->tick, tick, AL_ATOMIC_RELAXED);
+ }
+ } else {
+ clock->offset += tick - clock->paused_at;
+ }
+ al_atomic_store(f64)(&clock->pause, RUNNING, AL_ATOMIC_RELEASE);
+ clock->paused_at = -1.0;
}
-void camu_clock_pause(struct camu_clock *clock, f64 base)
+bool camu_clock_is_paused(struct camu_clock *clock)
{
- f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_ACQUIRE);
- if (tick == 0.0) return;
- if (base == -1.0) base = aki_get_tick() - tick;
- al_atomic_f64_store(&clock->tick, 0.0, AL_ATOMIC_RELEASE);
- clock->pause_at = base;
- al_atomic_f64_store(&clock->base, clock->pause_at, AL_ATOMIC_RELAXED);
+ return al_atomic_load(f64)(&clock->pause, AL_ATOMIC_RELAXED) == PAUSED;
}
-void camu_clock_arm_pause(struct camu_clock *clock, f64 pts)
+void camu_clock_end(struct camu_clock *clock)
{
- clock->pause_at = pts;
+ al_atomic_store(f64)(&clock->pause, ENDED, AL_ATOMIC_RELAXED);
}
-bool camu_clock_is_paused(struct camu_clock *clock)
+bool camu_clock_is_ended(struct camu_clock *clock)
{
- f64 tick = al_atomic_f64_load(&clock->tick, AL_ATOMIC_RELAXED);
- if (tick == 0.0) return true;
- 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);
- return true;
- }
- return false;
+ return al_atomic_load(f64)(&clock->pause, AL_ATOMIC_RELAXED) == ENDED;
}
f64 camu_clock_get_base_pts(struct camu_clock *clock)
{
- return al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
+ return clock->base;
}
f64 camu_clock_get_pts(struct camu_clock *clock, f64 offset)
{
- 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;
+ f64 pause = al_atomic_load(f64)(&clock->pause, AL_ATOMIC_ACQUIRE);
+ if (pause == PAUSED || pause == ENDED) return -1.0;
+ f64 current = aki_get_tick();
+ f64 tick = al_atomic_load(f64)(&clock->tick, AL_ATOMIC_RELAXED);
+ if (tick == -1.0) {
+ tick = al_atomic_compare_and_swap(f64)(&clock->tick, -1.0, current);
+ if (tick == -1.0) tick = current;
+ }
+ // pause > 0.0 means running or pause armed.
+ if (pause > 0.0 && current > pause) {
+ current = pause;
+ pause = al_atomic_compare_and_swap(f64)(&clock->pause, pause, PAUSED);
+ if (pause != PAUSED) {
+ clock->callback(clock->userdata, CAMU_CLOCK_PAUSED);
+ }
+ }
+ return clock->base + (current - (tick + clock->offset)) + offset;
}
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
index 6f505be..17edb0d 100644
--- a/src/buffer/clock.h
+++ b/src/buffer/clock.h
@@ -1,27 +1,51 @@
#pragma once
-#include <aki/thread.h>
+#include <al/types.h>
#include <al/atomic.h>
-#include "../shrub/handler.h"
+// Synced:
+// () = other clients only
+// Pause:
+// - User input -> start fade and request pause from server
+// - Get response -> (start fade) arm pause or reseek if late
+// Resume:
+// - User input -> request resume from server
+// - Get response -> arm resume, apply offset to queued
+// Next/Prev:
+// - User input -> start fade request next from server
+// - Get response -> (start fade) arm pause, arm resume on next, then swap on pause
+//
+// Local:
+// Pause:
+// - User input -> start fade and pause immediately
+// Reumse:
+// - User input -> resume and fade in
+// Next/Prev:
+// - User input -> swap
+
+enum {
+ CAMU_CLOCK_PAUSED = 0
+};
struct camu_clock {
- atomic_f64 base, tick;
- atomic_u64 start;
- u64 ts;
- s64 offset;
- u64 delay;
- f64 pause_at;
+ f64 base;
+ f64 offset;
+ atomic(f64) tick;
+ atomic(f64) pause;
+ f64 paused_at;
+ void (*callback)(void *, u8);
+ void *userdata;
};
-void camu_clock_set(struct camu_clock *clock, struct shrb_seek_req *req);
-bool camu_clock_is_late(struct camu_clock *clock);
-void camu_clock_ready(struct camu_clock *clock);
-void camu_clock_offset(struct camu_clock *clock, s64 offset);
-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, f64 base);
-void camu_clock_arm_pause(struct camu_clock *clock, f64 pts);
+void camu_clock_init(struct camu_clock *clock, void (*callback)(void *, u8), void *userdata);
+void camu_clock_set(struct camu_clock *clock, f64 base);
+void camu_clock_offset(struct camu_clock *clock, f64 amount);
+
+void camu_clock_pause(struct camu_clock *clock, u64 ts);
+void camu_clock_resume(struct camu_clock *clock, u64 ts);
bool camu_clock_is_paused(struct camu_clock *clock);
+void camu_clock_end(struct camu_clock *clock);
+bool camu_clock_is_ended(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/frame_queue.h b/src/buffer/frame_queue.h
index 512b7f2..0175668 100644
--- a/src/buffer/frame_queue.h
+++ b/src/buffer/frame_queue.h
@@ -14,7 +14,7 @@ enum {
};
struct camu_frame_queue {
- void (*push)(struct camu_frame_queue *, struct camu_frame *, f64);
+ void (*push)(struct camu_frame_queue *, struct camu_codec_frame *, f64);
#ifdef CAMU_HAVE_FFMPEG
void (*push_av_frame)(struct camu_frame_queue *, AVFrame *, f64);
#endif
diff --git a/src/buffer/peak_buffer.c b/src/buffer/peak_buffer.c
index 2856a78..005e77f 100644
--- a/src/buffer/peak_buffer.c
+++ b/src/buffer/peak_buffer.c
@@ -2,10 +2,10 @@
// https://github.com/MusicPlayerDaemon/MPD/blob/c71e586c530e1d066efde7c2ca40c363f74100c7/src/util/PeakBuffer.cxx
-void camu_peak_buffer_init(struct camu_peak_buffer *buf)
+void camu_peak_buffer_init(struct camu_peak_buffer *buf, size_t size)
{
aki_buffer_init(&buf->buf);
- aki_buffer_ensure_space(&buf->buf, 128 * 1024);
+ aki_buffer_ensure_space(&buf->buf, size);
}
void camu_peak_buffer_push(struct camu_peak_buffer *buf, u8 *data, size_t size)
@@ -13,9 +13,13 @@ void camu_peak_buffer_push(struct camu_peak_buffer *buf, u8 *data, size_t size)
aki_buffer_append(&buf->buf, data, size);
}
-u8 *camu_peak_buffer_flush(struct camu_peak_buffer *buf, size_t *size)
+size_t camu_peak_buffer_size(struct camu_peak_buffer *buf)
+{
+ return aki_buffer_get_size(&buf->buf);
+}
+
+u8 *camu_peak_buffer_flush(struct camu_peak_buffer *buf)
{
- *size = aki_buffer_get_size(&buf->buf);
aki_buffer_set_size(&buf->buf, 0);
return aki_buffer_get_ptr(&buf->buf, 0);
}
diff --git a/src/buffer/peak_buffer.h b/src/buffer/peak_buffer.h
index b8d9934..c518ade 100644
--- a/src/buffer/peak_buffer.h
+++ b/src/buffer/peak_buffer.h
@@ -7,7 +7,8 @@ struct camu_peak_buffer {
struct aki_buffer buf;
};
-void camu_peak_buffer_init(struct camu_peak_buffer *buf);
+void camu_peak_buffer_init(struct camu_peak_buffer *buf, size_t size);
void camu_peak_buffer_push(struct camu_peak_buffer *buf, u8 *data, size_t size);
-u8 *camu_peak_buffer_flush(struct camu_peak_buffer *buf, size_t *size);
+size_t camu_peak_buffer_size(struct camu_peak_buffer *buf);
+u8 *camu_peak_buffer_flush(struct camu_peak_buffer *buf);
void camu_peak_buffer_free(struct camu_peak_buffer *buf);
diff --git a/src/buffer/video.c b/src/buffer/video.c
index 3b9be20..b2381b2 100644
--- a/src/buffer/video.c
+++ b/src/buffer/video.c
@@ -4,35 +4,39 @@
#include "common.h"
#include "common_internal.h"
-//#define CAMU_VIDEO_BUFFER_FORCE_SCALER
-
#define BUFFER_WATERMARK_LOW 8 // frames.
-#define BUFFER_WATERMARK_BUFFERED 8
-#define BUFFER_WATERMARK_HIGH 14
+#define BUFFER_WATERMARK_BUFFERED 5
+#define BUFFER_WATERMARK_HIGH 12
#define BUFFER_WATERMARK_RESET BUFFER_WATERMARK_HIGH + 10.
-bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock, f64 offset,
- struct camu_renderer *renderer)
+bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock, struct camu_renderer *renderer)
{
buf->clock = clock;
- buf->offset = offset;
+ buf->latency = 0.0;
+ buf->pts = -1.0;
+ // Defaulting single_frame to true should simplify non-configured buffers in sink.
+ buf->single_frame = true;
buf->queue = renderer->create_queue(renderer);
buf->buffered = false;
- al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
+ al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
#ifdef CAMU_SCREEN_THREADED
- al_atomic_bool_store(&buf->ref, false, AL_ATOMIC_RELAXED);
+ al_atomic_store(u8)(&buf->ref, 0, AL_ATOMIC_RELAXED);
#endif
buf->view.mode = CAMU_VIEW_NONE;
return true;
}
-bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stream *stream)
+bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_codec_stream *stream)
{
buf->stream = stream;
- switch (stream->type) {
+ switch (stream->mode) {
case CAMU_NORMAL: {
+ s32 width = buf->stream->video.width;
+ s32 height = buf->stream->video.height;
buf->single_frame = true;
buf->avg_frame_duration = 0.0;
+ const char *format_name = camu_pixel_format_name(buf->stream->video.format);
+ al_log_info("video_buffer", "Stream: %s (%dx%d) %s.", format_name, width, height, "IMAGE");
break;
}
#ifdef CAMU_HAVE_FFMPEG
@@ -44,8 +48,8 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stre
buf->fmt.in_format = stream->av.stream->codecpar->format;
buf->single_frame = stream->av.stream->duration == 0 ||
stream->av.stream->avg_frame_rate.den == 0;
- AVRational frame_rate = buf->single_frame ? (AVRational){ 0, 1 } :
- stream->av.stream->avg_frame_rate;
+ AVRational frame_rate = buf->single_frame ? ((AVRational){ 0, 1 })
+ : stream->av.stream->avg_frame_rate;
buf->avg_frame_duration = av_q2d(av_inv_q(frame_rate));
const char *format_name = av_get_pix_fmt_name(buf->fmt.in_format);
al_log_info("video_buffer", "Stream: %s (%dx%d) %s %.2ffps.", format_name,
@@ -56,7 +60,7 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stre
buf->fmt.req_format = AV_PIX_FMT_RGBA;
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,
+ al_log_info("video_buffer", "Scaling to: %s (%dx%d)", format_name,
buf->fmt.req_width, buf->fmt.req_height);
} else {
return false;
@@ -69,6 +73,11 @@ bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stre
return true;
}
+void camu_video_buffer_set_latency(struct camu_video_buffer *buf, f64 latency)
+{
+ buf->latency = latency;
+}
+
bool camu_video_buffer_is_single_frame(struct camu_video_buffer *buf)
{
return buf->single_frame;
@@ -78,12 +87,12 @@ 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)) {
+ al_log_debug("video_buffer", "Buffered (watermark: %d frames).", count);
// 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 (buf->single_frame) al_atomic_store(u8)(&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);
@@ -95,12 +104,13 @@ static void push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame
f64 pts = 0.0;
if (!buf->single_frame) {
AVStream *stream = buf->stream->av.stream;
- pts = (frame->best_effort_timestamp - stream->start_time) * av_q2d(stream->time_base);
+ pts = frame->best_effort_timestamp * av_q2d(stream->time_base);
if (pts + buf->avg_frame_duration < camu_clock_get_base_pts(buf->clock)) {
av_frame_free(&frame);
return;
}
}
+ if (buf->pts == -1.0) buf->pts = pts;
#ifdef CAMU_VIDEO_BUFFER_FORCE_SCALER
if (buf->fmt.scaler_needed) {
if (!camu_ff_scaler_scale(&buf->scale, (const u8 **)frame->data, frame->linesize)) {
@@ -115,14 +125,14 @@ static void push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame
}
#endif
-void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_frame *frame)
+void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_codec_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);
+ camu_codec_frame_discard(frame);
return;
}
- switch (frame->type) {
+ switch (frame->mode) {
case CAMU_NORMAL:
buf->queue->push(buf->queue, frame, 0.0);
break;
@@ -141,7 +151,7 @@ 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);
+ al_atomic_store(u8)(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
}
// flush() always comes from the same thread as push().
@@ -149,26 +159,28 @@ void camu_video_buffer_flush(struct camu_video_buffer *buf)
{
buf->queue->flush(buf->queue);
if (!buf->buffered) {
+ al_log_debug("video_buffer", "Buffered (watermark: %d frames).", buf->queue->count(buf->queue));
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);
+ al_atomic_store(u8)(&buf->flow, FLUSHED, AL_ATOMIC_RELAXED);
}
bool camu_video_buffer_read(struct camu_video_buffer *buf, void *out)
{
- 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;
+ if (!buf->single_frame) {
+ f64 pts = camu_clock_get_pts(buf->clock, buf->latency);
+ if (pts > buf->pts) buf->pts = pts;
+ }
+ u8 flow = al_atomic_load(u8)(&buf->flow, AL_ATOMIC_RELAXED);
+ if (flow == SIGNALED && !buf->single_frame) {
+ buf->callback(buf->userdata, CAMU_BUFFER_EOF);
+ return false;
}
- u8 flow = al_atomic_u8_load(&buf->flow, AL_ATOMIC_RELAXED);
- u8 ret = buf->queue->read(buf->queue, base, out);
+ u8 ret = buf->queue->read(buf->queue, buf->pts, out);
if (flow == FLUSHED && (ret == CAMU_QUEUE_EOF || (buf->single_frame && ret == CAMU_QUEUE_OK))) {
buf->callback(buf->userdata, CAMU_BUFFER_EOF);
- al_atomic_u8_store(&buf->flow, SIGNALED, AL_ATOMIC_RELAXED);
+ al_atomic_store(u8)(&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);
diff --git a/src/buffer/video.h b/src/buffer/video.h
index 5e18b33..0e8bfff 100644
--- a/src/buffer/video.h
+++ b/src/buffer/video.h
@@ -1,36 +1,43 @@
#pragma once
+//#define CAMU_VIDEO_BUFFER_FORCE_SCALER
+
#include <al/types.h>
#include "../codec/codec.h"
-#include "../codec/ffmpeg/scaler.h"
#include "../render/renderer.h"
#include "../screen/screen.h"
#include "../screen/view.h"
+#ifdef CAMU_VIDEO_BUFFER_FORCE_SCALER
+#endif
+#include "../codec/ffmpeg/scaler.h"
#include "clock.h"
#include "frame_queue.h"
struct camu_video_buffer {
- struct camu_stream *stream;
+ struct camu_codec_stream *stream;
struct camu_clock *clock;
- f32 start_time;
- f32 avg_frame_duration;
- f64 offset;
+ f64 latency;
+ f64 pts;
bool single_frame;
+ f32 avg_frame_duration;
+#ifdef CAMU_VIDEO_BUFFER_FORCE_SCALER
struct camu_ff_scaler scale;
+#endif
+ // TODO: Terrible
struct camu_ff_scale_fmt fmt;
struct camu_frame_queue *queue;
bool buffered;
- atomic_u8 flow;
+ atomic(u8) flow;
#ifdef CAMU_SCREEN_THREADED
- atomic_bool ref;
+ atomic(u8) ref;
#endif
struct camu_view view;
@@ -38,11 +45,11 @@ struct camu_video_buffer {
void *userdata;
};
-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_init(struct camu_video_buffer *buf, struct camu_clock *clock, struct camu_renderer *renderer);
+bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_codec_stream *stream);
+void camu_video_buffer_set_latency(struct camu_video_buffer *buf, f64 latency);
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_push(struct camu_video_buffer *buf, struct camu_codec_frame *frame);
void camu_video_buffer_reset(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);
diff --git a/src/buffer/volume.h b/src/buffer/volume.h
new file mode 100644
index 0000000..f785eb9
--- /dev/null
+++ b/src/buffer/volume.h
@@ -0,0 +1,36 @@
+#include "../codec/codec.h"
+
+AL_UNUSED_FUNCTION_PUSH
+
+static inline f32 apply_audio_f32(f32 *data, size_t sample_count, s32 channel_count, f32 volume, f32 step)
+{
+ for (u32 i = 0; i < sample_count; i++) {
+ if ((i + 1) % channel_count == 0) {
+ volume = AL_CLAMP(volume + step, 0.f, 1.f);
+ }
+ data[i] *= volume;
+ }
+ return volume;
+}
+
+static inline f32 apply_volume(u8 *data, size_t size, struct camu_audio_format *fmt, f32 volume, f32 step)
+{
+ size_t sample_count;
+ switch (camu_audio_format_bytes_per_sample(fmt)) {
+ case 4:
+ sample_count = size / 4;
+ if (fmt->format == CAMU_SAMPLE_FORMAT_S32) {
+ } else if (fmt->format == CAMU_SAMPLE_FORMAT_FLT) {
+ return apply_audio_f32((f32 *)data, sample_count, fmt->channel_count, volume, step);
+ }
+ break;
+ case 2:
+ sample_count = size / 2;
+ if (fmt->format == CAMU_SAMPLE_FORMAT_S16) {
+ }
+ break;
+ }
+ return volume;
+}
+
+AL_UNUSED_FUNCTION_POP