From 60b4ebfbf3be78dba9dc7c65ab2bdaa0b218c0c2 Mon Sep 17 00:00:00 2001 From: Andrew Opalach Date: Mon, 21 Oct 2024 19:22:50 -0400 Subject: Everything before initial synced list Signed-off-by: Andrew Opalach --- src/buffer/audio.c | 427 ++++++++++++++++++++++------------------------- src/buffer/audio.h | 36 ++-- src/buffer/clock.c | 151 ++++++++--------- src/buffer/clock.h | 56 +++++-- src/buffer/frame_queue.h | 2 +- src/buffer/peak_buffer.c | 12 +- src/buffer/peak_buffer.h | 5 +- src/buffer/video.c | 74 ++++---- src/buffer/video.h | 29 ++-- src/buffer/volume.h | 36 ++++ 10 files changed, 444 insertions(+), 384 deletions(-) create mode 100644 src/buffer/volume.h (limited to 'src/buffer') 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 +#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 #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 +#include #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 +#include #include -#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 #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 -- cgit v1.2.3-101-g0448