Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 3 additions & 8 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -60,9 +60,8 @@ info.format.duration_s
The Rust crate is built on rsmpeg's safe wrappers with
`#![deny(unsafe_code)]` at the root.
`native/exmpeg_native/src/ffi_helpers.rs` is the only module that
contains `unsafe`; everything else, including the progress emitter that
reconstructs an `Env<'_>` through those helpers, stays outside it. The
quarantined operations are the ones rsmpeg does not yet expose safely:
contains `unsafe`; everything else stays outside it. The quarantined
operations are the ones rsmpeg does not yet expose safely:

- clearing `AVCodecParameters.codec_tag` (a single primitive store on a
unique `&mut` borrow),
Expand All @@ -72,11 +71,7 @@ quarantined operations are the ones rsmpeg does not yet expose safely:
`AVFormatContextOutput.metadata` (libavformat takes ownership),
- comparing two raw `AVChannelLayout`s through
`av_channel_layout_compare`, which reads both and keeps no pointer, to
decide whether audio extraction can skip resampling,
- rebuilding an `Env<'_>` from the raw `NIF_ENV` captured at the entry
point, so a long-running operation can emit throttled
`{:exmpeg_progress, ...}` messages without an `OwnedEnv` (which panics
on dirty-scheduler threads).
decide whether audio extraction can skip resampling.

Every `unsafe` block names its invariant in a `SAFETY:` comment, and unit
tests in the same module exercise the round-trips.
Expand Down
18 changes: 6 additions & 12 deletions native/exmpeg_native/src/cancel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,6 @@ use std::time::{Duration, Instant};

use rustler::Env;
use rustler::types::LocalPid;
use rustler::wrapper::NIF_ENV;

use crate::errors::NativeError;

Expand All @@ -43,22 +42,18 @@ const CHECK_INTERVAL: Duration = Duration::from_millis(100);

/// Watches whether the calling BEAM process is still alive so a
/// long-running NIF can bail out when its caller dies.
pub(crate) struct CancelGuard {
pub(crate) struct CancelGuard<'a> {
pid: LocalPid,
/// Raw `NIF_ENV` pointer for the calling process, reconstructed into
/// an `Env<'_>` for each liveness check the same way `ProgressEmitter`
/// reconstructs it to send messages. Valid for the lifetime of the
/// NIF call that constructed this guard.
env_ptr: NIF_ENV,
env: Env<'a>,
last_check: Option<Instant>,
}

impl CancelGuard {
impl<'a> CancelGuard<'a> {
/// Capture the calling process from the entry-point env.
pub(crate) fn new(env: Env<'_>) -> Self {
pub(crate) fn new(env: Env<'a>) -> Self {
Self {
pid: env.pid(),
env_ptr: env.as_c_arg(),
env,
last_check: None,
}
}
Expand All @@ -77,8 +72,7 @@ impl CancelGuard {
}
self.last_check = Some(now);

let env = crate::ffi_helpers::reconstruct_env(self.env_ptr);
if env.is_process_alive(self.pid) {
if self.env.is_process_alive(self.pid) {
Ok(())
} else {
Err(NativeError::new(
Expand Down
2 changes: 1 addition & 1 deletion native/exmpeg_native/src/concat.rs
Original file line number Diff line number Diff line change
Expand Up @@ -172,7 +172,7 @@ fn process_input(
pts_offset: &[i64],
next_min_dts: &mut [i64],
packets_written: &mut u64,
cancel: &mut CancelGuard,
cancel: &mut CancelGuard<'_>,
) -> Result<(), NativeError> {
// Move the input to a zero origin before the cumulative offset, as
// `ffmpeg -f concat` does. An MPEG-TS capture or an MP4 with an
Expand Down
2 changes: 1 addition & 1 deletion native/exmpeg_native/src/extract_frame.rs
Original file line number Diff line number Diff line change
Expand Up @@ -237,7 +237,7 @@ fn decode_target_frame(
video_index: usize,
time_base: ffi::AVRational,
target_s: f64,
cancel: &mut CancelGuard,
cancel: &mut CancelGuard<'_>,
) -> Result<AVFrame, NativeError> {
let target_pts =
(target_s * f64::from(time_base.den) / f64::from(time_base.num)).round() as i64;
Expand Down
11 changes: 0 additions & 11 deletions native/exmpeg_native/src/ffi_helpers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,17 +20,6 @@ use rsmpeg::avutil::{AVAudioFifo, AVDictionary, AVFrame};
use rsmpeg::error::RsmpegError;
use rsmpeg::ffi;

/// Reconstruct an `Env` struct from a raw `NIF_ENV` pointer.
///
/// SAFETY: The raw `NIF_ENV` pointer must be a valid environment pointer handed to the NIF
/// by Rustler, and the returned `Env` must not outlive the NIF execution context.
pub(crate) fn reconstruct_env<'a>(c_env: rustler::wrapper::NIF_ENV) -> rustler::Env<'a> {
// SAFETY: We assume the caller provides a valid `c_env` pointer from the NIF context.
// Creating an Env is unsafe because it lets the caller create arbitrary lifetime references,
// which is sound as long as the returned Env does not escape the NIF call.
unsafe { rustler::Env::new(&(), c_env) }
}

/// Whether two channel layouts are semantically identical - the same
/// channels in the same positions (not just the same count). Returns
/// `false` on the rare comparison error so callers fall back to the safe
Expand Down
38 changes: 16 additions & 22 deletions native/exmpeg_native/src/progress.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,20 +7,14 @@
//! sent unconditionally via `finish` so the caller sees the closing
//! counts.
//!
//! Implementation note: `Env::send` from inside a NIF requires the
//! calling NIF's env (a process-bound, scheduler-thread-managed env).
//! `OwnedEnv::send_and_clear` panics on managed threads. We capture the
//! raw `NIF_ENV` pointer from the entry-point env at construction time
//! and reconstruct an `Env<'_>` for each emission. This is sound: the
//! pointer is valid for the duration of the NIF call (Rustler
//! guarantees this), and the emitter cannot outlive the call because
//! the NIF function consumes ownership of it before returning.
//! The emitter sends through the calling NIF's `Env<'a>`:
//! `OwnedEnv::send_and_clear` panics on BEAM-managed threads. The lifetime
//! keeps the emitter inside the NIF call.

use std::time::{Duration, Instant};

use rsmpeg::ffi;
use rustler::types::LocalPid;
use rustler::wrapper::NIF_ENV;
use rustler::{Encoder, Env, NifMap};

mod atoms {
Expand All @@ -47,34 +41,32 @@ pub(crate) struct ProgressUpdate {
pub(crate) total_duration_s: f64,
}

pub(crate) struct ProgressEmitter {
inner: Option<Inner>,
pub(crate) struct ProgressEmitter<'a> {
inner: Option<Inner<'a>>,
}

struct Inner {
struct Inner<'a> {
pid: LocalPid,
/// Raw `NIF_ENV` pointer for the calling process. Valid for the
/// lifetime of the NIF call that constructed this emitter.
env_ptr: NIF_ENV,
env: Env<'a>,
op: &'static str,
total_duration_s: f64,
last_emit: Option<Instant>,
}

impl ProgressEmitter {
impl<'a> ProgressEmitter<'a> {
/// Build an emitter. `env` is the calling NIF's environment (used
/// to send messages from the dirty-scheduler thread). `pid` is the
/// BEAM pid messages go to; if `None`, the emitter is a no-op.
pub(crate) fn new(
env: Env<'_>,
env: Env<'a>,
pid: Option<LocalPid>,
op: &'static str,
total_duration_s: f64,
) -> Self {
Self {
inner: pid.map(|pid| Inner {
pid,
env_ptr: env.as_c_arg(),
env,
op,
total_duration_s,
last_emit: None,
Expand All @@ -83,7 +75,7 @@ impl ProgressEmitter {
}

pub(crate) fn from_av_duration(
env: Env<'_>,
env: Env<'a>,
pid: Option<LocalPid>,
op: &'static str,
av_duration_ticks: i64,
Expand Down Expand Up @@ -122,15 +114,17 @@ impl ProgressEmitter {
}
}

fn send(inner: &Inner, packets_written: u64, current_pts_s: f64) {
fn send(inner: &Inner<'_>, packets_written: u64, current_pts_s: f64) {
let update = ProgressUpdate {
op: inner.op.to_owned(),
packets_written,
current_pts_s,
total_duration_s: inner.total_duration_s,
};
let env = crate::ffi_helpers::reconstruct_env(inner.env_ptr);
// Errors here are advisory: process gone, mailbox full, etc. Never
// block a transcode on a slow subscriber.
let _ = env.send(&inner.pid, (atoms::exmpeg_progress(), update).encode(env));
let _ = inner.env.send(
&inner.pid,
(atoms::exmpeg_progress(), update).encode(inner.env),
);
}
4 changes: 2 additions & 2 deletions native/exmpeg_native/src/transcode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -726,7 +726,7 @@ fn process_video_packet(
output: &mut AVFormatContextOutput,
out_time_bases: &[ffi::AVRational],
packets_written: &mut u64,
cancel: &mut CancelGuard,
cancel: &mut CancelGuard<'_>,
) -> Result<(), NativeError> {
let StreamPipeline::Video {
out_idx,
Expand Down Expand Up @@ -827,7 +827,7 @@ fn drain_filter_and_encode(
out_time_bases: &[ffi::AVRational],
last_out_pts: &mut i64,
packets_written: &mut u64,
cancel: &mut CancelGuard,
cancel: &mut CancelGuard<'_>,
) -> Result<(), NativeError> {
loop {
cancel.check()?;
Expand Down