summaryrefslogtreecommitdiff
path: root/src/buffer
diff options
context:
space:
mode:
Diffstat (limited to 'src/buffer')
-rw-r--r--src/buffer/audio.c361
-rw-r--r--src/buffer/audio.h59
-rw-r--r--src/buffer/clock.c53
-rw-r--r--src/buffer/clock.h16
-rw-r--r--src/buffer/common.h9
-rw-r--r--src/buffer/frame_queue.h28
-rw-r--r--src/buffer/meson.build8
-rw-r--r--src/buffer/peak_buffer.c24
-rw-r--r--src/buffer/peak_buffer.h13
-rw-r--r--src/buffer/video.c148
-rw-r--r--src/buffer/video.h42
11 files changed, 761 insertions, 0 deletions
diff --git a/src/buffer/audio.c b/src/buffer/audio.c
new file mode 100644
index 0000000..34d3c65
--- /dev/null
+++ b/src/buffer/audio.c
@@ -0,0 +1,361 @@
+#include <al/log.h>
+
+#include "audio.h"
+
+#define BUFFER_USEC (10 * 1000000L)
+#define BUFFER_WATERMARK_LOW (3 * 1000000L) // Must be a most half of the buffer size.
+#define BUFFER_WATERMARK_HIGH (4 * 1000000L)
+
+#define DESYNC_PTS 0.022
+
+#ifdef CAMU_AUDIO_BUFFER_FADE
+#define FADE 0.9999
+#define FADE_LENGTH 6
+#define FADE_STEP(rate) ((FADE / (rate)) / (FADE_LENGTH - 1))
+#endif
+
+enum {
+ FLOWING = 0,
+ FLUSHED,
+ SIGNALED
+};
+
+enum {
+ PAUSE_PRE = 0,
+ PAUSE_UNPAUSED,
+ PAUSE_IGNORE_DESYNC,
+#ifdef CAMU_AUDIO_BUFFER_FADE
+ PAUSE_FADING,
+#endif
+ PAUSE_PLAYING
+};
+
+static bool setup_optimal_resampler(struct camu_audio_buffer *buf, struct camu_lav_resampler *resamp,
+ struct camu_lav_resample_fmt *fmt)
+{
+ /*
+ 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_sample_rate = fmt->in_sample_rate;
+ */
+
+ buf->fade_rate = buf->fmt.req_sample_rate / 10;
+ buf->bytes_per_sample = (s32)camu_lav_resample_fmt_bytes_per_sample(fmt);
+
+ return camu_lav_resampler_init(resamp, fmt);
+}
+
+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->pts = camu_mixer_get_latency(buf->mixer);
+ buf->buffered = false;
+ buf->pause = PAUSE_PRE;
+ //buf->pause = PAUSE_IGNORE_DESYNC;
+#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
+ return true;
+}
+
+bool camu_audio_buffer_configure(struct camu_audio_buffer *buf, struct camu_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_mixer_pick_format(buf->mixer, &buf->fmt);
+ av_channel_layout_default(&buf->fmt.req_channel_layout, buf->fmt.req_channel_count);
+ if (!setup_optimal_resampler(buf, &buf->resamp, &buf->fmt)) {
+ return false;
+ }
+ 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);
+ 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);
+ buf->size = camu_lav_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);
+ camu_peak_buffer_init(&buf->peak);
+ buf->stream = stream;
+ return true;
+}
+
+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);
+ }
+
+ f64 duration = camu_lav_resample_fmt_samples_to_sec(&buf->fmt, sample_count);
+ if (pts + duration < camu_clock_get_base_pts(buf->clock)) {
+ return false;
+ }
+
+ 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 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;
+ }
+ 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. We could attempt to recover by
+ // re-seeking the whole stream.
+ camu_peak_buffer_flush(&buf->peak, &size);
+ }
+ // 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.
+ al_atomic_size_t_store(&buf->continue_mark, buf->watermark.low + peak, AL_ATOMIC_RELAXED);
+ buf->callback(buf->userdata, CAMU_BUFFER_STOP);
+ 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, data[0], have);
+
+ return true;
+}
+
+#ifdef HAVE_FFMPEG
+static void push_av_frame_internal(struct camu_audio_buffer *buf, AVFrame *frame)
+{
+ if (al_atomic_u8_load(&buf->flow, AL_ATOMIC_RELAXED) == FLOWING) {
+ u8 **data = frame->data;
+ s32 sample_count = frame->nb_samples;
+ AVStream *stream = buf->stream->av.stream;
+ f64 pts = (frame->best_effort_timestamp - stream->start_time) * av_q2d(stream->time_base);
+ push_internal(buf, data, sample_count, pts);
+ }
+ av_frame_free(&frame);
+}
+#endif
+
+void camu_audio_buffer_push(struct camu_audio_buffer *buf, struct camu_frame *frame)
+{
+ switch (frame->type) {
+ case CAMU_NORMAL:
+ break;
+#ifdef HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT: {
+ push_av_frame_internal(buf, frame->av.frame);
+ break;
+ }
+#endif
+ }
+ al_free(frame);
+}
+
+void camu_audio_buffer_reset(struct camu_audio_buffer *buf)
+{
+ al_atomic_size_t_store(&buf->continue_mark, 0, AL_ATOMIC_RELAXED);
+ al_atomic_u8_store(&buf->flow, FLOWING, AL_ATOMIC_RELAXED);
+ buf->pts = camu_clock_get_base_pts(buf->clock);
+ buf->buffered = false;
+ al_ring_buffer_reset(&buf->rb);
+}
+
+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;
+ }
+}
+
+#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)
+
+static void handle_fade(struct camu_audio_buffer *buf, u8 *data, size_t size, bool out)
+{
+ u8 bits[MAX_BYTES];
+ u8 *ptr = bits + (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);
+ }
+ data += buf->bytes_per_sample;
+ }
+}
+#endif
+
+size_t camu_audio_buffer_read(struct camu_audio_buffer *buf, u8 *data, size_t req)
+{
+ f64 pts = camu_clock_get_pts(buf->clock);
+ if (pts == -1.0) {
+#ifdef CAMU_AUDIO_BUFFER_FADE
+ if (buf->pause == PAUSE_PLAYING) {
+ buf->volume = FADE;
+ buf->fade_period = FADE_LENGTH;
+ buf->fade_offset = 0;
+ buf->pause = PAUSE_FADING;
+ } else if (buf->fade_period == 0) {
+ al_memset(data, 0, req);
+ if (buf->pause != PAUSE_UNPAUSED && buf->pause != PAUSE_PRE) {
+ buf->pause = PAUSE_UNPAUSED;
+ buf->callback(buf->userdata, CAMU_BUFFER_PAUSED);
+ }
+ return req;
+ }
+ } else if (buf->pause == PAUSE_FADING || buf->pause == PAUSE_UNPAUSED) {
+ buf->volume = 1.0 - FADE;
+ buf->fade_period = FADE_LENGTH;
+ buf->fade_offset = 0;
+ buf->pause = PAUSE_PLAYING;
+#else
+ al_memset(data, 0, req);
+ if (buf->pause != PAUSE_UNPAUSED && buf->pause != PAUSE_PRE) {
+ buf->pause = PAUSE_UNPAUSED;
+ buf->callback(buf->userdata, CAMU_BUFFER_PAUSED);
+ }
+ return req;
+#endif
+ }
+ size_t ret, signal = req;
+ size_t size = al_ring_buffer_occupied(&buf->rb);
+ if (size < buf->fade_offset) size = 0;
+ else size -= buf->fade_offset;
+#ifdef CAMU_AUDIO_BUFFER_FADE
+ if (buf->pause != PAUSE_FADING) {
+#endif
+ pts -= buf->pts;
+ if (UNLIKELY(buf->pause == PAUSE_PRE || buf->pause == PAUSE_UNPAUSED)) {
+ if (pts > 0.0) {
+ ret = camu_lav_resample_fmt_sec_to_bytes(&buf->fmt, pts);
+ ret = AL_MIN(ret, size);
+ al_log_debug("audio_buffer", "Skipping %.5fs of audio.", pts);
+ ret = al_ring_buffer_discard(&buf->rb, ret);
+ size -= ret;
+ buf->pts += camu_lav_resample_fmt_bytes_to_sec(&buf->fmt, ret);
+ } else if (pts < 0.0) {
+ pts = -pts;
+ ret = camu_lav_resample_fmt_sec_to_bytes(&buf->fmt, pts);
+ ret = AL_MIN(ret, req);
+ al_log_debug("audio_buffer", "Delaying audio by %.5fs.", pts);
+ al_memset(data, 0, ret);
+ data += ret;
+ req -= ret;
+ }
+ }
+ 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) {
+ if (flow == FLUSHED) {
+ signal = size;
+ al_atomic_u8_store(&buf->flow, SIGNALED, AL_ATOMIC_RELAXED);
+ al_log_debug("audio_buffer", "Flushed (signal: %li).", signal);
+ buf->callback(buf->userdata, CAMU_BUFFER_EOF);
+ } else {
+ al_memset(data + size, 0, req - size);
+ if (flow == FLOWING) {
+ al_log_warn("audio_buffer", "Underrun (req: %i, have: %i).", req, size);
+ }
+ }
+ req = size;
+ }
+ if (req > 0) {
+#ifdef CAMU_AUDIO_BUFFER_FADE
+ if (buf->pause == PAUSE_FADING) {
+ ret = al_ring_buffer_peek(&buf->rb, data, buf->fade_offset, req);
+ buf->fade_offset += ret;
+ } else {
+#endif
+ ret = al_ring_buffer_read(&buf->rb, data, req);
+ buf->pts += camu_lav_resample_fmt_bytes_to_sec(&buf->fmt, ret);
+ if (flow == FLOWING) {
+ ret = al_atomic_size_t_load(&buf->continue_mark, AL_ATOMIC_RELAXED);
+ if (ret && ((buf->size - size) - req) >= ret) {
+ buf->callback(buf->userdata, CAMU_BUFFER_CONTINUE);
+ }
+ }
+#ifdef CAMU_AUDIO_BUFFER_FADE
+ }
+ if (buf->fade_period > 0) {
+ buf->fade_period--;
+ handle_fade(buf, data, req, buf->pause == PAUSE_FADING);
+ }
+#endif
+ }
+ return signal;
+}
+
+void camu_audio_buffer_free(struct camu_audio_buffer *buf)
+{
+ if (buf->fmt.resampler_needed) {
+ camu_lav_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
new file mode 100644
index 0000000..1a51807
--- /dev/null
+++ b/src/buffer/audio.h
@@ -0,0 +1,59 @@
+#pragma once
+
+#include <al/types.h>
+#include <al/ring_buffer.h>
+
+#include "../codec/codec.h"
+#include "../codec/libav/resampler.h"
+#include "../mixer/mixer.h"
+#include "../mixer/audio.h"
+
+#include "clock.h"
+#include "common.h"
+#include "peak_buffer.h"
+
+#define CAMU_AUDIO_BUFFER_FADE
+
+struct camu_audio_buffer {
+ struct camu_stream *stream;
+ struct camu_mixer *mixer;
+
+ f64 pts;
+ u8 pause;
+ struct camu_clock *clock;
+
+ f64 volume;
+ s32 fade_period;
+ size_t fade_offset;
+ s32 fade_rate;
+
+ s32 bytes_per_sample;
+ struct camu_lav_resampler resamp;
+ struct camu_lav_resample_fmt fmt;
+
+ atomic_u8 flow;
+
+ u8 *data;
+ size_t size;
+ struct al_ring_buffer rb;
+ struct { size_t low, high; } watermark;
+ bool buffered;
+
+ struct camu_peak_buffer peak;
+ atomic_size_t continue_mark;
+
+#ifdef CAMU_MIXER_THREADED
+ atomic_bool ref;
+#endif
+
+ void (*callback)(void *, u8);
+ void *userdata;
+};
+
+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_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);
+void camu_audio_buffer_free(struct camu_audio_buffer *buf);
diff --git a/src/buffer/clock.c b/src/buffer/clock.c
new file mode 100644
index 0000000..5f4133e
--- /dev/null
+++ b/src/buffer/clock.c
@@ -0,0 +1,53 @@
+#include "clock.h"
+
+void camu_clock_init(struct camu_clock *clock)
+{
+ f64 tick = aki_get_tick();
+ al_atomic_f64_store(&clock->pause, tick, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->base, 0.0, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->start, 0.0, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->offset, 0.0, AL_ATOMIC_RELAXED);
+}
+
+void camu_clock_pause(struct camu_clock *clock)
+{
+ al_atomic_f64_store(&clock->pause, aki_get_tick(), AL_ATOMIC_RELAXED);
+}
+
+void camu_clock_resume(struct camu_clock *clock)
+{
+ f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED);
+ f64 start = al_atomic_f64_load(&clock->start, AL_ATOMIC_RELAXED);
+ f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
+ if (start == 0.0) start = (pause - base);
+ al_atomic_f64_store(&clock->start, aki_get_tick() - (pause - start), AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->pause, -1.0, AL_ATOMIC_RELAXED);
+}
+
+void camu_clock_seek(struct camu_clock *clock, f64 pos)
+{
+ al_atomic_f64_store(&clock->start, 0.0, AL_ATOMIC_RELAXED);
+ al_atomic_f64_store(&clock->base, pos, AL_ATOMIC_RELAXED);
+ f64 tick = aki_get_tick();
+ al_atomic_f64_store(&clock->pause, tick, AL_ATOMIC_RELAXED);
+}
+
+bool camu_clock_is_paused(struct camu_clock *clock)
+{
+ f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED);
+ return pause != -1.0;
+}
+
+f64 camu_clock_get_base_pts(struct camu_clock *clock)
+{
+ f64 base = al_atomic_f64_load(&clock->base, AL_ATOMIC_RELAXED);
+ return base;
+}
+
+f64 camu_clock_get_pts(struct camu_clock *clock)
+{
+ f64 pause = al_atomic_f64_load(&clock->pause, AL_ATOMIC_RELAXED);
+ if (pause != -1.0) return -1.0;
+ f64 start = al_atomic_f64_load(&clock->start, AL_ATOMIC_RELAXED);
+ return aki_get_tick() - start;
+}
diff --git a/src/buffer/clock.h b/src/buffer/clock.h
new file mode 100644
index 0000000..072f7d3
--- /dev/null
+++ b/src/buffer/clock.h
@@ -0,0 +1,16 @@
+#pragma once
+
+#include <aki/thread.h>
+#include <al/atomic.h>
+
+struct camu_clock {
+ atomic_f64 base, start, pause, offset;
+};
+
+void camu_clock_init(struct camu_clock *clock);
+void camu_clock_pause(struct camu_clock *clock);
+void camu_clock_resume(struct camu_clock *clock);
+void camu_clock_seek(struct camu_clock *clock, f64 pos);
+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);
diff --git a/src/buffer/common.h b/src/buffer/common.h
new file mode 100644
index 0000000..8a99c96
--- /dev/null
+++ b/src/buffer/common.h
@@ -0,0 +1,9 @@
+#pragma once
+
+enum {
+ CAMU_BUFFER_BUFFERED = 0,
+ CAMU_BUFFER_STOP,
+ CAMU_BUFFER_CONTINUE,
+ CAMU_BUFFER_PAUSED,
+ CAMU_BUFFER_EOF
+};
diff --git a/src/buffer/frame_queue.h b/src/buffer/frame_queue.h
new file mode 100644
index 0000000..35a84e4
--- /dev/null
+++ b/src/buffer/frame_queue.h
@@ -0,0 +1,28 @@
+#pragma once
+
+#ifdef HAVE_FFMPEG
+#include <libavutil/frame.h>
+#endif
+
+#include "../codec/codec.h"
+
+#include "common.h"
+
+enum {
+ CAMU_QUEUE_OK = 0,
+ CAMU_QUEUE_MORE,
+ CAMU_QUEUE_EOF,
+ CAMU_QUEUE_ERR
+};
+
+struct camu_frame_queue {
+ void (*push)(struct camu_frame_queue *, struct camu_frame *, f64);
+#ifdef HAVE_FFMPEG
+ void (*push_av_frame)(struct camu_frame_queue *, AVFrame *, f64);
+#endif
+ void (*flush)(struct camu_frame_queue *);
+ s32 (*count)(struct camu_frame_queue *);
+ u8 (*read)(struct camu_frame_queue *, f64, void *);
+ void (*reset)(struct camu_frame_queue *);
+ void (*free)(struct camu_frame_queue **);
+};
diff --git a/src/buffer/meson.build b/src/buffer/meson.build
new file mode 100644
index 0000000..1fb3cfa
--- /dev/null
+++ b/src/buffer/meson.build
@@ -0,0 +1,8 @@
+if meson.is_subproject() # TMP
+ buffer_src = ['audio.c', 'clock.c', 'peak_buffer.c']
+else
+ buffer_src = ['video.c', 'audio.c', 'clock.c', 'peak_buffer.c']
+endif
+buffer_deps = [common_deps]
+buffer = declare_dependency(sources: buffer_src,
+ dependencies: buffer_deps)
diff --git a/src/buffer/peak_buffer.c b/src/buffer/peak_buffer.c
new file mode 100644
index 0000000..22a007d
--- /dev/null
+++ b/src/buffer/peak_buffer.c
@@ -0,0 +1,24 @@
+#include "peak_buffer.h"
+
+void camu_peak_buffer_init(struct camu_peak_buffer *buf)
+{
+ aki_buffer_init(&buf->buf);
+ aki_buffer_ensure_space(&buf->buf, 128 * 1024);
+}
+
+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 = aki_buffer_get_size(&buf->buf);
+ aki_buffer_set_size(&buf->buf, 0);
+ return aki_buffer_get_ptr(&buf->buf, 0);
+}
+
+void camu_peak_buffer_free(struct camu_peak_buffer *buf)
+{
+ aki_buffer_free(&buf->buf);
+}
diff --git a/src/buffer/peak_buffer.h b/src/buffer/peak_buffer.h
new file mode 100644
index 0000000..b8d9934
--- /dev/null
+++ b/src/buffer/peak_buffer.h
@@ -0,0 +1,13 @@
+#pragma once
+
+#include <al/types.h>
+#include <aki/common.h>
+
+struct camu_peak_buffer {
+ struct aki_buffer buf;
+};
+
+void camu_peak_buffer_init(struct camu_peak_buffer *buf);
+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);
+void camu_peak_buffer_free(struct camu_peak_buffer *buf);
diff --git a/src/buffer/video.c b/src/buffer/video.c
new file mode 100644
index 0000000..d5d2c12
--- /dev/null
+++ b/src/buffer/video.c
@@ -0,0 +1,148 @@
+#include <al/log.h>
+
+#include "video.h"
+
+//#define CAMU_VIDEO_BUFFER_FORCE_SCALER
+
+bool camu_video_buffer_init(struct camu_video_buffer *buf, struct camu_clock *clock,
+ struct camu_renderer *renderer)
+{
+ buf->queue = renderer->create_queue(renderer);
+ buf->buffered = false;
+ buf->clock = clock;
+ buf->last_pts = 0.0;
+#ifdef CAMU_SCREEN_THREADED
+ al_atomic_bool_store(&buf->ref, false, AL_ATOMIC_RELAXED);
+#endif
+ return true;
+}
+
+bool camu_video_buffer_configure(struct camu_video_buffer *buf, struct camu_stream *stream)
+{
+ buf->stream = stream;
+ switch (stream->type) {
+ case CAMU_NORMAL: {
+ buf->single_frame = true;
+ buf->avg_frame_duration = 0.0;
+ break;
+ }
+#ifdef HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT: {
+ s32 width = stream->av.stream->codecpar->width;
+ s32 height = stream->av.stream->codecpar->height;
+ buf->fmt.in_width = buf->stream->video.width = width;
+ buf->fmt.in_height = buf->stream->video.height = height;
+ buf->single_frame = stream->av.stream->duration == 0 ||
+ stream->av.stream->avg_frame_rate.den == 0;
+ buf->avg_frame_duration = buf->single_frame ? 0.0 :
+ av_q2d(av_inv_q(stream->av.stream->avg_frame_rate));
+ buf->fmt.in_format = stream->av.stream->codecpar->format;
+ 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, width, height, buf->single_frame ? "IMAGE" : "VIDEO",
+ buf->single_frame ? 0.0 : 1.0 / buf->avg_frame_duration);
+#ifdef CAMU_VIDEO_BUFFER_FORCE_SCALER
+ 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)) {
+ return false;
+ }
+#endif
+ break;
+ }
+#endif
+ }
+ return true;
+}
+
+static void after_push_internal(struct camu_video_buffer *buf)
+{
+ s32 count = buf->queue->count(buf->queue);
+ if (!buf->buffered && (buf->single_frame || count >= 15)) {
+ buf->callback(buf->userdata, CAMU_BUFFER_BUFFERED);
+ buf->buffered = true;
+ }
+ if (count >= 40) buf->queue->reset(buf->queue);
+ else if (count >= 30) buf->callback(buf->userdata, CAMU_BUFFER_STOP);
+ if (buf->single_frame) camu_video_buffer_flush(buf);
+}
+
+#ifdef HAVE_FFMPEG
+static void push_av_frame_internal(struct camu_video_buffer *buf, AVFrame *frame)
+{
+ AVStream *stream = buf->stream->av.stream;
+ f64 pts = buf->single_frame ? 0.0 :
+ (frame->best_effort_timestamp - stream->start_time) * av_q2d(stream->time_base);
+ if (pts + buf->avg_frame_duration < camu_clock_get_base_pts(buf->clock)) {
+ av_frame_free(&frame);
+ return;
+ }
+ AVFrame *scaled_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)) {
+ return;
+ }
+ scaled_frame = buf->scale.frame;
+ } else {
+ scaled_frame = frame;
+ }
+#else
+ scaled_frame = frame;
+#endif
+ if (scaled_frame) scaled_frame->opaque = buf;
+ buf->queue->push_av_frame(buf->queue, scaled_frame, pts);
+}
+#endif
+
+void camu_video_buffer_push(struct camu_video_buffer *buf, struct camu_frame *frame)
+{
+ switch (frame->type) {
+ case CAMU_NORMAL:
+ buf->queue->push(buf->queue, frame, 0.0);
+ break;
+#ifdef HAVE_FFMPEG
+ case CAMU_FFMPEG_COMPAT:
+ push_av_frame_internal(buf, frame->av.frame);
+ al_free(frame);
+ break;
+#endif
+ }
+ after_push_internal(buf);
+}
+
+void camu_video_buffer_reset(struct camu_video_buffer *buf)
+{
+ buf->queue->reset(buf->queue);
+ buf->buffered = false;
+}
+
+bool camu_video_buffer_is_single_frame(struct camu_video_buffer *buf)
+{
+ return buf->single_frame;
+}
+
+void camu_video_buffer_flush(struct camu_video_buffer *buf)
+{
+ buf->queue->flush(buf->queue);
+}
+
+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);
+ if (pts == -1.0) pts = buf->last_pts;
+ buf->last_pts = pts;
+ u8 ret = buf->queue->read(buf->queue, pts, out);
+ if (ret == CAMU_QUEUE_EOF || (ret == CAMU_QUEUE_OK && buf->single_frame)) {
+ buf->callback(buf->userdata, CAMU_BUFFER_EOF);
+ } else if (buf->queue->count(buf->queue) <= 20) {
+ buf->callback(buf->userdata, CAMU_BUFFER_CONTINUE);
+ }
+ return ret == CAMU_QUEUE_OK || ret == CAMU_QUEUE_MORE;
+}
+
+void camu_video_buffer_free(struct camu_video_buffer *buf)
+{
+ buf->queue->free(&buf->queue);
+}
diff --git a/src/buffer/video.h b/src/buffer/video.h
new file mode 100644
index 0000000..534f797
--- /dev/null
+++ b/src/buffer/video.h
@@ -0,0 +1,42 @@
+#pragma once
+
+#include <al/types.h>
+
+#include "../codec/codec.h"
+#include "../codec/libav/scaler.h"
+#include "../render/renderer.h"
+#include "../screen/screen.h"
+
+#include "clock.h"
+#include "frame_queue.h"
+
+struct camu_video_buffer {
+ struct camu_stream *stream;
+
+ struct camu_clock *clock;
+ bool single_frame;
+ f64 last_pts;
+ f64 avg_frame_duration;
+
+ struct camu_lav_scaler scale;
+ struct camu_lav_scale_fmt fmt;
+
+ struct camu_frame_queue *queue;
+ bool buffered;
+
+#ifdef CAMU_SCREEN_THREADED
+ atomic_bool ref;
+#endif
+
+ void (*callback)(void *, u8);
+ void *userdata;
+};
+
+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_stream *stream);
+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);