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