pub use crate::error::PacketSinkError;
use crate::core::scheduler::ffmpeg_scheduler::{is_stopping, FfmpegScheduler, Running};
use crate::core::scheduler::owned_run_iter::OwnedRunIter;
use ffmpeg_sys_next::{AVCodecID, AVMediaType, AVRational};
use std::num::NonZeroUsize;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, OnceLock};
use std::time::Duration;
#[cfg(test)]
mod bench_nal_scan;
pub(crate) mod codec;
mod job_failure;
pub(crate) mod nal_framing;
pub(crate) mod side_data;
pub(crate) mod strict;
pub(crate) mod timeline;
pub use job_failure::{JobFailureKind, JobFailureSummary};
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum PacketSinkTier {
#[default]
Strict,
}
#[derive(Debug, Clone)]
pub struct PacketCallbackError {
message: String,
source: Option<Arc<dyn std::error::Error + Send + Sync + 'static>>,
pub(crate) kind: CallbackFailureKind,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum CallbackFailureKind {
Failure,
Disconnected,
Cancelled,
JobStopped,
}
impl PacketCallbackError {
pub fn new(message: impl Into<String>) -> Self {
Self {
message: message.into(),
source: None,
kind: CallbackFailureKind::Failure,
}
}
pub fn with_source(
message: impl Into<String>,
source: impl std::error::Error + Send + Sync + 'static,
) -> Self {
Self {
message: message.into(),
source: Some(Arc::new(source)),
kind: CallbackFailureKind::Failure,
}
}
pub(crate) fn disconnected() -> Self {
Self {
message: "packet-sink channel receiver dropped".to_string(),
source: None,
kind: CallbackFailureKind::Disconnected,
}
}
pub(crate) fn job_stopped() -> Self {
Self {
message: "job failed elsewhere; blocking send abandoned".to_string(),
source: None,
kind: CallbackFailureKind::JobStopped,
}
}
pub(crate) fn cancelled() -> Self {
Self {
message: "job stopping; blocking send cancelled".to_string(),
source: None,
kind: CallbackFailureKind::Cancelled,
}
}
}
impl std::fmt::Display for PacketCallbackError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.message)
}
}
impl std::error::Error for PacketCallbackError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
self.source
.as_ref()
.map(|s| s.as_ref() as &(dyn std::error::Error + 'static))
}
}
pub type PacketCallbackResult = Result<(), PacketCallbackError>;
pub trait PacketSinkHandler: Send + 'static {
fn on_stream_info(&mut self, _streams: &[PacketStreamInfo]) -> PacketCallbackResult {
Ok(())
}
fn on_packet(&mut self, packet: &PacketView<'_>) -> PacketCallbackResult;
fn on_end(&mut self) {}
fn on_job_failed(&mut self, _summary: &JobFailureSummary) {}
fn on_delivery_error(&mut self, _error: &PacketSinkError) {}
}
#[non_exhaustive]
#[derive(Debug, Clone)]
pub struct VideoPacketConfig {
pub(crate) stream_index: usize,
pub(crate) codec_id: AVCodecID,
pub(crate) codec_string: String,
pub(crate) profile: u8,
pub(crate) compatibility: u8,
pub(crate) level: u8,
pub(crate) codec_config: Vec<u8>,
pub(crate) time_base: AVRational,
pub(crate) width: i32,
pub(crate) height: i32,
pub(crate) sample_aspect_ratio: Option<AVRational>,
pub(crate) frame_rate: Option<AVRational>,
}
impl VideoPacketConfig {
pub fn stream_index(&self) -> usize {
self.stream_index
}
pub fn codec_id(&self) -> AVCodecID {
self.codec_id
}
pub fn codec_string(&self) -> &str {
&self.codec_string
}
pub fn profile(&self) -> u8 {
self.profile
}
pub fn compatibility(&self) -> u8 {
self.compatibility
}
pub fn level(&self) -> u8 {
self.level
}
pub fn codec_config(&self) -> &[u8] {
&self.codec_config
}
pub fn extradata(&self) -> &[u8] {
&self.codec_config
}
pub fn time_base(&self) -> AVRational {
self.time_base
}
pub fn width(&self) -> i32 {
self.width
}
pub fn height(&self) -> i32 {
self.height
}
pub fn sample_aspect_ratio(&self) -> Option<AVRational> {
self.sample_aspect_ratio
}
pub fn frame_rate(&self) -> Option<AVRational> {
self.frame_rate
}
}
#[non_exhaustive]
#[derive(Debug, Clone)]
pub struct AudioPacketConfig {
pub(crate) stream_index: usize,
pub(crate) codec_id: AVCodecID,
pub(crate) codec_string: String,
pub(crate) codec_config: Vec<u8>,
pub(crate) time_base: AVRational,
pub(crate) sample_rate: i32,
pub(crate) channels: i32,
pub(crate) channel_layout: String,
}
impl AudioPacketConfig {
pub fn stream_index(&self) -> usize {
self.stream_index
}
pub fn codec_id(&self) -> AVCodecID {
self.codec_id
}
pub fn codec_string(&self) -> &str {
&self.codec_string
}
pub fn codec_config(&self) -> &[u8] {
&self.codec_config
}
pub fn extradata(&self) -> &[u8] {
&self.codec_config
}
pub fn time_base(&self) -> AVRational {
self.time_base
}
pub fn sample_rate(&self) -> i32 {
self.sample_rate
}
pub fn channels(&self) -> i32 {
self.channels
}
pub fn channel_layout(&self) -> &str {
&self.channel_layout
}
}
#[non_exhaustive]
#[derive(Debug, Clone)]
pub enum PacketStreamInfo {
Video(VideoPacketConfig),
Audio(AudioPacketConfig),
}
impl PacketStreamInfo {
pub fn stream_index(&self) -> usize {
match self {
PacketStreamInfo::Video(v) => v.stream_index,
PacketStreamInfo::Audio(a) => a.stream_index,
}
}
pub fn media_type(&self) -> AVMediaType {
match self {
PacketStreamInfo::Video(_) => AVMediaType::AVMEDIA_TYPE_VIDEO,
PacketStreamInfo::Audio(_) => AVMediaType::AVMEDIA_TYPE_AUDIO,
}
}
pub fn codec_id(&self) -> AVCodecID {
match self {
PacketStreamInfo::Video(v) => v.codec_id,
PacketStreamInfo::Audio(a) => a.codec_id,
}
}
pub fn codec_string(&self) -> &str {
match self {
PacketStreamInfo::Video(v) => &v.codec_string,
PacketStreamInfo::Audio(a) => &a.codec_string,
}
}
pub fn codec_config(&self) -> &[u8] {
match self {
PacketStreamInfo::Video(v) => &v.codec_config,
PacketStreamInfo::Audio(a) => &a.codec_config,
}
}
pub fn extradata(&self) -> &[u8] {
self.codec_config()
}
pub fn time_base(&self) -> AVRational {
match self {
PacketStreamInfo::Video(v) => v.time_base,
PacketStreamInfo::Audio(a) => a.time_base,
}
}
pub fn video(&self) -> Option<&VideoPacketConfig> {
match self {
PacketStreamInfo::Video(v) => Some(v),
_ => None,
}
}
pub fn audio(&self) -> Option<&AudioPacketConfig> {
match self {
PacketStreamInfo::Audio(a) => Some(a),
_ => None,
}
}
}
fn ticks_to_us(ticks: i64, time_base: AVRational) -> i64 {
unsafe {
ffmpeg_sys_next::av_rescale_q(
ticks,
time_base,
AVRational {
num: 1,
den: 1_000_000,
},
)
}
}
#[non_exhaustive]
#[derive(Debug)]
pub struct PacketView<'a> {
pub(crate) stream_index: usize,
pub(crate) pts: i64,
pub(crate) dts: i64,
pub(crate) duration: i64,
pub(crate) time_base: AVRational,
pub(crate) is_key: bool,
pub(crate) applied_offset: i64,
pub(crate) data: &'a [u8],
}
impl<'a> PacketView<'a> {
pub fn stream_index(&self) -> usize {
self.stream_index
}
pub fn pts(&self) -> i64 {
self.pts
}
pub fn dts(&self) -> i64 {
self.dts
}
pub fn duration(&self) -> i64 {
self.duration
}
pub fn time_base(&self) -> AVRational {
self.time_base
}
pub fn pts_us(&self) -> i64 {
ticks_to_us(self.pts, self.time_base)
}
pub fn dts_us(&self) -> i64 {
ticks_to_us(self.dts, self.time_base)
}
pub fn duration_us(&self) -> i64 {
ticks_to_us(self.duration, self.time_base)
}
pub fn applied_offset_us(&self) -> i64 {
ticks_to_us(self.applied_offset, self.time_base)
}
pub fn is_key(&self) -> bool {
self.is_key
}
pub fn applied_offset(&self) -> i64 {
self.applied_offset
}
pub fn data(&self) -> &'a [u8] {
self.data
}
}
pub(crate) type StreamInfoFn =
Box<dyn FnMut(&[PacketStreamInfo]) -> PacketCallbackResult + Send>;
pub(crate) type PacketFn =
Box<dyn for<'a> FnMut(&PacketView<'a>) -> PacketCallbackResult + Send>;
pub(crate) type EndFn = Box<dyn FnMut() + Send>;
pub(crate) type JobFailedFn = Box<dyn FnMut(&JobFailureSummary) + Send>;
pub(crate) type DeliveryErrorFn = Box<dyn FnMut(&PacketSinkError) + Send>;
enum SinkDispatch {
Closures {
on_stream_info: Option<StreamInfoFn>,
on_packet: PacketFn,
on_end: Option<EndFn>,
on_job_failed: Option<JobFailedFn>,
on_delivery_error: Option<DeliveryErrorFn>,
},
Handler(Box<dyn PacketSinkHandler>),
}
pub(crate) struct JobStopObservables {
pub(crate) status: Arc<AtomicUsize>,
pub(crate) result: Arc<std::sync::Mutex<Option<crate::error::Result<()>>>>,
}
pub(crate) type CancellationSlot = Arc<OnceLock<JobStopObservables>>;
pub struct PacketSink {
pub(crate) tier: PacketSinkTier,
dispatch: SinkDispatch,
pub(crate) cancellation: Option<CancellationSlot>,
}
impl std::fmt::Debug for PacketSink {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PacketSink")
.field("tier", &self.tier)
.finish_non_exhaustive()
}
}
impl PacketSink {
pub fn tier(&self) -> PacketSinkTier {
self.tier
}
pub fn builder<F>(on_packet: F) -> PacketSinkBuilder
where
F: for<'a> FnMut(&PacketView<'a>) -> PacketCallbackResult + Send + 'static,
{
PacketSinkBuilder {
tier: PacketSinkTier::Strict,
on_stream_info: None,
on_packet: Box::new(on_packet),
on_end: None,
on_job_failed: None,
on_delivery_error: None,
}
}
pub fn discard() -> PacketSink {
PacketSink::builder(|_| Ok(())).build()
}
pub fn from_handler<H: PacketSinkHandler>(handler: H) -> PacketSink {
PacketSink {
tier: PacketSinkTier::Strict,
dispatch: SinkDispatch::Handler(Box::new(handler)),
cancellation: None,
}
}
pub fn channel(capacity: NonZeroUsize) -> (PacketSink, PacketSinkReceiver) {
let (tx, rx) = crossbeam_channel::bounded::<PacketSinkEvent>(capacity.get());
let cancellation: CancellationSlot = Arc::new(OnceLock::new());
let info_tx = tx.clone();
let info_cancel = cancellation.clone();
let pkt_tx = tx.clone();
let pkt_cancel = cancellation.clone();
let end_tx = tx.clone();
let job_failed_tx = tx.clone();
let err_tx = tx;
let mut sink = PacketSink::builder(move |packet: &PacketView<'_>| {
send_with_cancellation(
&pkt_tx,
&pkt_cancel,
PacketSinkEvent::Packet(EncodedPacket::from_view(packet)),
)
})
.on_stream_info(move |infos: &[PacketStreamInfo]| {
send_with_cancellation(
&info_tx,
&info_cancel,
PacketSinkEvent::StreamInfo(infos.to_vec()),
)
})
.on_end(move || {
let _ = end_tx.try_send(PacketSinkEvent::End);
})
.on_job_failed(move |summary: &JobFailureSummary| {
let free = job_failed_tx
.capacity()
.unwrap_or(usize::MAX)
.saturating_sub(job_failed_tx.len());
if free >= 2 {
let _ = job_failed_tx.try_send(PacketSinkEvent::JobFailure(summary.clone()));
}
})
.on_delivery_error(move |e: &PacketSinkError| {
let _ = err_tx.try_send(PacketSinkEvent::Error(e.clone()));
})
.build();
sink.cancellation = Some(cancellation.clone());
(
sink,
PacketSinkReceiver {
inner: rx,
token: cancellation,
},
)
}
pub(crate) fn dispatch_stream_info(
&mut self,
infos: &[PacketStreamInfo],
) -> PacketCallbackResult {
match &mut self.dispatch {
SinkDispatch::Closures { on_stream_info, .. } => match on_stream_info {
Some(f) => f(infos),
None => Ok(()),
},
SinkDispatch::Handler(h) => h.on_stream_info(infos),
}
}
pub(crate) fn dispatch_packet(&mut self, packet: &PacketView<'_>) -> PacketCallbackResult {
match &mut self.dispatch {
SinkDispatch::Closures { on_packet, .. } => on_packet(packet),
SinkDispatch::Handler(h) => h.on_packet(packet),
}
}
pub(crate) fn dispatch_end(&mut self) {
match &mut self.dispatch {
SinkDispatch::Closures { on_end, .. } => {
if let Some(f) = on_end {
f()
}
}
SinkDispatch::Handler(h) => h.on_end(),
}
}
pub(crate) fn dispatch_job_failed(&mut self, summary: &JobFailureSummary) {
match &mut self.dispatch {
SinkDispatch::Closures { on_job_failed, .. } => {
if let Some(f) = on_job_failed {
f(summary)
}
}
SinkDispatch::Handler(h) => h.on_job_failed(summary),
}
}
pub(crate) fn dispatch_delivery_error(&mut self, error: &PacketSinkError) {
match &mut self.dispatch {
SinkDispatch::Closures {
on_delivery_error, ..
} => {
if let Some(f) = on_delivery_error {
f(error)
}
}
SinkDispatch::Handler(h) => h.on_delivery_error(error),
}
}
pub(crate) fn dispose_contained(self) -> bool {
let Self {
tier: _,
dispatch,
cancellation,
} = self;
let mut panicked = false;
match dispatch {
SinkDispatch::Closures {
on_stream_info,
on_packet,
on_end,
on_job_failed,
on_delivery_error,
} => {
if let Some(f) = on_stream_info {
panicked |= drop_contained(f);
}
panicked |= drop_contained(on_packet);
if let Some(f) = on_end {
panicked |= drop_contained(f);
}
if let Some(f) = on_job_failed {
panicked |= drop_contained(f);
}
if let Some(f) = on_delivery_error {
panicked |= drop_contained(f);
}
}
SinkDispatch::Handler(handler) => {
panicked |= drop_contained(handler);
}
}
if let Some(slot) = cancellation {
panicked |= drop_contained(slot);
}
panicked
}
}
fn drop_contained<T>(value: T) -> bool {
match std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || drop(value))) {
Ok(()) => false,
Err(payload) => {
dispose_panic_payload(payload);
true
}
}
}
pub(crate) fn dispose_panic_payload(payload: Box<dyn std::any::Any + Send>) {
let mut payload = payload;
for _ in 0..4 {
match std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || drop(payload))) {
Ok(()) => return,
Err(next) => payload = next,
}
}
std::mem::forget(payload);
}
fn send_with_cancellation(
tx: &crossbeam_channel::Sender<PacketSinkEvent>,
cancellation: &CancellationSlot,
event: PacketSinkEvent,
) -> PacketCallbackResult {
let mut event = match tx.try_send(event) {
Ok(()) => return Ok(()),
Err(crossbeam_channel::TrySendError::Disconnected(_)) => {
return Err(PacketCallbackError::disconnected());
}
Err(crossbeam_channel::TrySendError::Full(back)) => back,
};
loop {
match tx.send_timeout(event, Duration::from_millis(50)) {
Ok(()) => return Ok(()),
Err(crossbeam_channel::SendTimeoutError::Timeout(back)) => {
event = back;
if let Some(observables) = cancellation.get() {
if is_stopping(observables.status.load(Ordering::Acquire)) {
let failed = observables
.result
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.as_ref()
.is_some_and(|result| result.is_err());
return Err(if failed {
PacketCallbackError::job_stopped()
} else {
PacketCallbackError::cancelled()
});
}
}
}
Err(crossbeam_channel::SendTimeoutError::Disconnected(_)) => {
return Err(PacketCallbackError::disconnected());
}
}
}
}
pub struct PacketSinkBuilder {
tier: PacketSinkTier,
on_stream_info: Option<StreamInfoFn>,
on_packet: PacketFn,
on_end: Option<EndFn>,
on_job_failed: Option<JobFailedFn>,
on_delivery_error: Option<DeliveryErrorFn>,
}
impl PacketSinkBuilder {
pub fn on_stream_info<F>(mut self, f: F) -> Self
where
F: FnMut(&[PacketStreamInfo]) -> PacketCallbackResult + Send + 'static,
{
self.on_stream_info = Some(Box::new(f));
self
}
pub fn on_end<F>(mut self, f: F) -> Self
where
F: FnMut() + Send + 'static,
{
self.on_end = Some(Box::new(f));
self
}
pub fn on_job_failed<F>(mut self, f: F) -> Self
where
F: FnMut(&JobFailureSummary) + Send + 'static,
{
self.on_job_failed = Some(Box::new(f));
self
}
pub fn on_delivery_error<F>(mut self, f: F) -> Self
where
F: FnMut(&PacketSinkError) + Send + 'static,
{
self.on_delivery_error = Some(Box::new(f));
self
}
pub fn build(self) -> PacketSink {
PacketSink {
tier: self.tier,
dispatch: SinkDispatch::Closures {
on_stream_info: self.on_stream_info,
on_packet: self.on_packet,
on_end: self.on_end,
on_job_failed: self.on_job_failed,
on_delivery_error: self.on_delivery_error,
},
cancellation: None,
}
}
}
#[non_exhaustive]
#[derive(Debug, Clone)]
pub struct EncodedPacket {
pub(crate) stream_index: usize,
pub(crate) pts: i64,
pub(crate) dts: i64,
pub(crate) duration: i64,
pub(crate) time_base: AVRational,
pub(crate) is_key: bool,
pub(crate) applied_offset: i64,
pub(crate) data: Vec<u8>,
}
impl EncodedPacket {
fn from_view(view: &PacketView<'_>) -> Self {
Self {
stream_index: view.stream_index,
pts: view.pts,
dts: view.dts,
duration: view.duration,
time_base: view.time_base,
is_key: view.is_key,
applied_offset: view.applied_offset,
data: view.data.to_vec(),
}
}
pub fn stream_index(&self) -> usize {
self.stream_index
}
pub fn pts(&self) -> i64 {
self.pts
}
pub fn dts(&self) -> i64 {
self.dts
}
pub fn duration(&self) -> i64 {
self.duration
}
pub fn time_base(&self) -> AVRational {
self.time_base
}
pub fn pts_us(&self) -> i64 {
ticks_to_us(self.pts, self.time_base)
}
pub fn dts_us(&self) -> i64 {
ticks_to_us(self.dts, self.time_base)
}
pub fn duration_us(&self) -> i64 {
ticks_to_us(self.duration, self.time_base)
}
pub fn applied_offset_us(&self) -> i64 {
ticks_to_us(self.applied_offset, self.time_base)
}
pub fn is_key(&self) -> bool {
self.is_key
}
pub fn applied_offset(&self) -> i64 {
self.applied_offset
}
pub fn data(&self) -> &[u8] {
&self.data
}
pub fn into_data(self) -> Vec<u8> {
self.data
}
}
#[non_exhaustive]
#[derive(Debug, Clone)]
pub enum PacketSinkEvent {
StreamInfo(Vec<PacketStreamInfo>),
Packet(EncodedPacket),
End,
JobFailure(JobFailureSummary),
Error(PacketSinkError),
}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PacketRecvError {
Disconnected,
}
impl std::fmt::Display for PacketRecvError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("packet-sink channel disconnected")
}
}
impl std::error::Error for PacketRecvError {}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PacketTryRecvError {
Empty,
Disconnected,
}
impl std::fmt::Display for PacketTryRecvError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
PacketTryRecvError::Empty => f.write_str("no packet-sink event queued"),
PacketTryRecvError::Disconnected => f.write_str("packet-sink channel disconnected"),
}
}
}
impl std::error::Error for PacketTryRecvError {}
#[non_exhaustive]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PacketRecvTimeoutError {
Timeout,
Disconnected,
}
impl std::fmt::Display for PacketRecvTimeoutError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
PacketRecvTimeoutError::Timeout => {
f.write_str("timed out waiting for a packet-sink event")
}
PacketRecvTimeoutError::Disconnected => {
f.write_str("packet-sink channel disconnected")
}
}
}
}
impl std::error::Error for PacketRecvTimeoutError {}
pub struct PacketEventsPairingError {
pub receiver: PacketSinkReceiver,
pub scheduler: FfmpegScheduler<Running>,
}
impl std::fmt::Debug for PacketEventsPairingError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PacketEventsPairingError")
.finish_non_exhaustive()
}
}
impl std::fmt::Display for PacketEventsPairingError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(
"packet-sink receiver paired with a scheduler that is not running its sink",
)
}
}
impl std::error::Error for PacketEventsPairingError {}
pub struct PacketSinkReceiver {
inner: crossbeam_channel::Receiver<PacketSinkEvent>,
token: CancellationSlot,
}
impl PacketSinkReceiver {
pub fn recv(&self) -> Result<PacketSinkEvent, PacketRecvError> {
self.inner.recv().map_err(|_| PacketRecvError::Disconnected)
}
pub fn try_recv(&self) -> Result<PacketSinkEvent, PacketTryRecvError> {
self.inner.try_recv().map_err(|e| match e {
crossbeam_channel::TryRecvError::Empty => PacketTryRecvError::Empty,
crossbeam_channel::TryRecvError::Disconnected => PacketTryRecvError::Disconnected,
})
}
pub fn recv_timeout(
&self,
timeout: Duration,
) -> Result<PacketSinkEvent, PacketRecvTimeoutError> {
self.inner.recv_timeout(timeout).map_err(|e| match e {
crossbeam_channel::RecvTimeoutError::Timeout => PacketRecvTimeoutError::Timeout,
crossbeam_channel::RecvTimeoutError::Disconnected => {
PacketRecvTimeoutError::Disconnected
}
})
}
pub fn iter(&self) -> impl Iterator<Item = PacketSinkEvent> + '_ {
self.inner.iter()
}
#[allow(clippy::result_large_err)]
pub fn into_events(
self,
scheduler: FfmpegScheduler<Running>,
) -> Result<PacketEventIter, PacketEventsPairingError> {
if !scheduler.runs_packet_sink(&self.token) {
return Err(PacketEventsPairingError {
receiver: self,
scheduler,
});
}
Ok(PacketEventIter {
inner: OwnedRunIter::new(self.inner, scheduler, std::convert::identity),
saw_end: false,
saw_error: false,
})
}
}
pub struct PacketEventIter {
inner: OwnedRunIter<PacketSinkEvent>,
saw_end: bool,
saw_error: bool,
}
impl std::fmt::Debug for PacketEventIter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PacketEventIter").finish_non_exhaustive()
}
}
impl Iterator for PacketEventIter {
type Item = Result<PacketSinkEvent, crate::error::Error>;
fn next(&mut self) -> Option<Self::Item> {
match self.inner.next() {
Some(Ok(event)) => {
if matches!(event, PacketSinkEvent::End) {
self.saw_end = true;
}
Some(Ok(event))
}
Some(Err(e)) => {
self.saw_error = true;
Some(Err(e))
}
None => {
if !self.saw_end && !self.saw_error {
self.saw_end = true;
Some(Ok(PacketSinkEvent::End))
} else {
None
}
}
}
}
}
impl std::iter::FusedIterator for PacketEventIter {}
pub(crate) struct PacketSinkPolicy {
pub(crate) global_header: bool,
pub(crate) variable_fps: bool,
pub(crate) no_timestamps: bool,
}
impl PacketSinkPolicy {
pub(crate) fn for_tier(tier: PacketSinkTier) -> Self {
match tier {
PacketSinkTier::Strict => Self {
global_header: true,
variable_fps: false,
no_timestamps: false,
},
}
}
pub(crate) fn oformat_flags(&self) -> i32 {
let mut flags = 0;
if self.global_header {
flags |= ffmpeg_sys_next::AVFMT_GLOBALHEADER;
}
if self.variable_fps {
flags |= ffmpeg_sys_next::AVFMT_VARIABLE_FPS;
}
if self.no_timestamps {
flags |= ffmpeg_sys_next::AVFMT_NOTIMESTAMPS;
}
flags
}
}
#[cfg(test)]
mod tests {
use super::*;
fn test_view(data: &[u8]) -> PacketView<'_> {
PacketView {
stream_index: 0,
pts: 10,
dts: 5,
duration: 1,
time_base: AVRational { num: 1, den: 25 },
is_key: false,
applied_offset: 3,
data,
}
}
#[test]
fn builder_requires_a_packet_consumer_and_discard_is_explicit() {
let mut sink = PacketSink::builder(|_| Ok(())).build();
assert_eq!(sink.tier, PacketSinkTier::Strict);
assert!(sink.dispatch_stream_info(&[]).is_ok());
let payload = [0u8, 0, 0, 1, 0x65];
assert!(sink.dispatch_packet(&test_view(&payload)).is_ok());
sink.dispatch_end();
sink.dispatch_delivery_error(&PacketSinkError::NoStreams);
let mut discard = PacketSink::discard();
assert!(discard.dispatch_packet(&test_view(&payload)).is_ok());
}
#[test]
fn every_construction_path_reports_the_strict_tier() {
assert_eq!(PacketSinkTier::default(), PacketSinkTier::Strict);
assert_eq!(
PacketSink::builder(|_| Ok(())).build().tier(),
PacketSinkTier::Strict
);
assert_eq!(PacketSink::discard().tier(), PacketSinkTier::Strict);
struct Accepting;
impl PacketSinkHandler for Accepting {
fn on_packet(&mut self, _packet: &PacketView<'_>) -> PacketCallbackResult {
Ok(())
}
}
assert_eq!(
PacketSink::from_handler(Accepting).tier(),
PacketSinkTier::Strict
);
let (sink, _receiver) = PacketSink::channel(NonZeroUsize::new(1).unwrap());
assert_eq!(sink.tier(), PacketSinkTier::Strict);
}
#[test]
fn handler_receives_serial_callbacks_with_shared_state() {
struct Counting {
packets: usize,
}
impl PacketSinkHandler for Counting {
fn on_packet(&mut self, _packet: &PacketView<'_>) -> PacketCallbackResult {
self.packets += 1;
if self.packets > 1 {
Err(PacketCallbackError::new("enough"))
} else {
Ok(())
}
}
}
let mut sink = PacketSink::from_handler(Counting { packets: 0 });
let payload = [0u8, 0, 0, 1, 0x65];
assert!(sink.dispatch_packet(&test_view(&payload)).is_ok());
let err = sink
.dispatch_packet(&test_view(&payload))
.expect_err("handler state must persist across calls");
assert_eq!(err.kind, CallbackFailureKind::Failure);
assert_eq!(err.to_string(), "enough");
}
#[test]
fn panic_payload_chains_are_disposed_without_escaping() {
struct ChainBomb(u32);
impl Drop for ChainBomb {
fn drop(&mut self) {
if self.0 > 0 {
std::panic::panic_any(ChainBomb(self.0 - 1));
}
}
}
dispose_panic_payload(Box::new(ChainBomb(3)));
dispose_panic_payload(Box::new(ChainBomb(64)));
}
#[test]
fn dispose_contained_destroys_every_box_across_multiple_drop_panics() {
use std::sync::atomic::AtomicBool;
struct DropBomb(Arc<AtomicBool>);
impl Drop for DropBomb {
fn drop(&mut self) {
self.0.store(true, Ordering::Release);
panic!("injected capture-destructor panic");
}
}
let flags: Vec<Arc<AtomicBool>> = (0..4).map(|_| Arc::new(AtomicBool::new(false))).collect();
let (b0, b1, b2, b3) = (
DropBomb(flags[0].clone()),
DropBomb(flags[1].clone()),
DropBomb(flags[2].clone()),
DropBomb(flags[3].clone()),
);
let sink = PacketSink::builder(move |_pkt| {
let _hold = &b0;
Ok(())
})
.on_end(move || {
let _hold = &b1;
})
.on_job_failed(move |_summary| {
let _hold = &b2;
})
.on_delivery_error(move |_e| {
let _hold = &b3;
})
.build();
assert!(
sink.dispose_contained(),
"four panicking capture destructors must be reported"
);
for (i, flag) in flags.iter().enumerate() {
assert!(
flag.load(Ordering::Acquire),
"callback box {i} was never destroyed"
);
}
struct BombHandler(Arc<AtomicBool>);
impl Drop for BombHandler {
fn drop(&mut self) {
self.0.store(true, Ordering::Release);
panic!("injected handler-destructor panic");
}
}
impl PacketSinkHandler for BombHandler {
fn on_packet(&mut self, _packet: &PacketView<'_>) -> PacketCallbackResult {
Ok(())
}
}
let destroyed = Arc::new(AtomicBool::new(false));
assert!(PacketSink::from_handler(BombHandler(destroyed.clone())).dispose_contained());
assert!(destroyed.load(Ordering::Acquire));
assert!(!PacketSink::builder(|_pkt| Ok(())).build().dispose_contained());
}
#[test]
fn callback_error_preserves_its_source() {
let io = std::io::Error::new(std::io::ErrorKind::BrokenPipe, "peer gone");
let err = PacketCallbackError::with_source("send failed", io);
assert_eq!(err.to_string(), "send failed");
let source = std::error::Error::source(&err).expect("source preserved");
assert!(source.to_string().contains("peer gone"));
}
#[test]
fn channel_adapter_forwards_events_in_order() {
let (mut sink, rx) = PacketSink::channel(NonZeroUsize::new(8).unwrap());
assert!(sink.dispatch_stream_info(&[]).is_ok());
let payload = [0u8, 0, 0, 1, 0x65];
assert!(sink.dispatch_packet(&test_view(&payload)).is_ok());
sink.dispatch_end();
match rx.recv().unwrap() {
PacketSinkEvent::StreamInfo(v) => assert!(v.is_empty()),
other => panic!("expected StreamInfo, got {other:?}"),
}
match rx.recv().unwrap() {
PacketSinkEvent::Packet(p) => {
assert_eq!(p.pts(), 10);
assert_eq!(p.dts(), 5);
assert_eq!(p.applied_offset(), 3);
assert_eq!(p.data(), &payload);
assert!(!p.is_key());
assert_eq!(p.pts_us(), 400_000);
assert_eq!(p.duration_us(), 40_000);
}
other => panic!("expected Packet, got {other:?}"),
}
assert!(matches!(rx.recv().unwrap(), PacketSinkEvent::End));
drop(sink);
assert!(matches!(rx.recv(), Err(PacketRecvError::Disconnected)));
}
#[test]
fn dropped_receiver_turns_sends_into_typed_disconnection() {
let (mut sink, rx) = PacketSink::channel(NonZeroUsize::new(1).unwrap());
drop(rx);
let payload = [0u8, 0, 0, 1, 0x65];
let err = sink
.dispatch_packet(&test_view(&payload))
.expect_err("send into a dropped receiver must fail");
assert_eq!(err.kind, CallbackFailureKind::Disconnected);
}
#[test]
fn blocked_channel_send_observes_cancellation() {
let (mut sink, rx) = PacketSink::channel(NonZeroUsize::new(1).unwrap());
let status = Arc::new(AtomicUsize::new(
crate::core::scheduler::ffmpeg_scheduler::STATUS_RUN,
));
let result = Arc::new(std::sync::Mutex::new(None));
sink.cancellation
.as_ref()
.expect("channel sinks carry a cancellation slot")
.set(JobStopObservables {
status: status.clone(),
result,
})
.ok();
let payload = [0u8, 0, 0, 1, 0x65];
assert!(sink.dispatch_packet(&test_view(&payload)).is_ok());
let flip = status.clone();
let flipper = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(120));
flip.store(
crate::core::scheduler::ffmpeg_scheduler::STATUS_END,
Ordering::Release,
);
});
let start = std::time::Instant::now();
let err = sink
.dispatch_packet(&test_view(&payload))
.expect_err("blocked send must cancel");
assert_eq!(err.kind, CallbackFailureKind::Cancelled);
assert!(
start.elapsed() < Duration::from_secs(5),
"cancellation must be prompt"
);
flipper.join().unwrap();
drop(rx);
}
#[test]
fn blocked_channel_send_classifies_failure_driven_stop() {
let (mut sink, rx) = PacketSink::channel(NonZeroUsize::new(1).unwrap());
let status = Arc::new(AtomicUsize::new(
crate::core::scheduler::ffmpeg_scheduler::STATUS_RUN,
));
let result: Arc<std::sync::Mutex<Option<crate::error::Result<()>>>> =
Arc::new(std::sync::Mutex::new(None));
sink.cancellation
.as_ref()
.expect("channel sinks carry a cancellation slot")
.set(JobStopObservables {
status: status.clone(),
result: result.clone(),
})
.ok();
let payload = [0u8, 0, 0, 1, 0x65];
assert!(sink.dispatch_packet(&test_view(&payload)).is_ok());
let flipper = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(120));
*result.lock().unwrap() = Some(Err(crate::error::Error::WorkerPanicked(
"muxer1:mpegts".to_string(),
)));
status.store(
crate::core::scheduler::ffmpeg_scheduler::STATUS_END,
Ordering::Release,
);
});
let err = sink
.dispatch_packet(&test_view(&payload))
.expect_err("blocked send must abandon on a failed job");
assert_eq!(err.kind, CallbackFailureKind::JobStopped);
flipper.join().unwrap();
drop(rx);
}
#[test]
fn recv_variants_distinguish_empty_timeout_disconnected() {
let (sink, rx) = PacketSink::channel(NonZeroUsize::new(1).unwrap());
assert_eq!(rx.try_recv().unwrap_err(), PacketTryRecvError::Empty);
assert_eq!(
rx.recv_timeout(Duration::from_millis(10)).unwrap_err(),
PacketRecvTimeoutError::Timeout
);
drop(sink);
assert_eq!(rx.try_recv().unwrap_err(), PacketTryRecvError::Disconnected);
assert_eq!(
rx.recv_timeout(Duration::from_millis(10)).unwrap_err(),
PacketRecvTimeoutError::Disconnected
);
}
}