use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, OnceLock};
use std::time::{Duration, Instant};
const STREAM_NOT_STARTED: i64 = i64::MIN;
const VIDEO_UNKNOWN: usize = usize::MAX;
const VIDEO_NONE: usize = usize::MAX - 1;
pub(crate) struct OutputTelemetry {
stream_watermark_us: Box<[AtomicI64]>,
stream_finished: Box<[AtomicBool]>,
video_stream: AtomicUsize,
video_packets: AtomicU64,
total_size: AtomicI64,
observed: AtomicBool,
#[cfg(test)]
size_probes: AtomicU64,
}
impl OutputTelemetry {
pub(crate) fn new(stream_count: usize) -> Self {
Self {
stream_watermark_us: (0..stream_count)
.map(|_| AtomicI64::new(STREAM_NOT_STARTED))
.collect(),
stream_finished: (0..stream_count).map(|_| AtomicBool::new(false)).collect(),
video_stream: AtomicUsize::new(VIDEO_UNKNOWN),
video_packets: AtomicU64::new(0),
total_size: AtomicI64::new(-1),
observed: AtomicBool::new(false),
#[cfg(test)]
size_probes: AtomicU64::new(0),
}
}
pub(crate) fn mark_observed(&self) {
self.observed.store(true, Ordering::Relaxed);
}
pub(crate) fn is_observed(&self) -> bool {
self.observed.load(Ordering::Relaxed)
}
#[cfg(test)]
pub(crate) fn note_perpacket_size_probe(&self) {
self.size_probes.fetch_add(1, Ordering::Relaxed);
}
#[cfg(test)]
pub(crate) fn perpacket_size_probes(&self) -> u64 {
self.size_probes.load(Ordering::Relaxed)
}
pub(crate) fn set_video_stream(&self, index: Option<usize>) {
self.video_stream
.store(index.unwrap_or(VIDEO_NONE), Ordering::Release);
}
pub(crate) fn record_written(&self, stream_index: usize, ts_us: Option<i64>) {
if stream_index >= self.stream_watermark_us.len() {
return;
}
if let Some(ts_us) = ts_us {
self.stream_watermark_us[stream_index].fetch_max(ts_us, Ordering::AcqRel);
}
if stream_index == self.video_stream.load(Ordering::Acquire) {
self.video_packets.fetch_add(1, Ordering::AcqRel);
}
}
pub(crate) fn mark_stream_finished(&self, stream_index: usize) {
if let Some(finished) = self.stream_finished.get(stream_index) {
finished.store(true, Ordering::Release);
}
}
pub(crate) fn mark_all_streams_finished(&self) {
for finished in &self.stream_finished {
finished.store(true, Ordering::Release);
}
}
pub(crate) fn set_total_size(&self, bytes: i64) {
if bytes >= 0 {
self.total_size.store(bytes, Ordering::Release);
}
}
pub(crate) fn out_time_us(&self) -> Option<i64> {
let mut active_min: Option<i64> = None;
let mut started_max: Option<i64> = None;
let mut all_finished = true;
for (watermark, finished) in self.stream_watermark_us.iter().zip(&self.stream_finished) {
let is_finished = finished.load(Ordering::Acquire);
let ts = watermark.load(Ordering::Acquire);
if ts != STREAM_NOT_STARTED {
started_max = Some(started_max.map_or(ts, |max: i64| max.max(ts)));
}
if is_finished {
continue;
}
all_finished = false;
if ts == STREAM_NOT_STARTED {
return None;
}
active_min = Some(active_min.map_or(ts, |min: i64| min.min(ts)));
}
if all_finished {
started_max
} else {
active_min
}
}
pub(crate) fn video_packets(&self) -> Option<u64> {
match self.video_stream.load(Ordering::Acquire) {
VIDEO_UNKNOWN | VIDEO_NONE => None,
_ => Some(self.video_packets.load(Ordering::Acquire)),
}
}
pub(crate) fn total_size(&self) -> Option<u64> {
let bytes = self.total_size.load(Ordering::Acquire);
(bytes >= 0).then_some(bytes as u64)
}
}
pub(crate) struct ProgressTracker {
started_at: OnceLock<Instant>,
completed_at: OnceLock<Instant>,
outputs: Box<[Arc<OutputTelemetry>]>,
demux_exited: Box<[Arc<AtomicBool>]>,
frame_source_exited: Box<[Arc<AtomicBool>]>,
}
impl ProgressTracker {
pub(crate) fn new(
outputs: Vec<Arc<OutputTelemetry>>,
demux_exited: Vec<Arc<AtomicBool>>,
frame_source_count: usize,
) -> Self {
Self {
started_at: OnceLock::new(),
completed_at: OnceLock::new(),
outputs: outputs.into_boxed_slice(),
demux_exited: demux_exited.into_boxed_slice(),
frame_source_exited: (0..frame_source_count)
.map(|_| Arc::new(AtomicBool::new(false)))
.collect(),
}
}
pub(crate) fn mark_started(&self) {
let _ = self.started_at.set(Instant::now());
}
pub(crate) fn mark_observed(&self) {
for output in &self.outputs {
output.mark_observed();
}
}
pub(crate) fn seal_completed(&self) {
let _ = self.completed_at.set(Instant::now());
}
pub(crate) fn is_completed(&self) -> bool {
self.completed_at.get().is_some()
}
pub(crate) fn elapsed(&self) -> Duration {
let Some(&started_at) = self.started_at.get() else {
return Duration::ZERO;
};
match self.completed_at.get() {
Some(&completed_at) => completed_at.saturating_duration_since(started_at),
None => started_at.elapsed(),
}
}
pub(crate) fn outputs(&self) -> &[Arc<OutputTelemetry>] {
&self.outputs
}
pub(crate) fn output_telemetry(&self, index: usize) -> Arc<OutputTelemetry> {
self.outputs[index].clone()
}
pub(crate) fn frame_source_exit_flag(&self, index: usize) -> Arc<AtomicBool> {
self.frame_source_exited[index].clone()
}
pub(crate) fn inputs_drained(&self) -> bool {
let producers = self.demux_exited.iter().chain(&self.frame_source_exited);
let mut any = false;
for exited in producers {
any = true;
if !exited.load(Ordering::Acquire) {
return false;
}
}
any
}
}