Skip to content
Draft
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
139 changes: 101 additions & 38 deletions InstantReplay.Externals/unienc/crates/unienc_apple_vt/src/audio/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ pub struct AudioToolboxEncoder {
}

pub struct AudioToolboxEncoderInput {
tx: mpsc::Sender<AudioPacket>,
tx: mpsc::UnboundedSender<AudioPacket>,
converter: AudioConverter,
max_output_packet_size: u32,
sample_rate: u32,
Expand All @@ -42,7 +42,7 @@ pub struct AudioToolboxEncoderInput {
unsafe impl Send for AudioToolboxEncoderInput {}

pub struct AudioToolboxEncoderOutput {
rx: mpsc::Receiver<AudioPacket>,
rx: mpsc::UnboundedReceiver<AudioPacket>,
}

#[derive(Encode, Decode, Clone)]
Expand Down Expand Up @@ -105,45 +105,13 @@ impl EncoderInput for AudioToolboxEncoderInput {
let mut sample = Some(&data);

while {
let num_output_packets = self.converter.fill_complex_buffer(
self.encode_round(
&mut sample,
false,
&mut output_buffer_data,
&mut packet_descs,
)?;

let magic_cookie = self
.converter
.get_property_raw(kAudioConverterCompressionMagicCookie)?;

let packet_descs = &packet_descs[..num_output_packets as usize];

for packet_desc in packet_descs {
// AudioToolbox can emit multiple output packets from a single input buffer (e.g. when a
// large batched buffer is pushed while the main thread is stalled by framerate jitter).
// Assigning every packet the input buffer's timestamp would produce duplicate / overlapping
// PTS that the muxer (AVAssetWriter) rejects. The encoder consumes a continuous PCM stream,
// so assign a contiguous per-packet timeline advancing by exactly one packet for each
// emitted packet. The running position is seeded to the first input timestamp at the start
// of `push`; forward discontinuities are materialized there as leading silence, so the
// position simply advances one packet at a time here.
let timestamp_in_samples = self
.output_position_in_samples
.expect("output_position_in_samples is initialized at the start of push");

let packet = AudioPacket {
data: output_buffer_data[packet_desc.mStartOffset as usize
..packet_desc.mStartOffset as usize + packet_desc.mDataByteSize as usize]
.to_vec(),
timestamp_in_samples,
sample_rate: self.sample_rate,
magic_cookie: magic_cookie.clone(),
};
self.tx.send(packet).await.map_err(AppleError::from)?;

self.output_position_in_samples =
Some(timestamp_in_samples + self.frames_per_packet);
}

sample.is_some()
} {}

Expand All @@ -153,6 +121,92 @@ impl EncoderInput for AudioToolboxEncoderInput {
}
}

impl AudioToolboxEncoderInput {
/// Runs a single conversion round and forwards every produced packet to the output channel,
/// advancing the contiguous per-packet output timeline. Returns the number of packets emitted.
fn encode_round(
&mut self,
sample: &mut Option<&AudioSample>,
end_of_stream: bool,
output_buffer_data: &mut [u8],
packet_descs: &mut [AudioStreamPacketDescription],
) -> Result<u32> {
let num_output_packets = self.converter.fill_complex_buffer(
sample,
end_of_stream,
output_buffer_data,
packet_descs,
)?;

let magic_cookie = self
.converter
.get_property_raw(kAudioConverterCompressionMagicCookie)?;

let packet_descs = &packet_descs[..num_output_packets as usize];

for packet_desc in packet_descs {
// AudioToolbox can emit multiple output packets from a single input buffer (e.g. when a
// large batched buffer is pushed while the main thread is stalled by framerate jitter).
// Assigning every packet the input buffer's timestamp would produce duplicate / overlapping
// PTS that the muxer (AVAssetWriter) rejects. The encoder consumes a continuous PCM stream,
// so assign a contiguous per-packet timeline advancing by exactly one packet for each
// emitted packet. The running position is seeded to the first input timestamp at the start
// of `push`; forward discontinuities are materialized there as leading silence, so the
// position simply advances one packet at a time here. The converter cannot produce packets
// before the first push seeds the position, so the fallback value is unreachable; it only
// exists to keep the end-of-stream drain (which runs in `drop`) panic-free.
let timestamp_in_samples = self.output_position_in_samples.unwrap_or_default();

let packet = AudioPacket {
data: output_buffer_data[packet_desc.mStartOffset as usize
..packet_desc.mStartOffset as usize + packet_desc.mDataByteSize as usize]
.to_vec(),
timestamp_in_samples,
sample_rate: self.sample_rate,
magic_cookie: magic_cookie.clone(),
};
self.tx.send(packet).map_err(AppleError::from)?;

self.output_position_in_samples = Some(timestamp_in_samples + self.frames_per_packet);
}

Ok(num_output_packets)
}

/// Drains packets still buffered inside the AudioConverter by signaling end of stream: complete
/// packets that were not emitted during `push` and the final partial (< frames-per-packet)
/// packet, which the converter pads and flushes. Without this, the tail of every recording
/// (up to a few tens of milliseconds of audio) is lost.
fn drain(&mut self) -> Result<()> {
let mut output_buffer_data = vec![0; self.max_output_packet_size as usize];
let mut packet_descs = vec![unsafe { std::mem::zeroed::<AudioStreamPacketDescription>() }];

loop {
let mut sample = None;
let num_output_packets = self.encode_round(
&mut sample,
true,
&mut output_buffer_data,
&mut packet_descs,
)?;
if num_output_packets == 0 {
break;
}
}

Ok(())
}
}

impl Drop for AudioToolboxEncoderInput {
fn drop(&mut self) {
// The input being dropped is the end-of-stream signal (`EncoderInput` has no explicit
// finish), so flush the converter here. Errors are ignored: the output side may already
// have been dropped, and panicking in drop is not an option.
_ = self.drain();
}
}

impl AudioToolboxEncoder {
pub fn new(options: &impl unienc_common::AudioEncoderOptions) -> Result<Self> {
let mut from = AudioStreamBasicDescription {
Expand Down Expand Up @@ -189,7 +243,7 @@ impl AudioToolboxEncoder {
&mut max_output_packet_size,
)?;

let (tx, rx) = mpsc::channel(32);
let (tx, rx) = mpsc::unbounded_channel();

Ok(Self {
input: AudioToolboxEncoderInput {
Expand Down Expand Up @@ -324,6 +378,7 @@ impl AudioConverter {
fn fill_complex_buffer(
&self,
sample: &mut Option<&AudioSample>,
end_of_stream: bool,
output_buffer: &mut [u8],
packet_descs: &mut [AudioStreamPacketDescription],
) -> Result<u32> {
Expand All @@ -348,7 +403,15 @@ impl AudioConverter {

let Some(sample) = sample.take() else {
*number_data_packets = 0;
return FillBufferResult::Skip as i32;
// Returning noErr with zero packets tells the converter the input stream has
// ended, making it flush buffered output (padding the final partial packet).
// Returning a non-zero code instead just pauses conversion until more input
// arrives on a later call.
return if end_of_stream {
FillBufferResult::Ok as i32
} else {
FillBufferResult::Skip as i32
};
};

let num_input_packets = sample.data.len() * 2 / self.from.mBytesPerPacket as usize;
Expand Down
111 changes: 85 additions & 26 deletions InstantReplay.Externals/unienc/crates/unienc_apple_vt/src/video/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,13 +24,23 @@ use unienc_common::{
use crate::{MetalTexture, common::UnsafeSendRetained, metal};
use unienc_common::TryFromUnityNativeTexturePointer;

/// Capacity of the bounded channel carrying encoded frames from the VideoToolbox callback to the
/// consumer. `push` reserves a slot before submitting each frame, so when the consumer falls behind
/// and the channel fills up, `push` suspends (backpressure) instead of the callback dropping
/// already-encoded frames (which would break the P-frame reference chain).
const OUTPUT_CHANNEL_CAPACITY: usize = 32;

/// A reserved output-channel slot handed to a single encode callback via the per-frame source-frame
/// ref-con. Completing the send through it can neither block the callback nor drop the frame.
type OutputPermit = mpsc::OwnedPermit<Result<VideoEncodedData>>;

pub struct VideoToolboxEncoder {
input: VideoToolboxEncoderInput,
output: VideoToolboxEncoderOutput,
}
pub struct VideoToolboxEncoderInput {
session: CompressionSession,
tx: Box<mpsc::Sender<VideoEncodedData>>,
tx: Box<mpsc::Sender<Result<VideoEncodedData>>>,
width: u32,
height: u32,
bitrate: u32,
Expand All @@ -43,7 +53,7 @@ struct CompressionSession {
unsafe impl Send for VideoToolboxEncoderInput {}

pub struct VideoToolboxEncoderOutput {
rx: mpsc::Receiver<VideoEncodedData>,
rx: mpsc::Receiver<Result<VideoEncodedData>>,
}

pub struct VideoEncodedData {
Expand Down Expand Up @@ -105,17 +115,33 @@ impl EncodedData for VideoEncodedData {
}

unsafe extern "C-unwind" fn handle_video_encode_output(
output_callback_ref_con: *mut c_void,
_source_frame_ref_con: *mut c_void,
_status: i32,
_output_callback_ref_con: *mut c_void,
source_frame_ref_con: *mut c_void,
status: i32,
_info_flags: VTEncodeInfoFlags,
sample_buffer: *mut CMSampleBuffer,
) {
let tx = unsafe { &*(output_callback_ref_con as *const mpsc::Sender<VideoEncodedData>) };
// Each submitted frame carries a pre-reserved output-channel slot (an OwnedPermit) as its
// source-frame ref-con. Completing the send through that permit means this callback never
// blocks (it runs on VideoToolbox's internal queue, which complete_frames() waits on) and
// never drops an encoded frame (dropping would break the reference chain of subsequent
// P-frames). Backpressure is instead applied in `push`, which awaits the reservation before
// submitting the frame, so a full channel throttles the producer rather than corrupting output.
let permit = unsafe { *Box::from_raw(source_frame_ref_con as *mut OutputPermit) };

// Propagate encode failures to the output side instead of silently ignoring them.
if let Err(err) = status.to_result() {
permit.send(Err(err));
return;
}

if let Some(sample_buffer) = unsafe { Retained::retain(sample_buffer) } {
_ = tx.try_send(VideoEncodedData::new(sample_buffer.into()));
} // otherwise dropped
permit.send(Ok(VideoEncodedData::new(sample_buffer.into())));
} else {
// The encoder itself dropped the frame (e.g. kVTEncodeInfo_FrameDropped); dropping the
// permit returns the reserved slot to the channel.
drop(permit);
}
}

unsafe extern "C-unwind" fn release_pixel_buffer(
Expand All @@ -141,6 +167,19 @@ impl EncoderInput for VideoToolboxEncoderInput {
type Data = VideoSample<MetalTexture>;

async fn push(&mut self, data: Self::Data) -> unienc_common::Result<()> {
// Reserve an output-channel slot before doing any encode work. When the consumer is behind
// and the channel is full, this await suspends `push`, propagating backpressure up the
// pipeline so the C# side drops *input* frames (via DroppingChannelInput) instead of us
// dropping already-encoded output. Holding the OwnedPermit (Send) across the blit await
// below is fine; if building the frame fails, dropping the permit releases the slot.
let permit = self
.tx
.as_ref()
.clone()
.reserve_owned()
.await
.map_err(AppleError::from)?;

let buffer = match data.frame {
unienc_common::VideoFrame::Bgra32(bgra32) => {
let buffer = bgra32.buffer;
Expand Down Expand Up @@ -218,30 +257,49 @@ impl EncoderInput for VideoToolboxEncoderInput {
}
};

// Hand the reserved slot to this frame's encode callback via the source-frame ref-con. No
// await happens past this point, so the raw pointer is never held across a suspension.
let permit_ptr = Box::into_raw(Box::new(permit));

let mut retry = 0;

loop {
let res = loop {
let res = unsafe {
self.session.inner.encode_frame(
&buffer,
CMTime::with_seconds(data.timestamp, 720),
kCMTimeInvalid,
None,
std::ptr::null_mut(),
permit_ptr as *mut c_void,
std::ptr::null_mut(),
)
};

if res == kVTInvalidSessionErr && retry == 0 {
// VTCompressionSession turns invalid when the app enters background on iOS
// retrying once
// VTCompressionSession turns invalid when the app enters background on iOS; retry
// once with a fresh session, re-passing the same reserved permit.
retry += 1;
self.session =
CompressionSession::new(self.width, self.height, self.bitrate, &*self.tx)?;
continue;
match CompressionSession::new(self.width, self.height, self.bitrate) {
Ok(session) => {
self.session = session;
continue;
}
Err(err) => {
// No callback will run for this frame; reclaim the permit to free the slot.
drop(unsafe { Box::from_raw(permit_ptr) });
return Err(err.into());
}
}
}

break res.to_result()?;
break res;
};

if let Err(err) = res.to_result() {
// The frame was not accepted, so the callback will never consume the permit; reclaim it
// to release the reserved slot.
drop(unsafe { Box::from_raw(permit_ptr) });
return Err(err.into());
}

Ok(())
Expand All @@ -252,7 +310,11 @@ impl EncoderOutput for VideoToolboxEncoderOutput {
type Data = VideoEncodedData;

async fn pull(&mut self) -> unienc_common::Result<Option<Self::Data>> {
Ok(self.rx.recv().await)
match self.rx.recv().await {
Some(Ok(data)) => Ok(Some(data)),
Some(Err(err)) => Err(err.into()),
None => Ok(None),
}
}
}

Expand All @@ -273,12 +335,7 @@ impl Drop for VideoToolboxEncoderInput {
}

impl CompressionSession {
fn new(
width: u32,
height: u32,
bitrate: u32,
tx: *const mpsc::Sender<VideoEncodedData>,
) -> Result<Self> {
fn new(width: u32, height: u32, bitrate: u32) -> Result<Self> {
let mut session: *mut VTCompressionSession = std::ptr::null_mut();

unsafe {
Expand All @@ -291,7 +348,9 @@ impl CompressionSession {
None,
None,
Some(handle_video_encode_output),
tx as *mut c_void,
// The output callback receives its channel slot per-frame via the source-frame
// ref-con (see `push`), so no session-level ref-con is needed.
std::ptr::null_mut(),
NonNull::new(&mut session).ok_or(AppleError::NonNullCreationFailed)?,
)
.to_result()?;
Expand Down Expand Up @@ -330,14 +389,14 @@ impl CompressionSession {

impl VideoToolboxEncoder {
pub fn new(options: &impl unienc_common::VideoEncoderOptions) -> Result<Self> {
let (tx, rx) = mpsc::channel(32);
let (tx, rx) = mpsc::channel(OUTPUT_CHANNEL_CAPACITY);
let tx = Box::new(tx);

let (width, height, bitrate) = (options.width(), options.height(), options.bitrate());

Ok(VideoToolboxEncoder {
input: VideoToolboxEncoderInput {
session: CompressionSession::new(width, height, bitrate, &*tx)?,
session: CompressionSession::new(width, height, bitrate)?,
tx,
width,
height,
Expand Down