mediaway 0.2.1

Convenience pipeline layer — composes encoder + container (+ device capture)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
//! [`EncodeSession`] — encoder + muxer composition (video, + optional audio track).

#![forbid(unsafe_code)]

use std::collections::VecDeque;

use crate::error::PipelineError;
use crate::filter::{FilterError, FrameFilter};
use mediaway_common::{AudioFrame, StreamInfo, VideoFrame, VideoFrameStorage};
use mediaway_container::{ContainerError, Mux, MuxOpen, mp4};
use mediaway_encoder::{AudioEncoder, VideoEncoder};
use mediaway_sw::apm::{AudioProcessor, VoiceActivityDetector};
use smallvec::SmallVec;

/// The optional audio side of an [`EncodeSession`] — present only when opened via
/// [`EncodeSession::open_with_audio`]. See
/// `adr/0003-audio-track-and-apm-integration.md`.
struct AudioTrack {
    encoder: Box<dyn AudioEncoder>,
    track_id: u32,
    /// AEC3 + NS + AGC2, if [`EncodeSession::attach_audio_processor`] was called.
    processor: Option<AudioProcessor>,
    /// RNN voice-activity detector, if [`EncodeSession::attach_vad`] was called.
    vad: Option<VoiceActivityDetector>,
    /// One score per processed 10ms block, drained by
    /// [`EncodeSession::poll_vad_score`].
    vad_scores: VecDeque<f32>,
}

/// Encode frames straight to container bytes.
///
/// Wraps one [`VideoEncoder`] + a muxer (single-track, or two-track when opened via
/// [`open_with_audio`](Self::open_with_audio)), draining `poll_packet` into the muxer on
/// every [`write_frame`](Self::write_frame)/[`write_audio_frame`](Self::write_audio_frame)
/// call instead of making callers write that loop themselves.
///
/// # Choosing a container
///
/// `M` is the muxer's track-registration phase and defaults to [`mp4::Muxer`], so
/// `EncodeSession::open(encoder)` still produces fragmented MP4 exactly as before. Name a
/// different [`MuxOpen`] to get a different container:
///
/// ```ignore
/// use mediaway_container::webm;
/// let session = EncodeSession::<_, webm::Muxer<webm::Open>>::open(encoder)?;
/// ```
///
/// Or hand over a pre-configured muxer with [`open_in`](Self::open_in), which is the only
/// way to reach options the muxer's own constructor exposes —
/// `mp4::Muxer::with_fragment_batch` was unreachable through this facade before.
///
/// Not every encoder pairs with every container: `add_track` is what rejects a codec the
/// container has no mapping for (`WebM` has no `CodecID` for H.264, for instance), so the
/// mismatch surfaces at [`open`](Self::open) rather than silently producing an unplayable
/// file.
///
/// Generic over `E` — works with a concrete unboxed encoder (e.g. Windows
/// `AutoVideoEncoder`) or `Box<dyn VideoEncoder>` (cross-platform dispatch via
/// [`crate::platform::AutoEncoder::open`]) without imposing a `Box` where the
/// caller doesn't already have one. The audio encoder is always `Box<dyn AudioEncoder>`
/// internally, regardless of `E` — see `adr/0003-audio-track-and-apm-integration.md`
/// § Struct shape for why this does not become a second generic parameter.
pub struct EncodeSession<E: VideoEncoder, M: MuxOpen = mp4::Muxer<mp4::Open>> {
    encoder: E,
    muxer: M::Live,
    track_id: u32,
    filters: SmallVec<[Box<dyn FrameFilter>; 4]>,
    audio: Option<AudioTrack>,
}

/// Fragmented-MP4 constructors.
///
/// These live on the concrete MP4 session rather than alongside
/// [`open_in`](EncodeSession::open_in) on the generic one, and that is deliberate: a
/// struct's default type parameter is *not* consulted when inferring an associated
/// function call, so a generic `EncodeSession::open` would force every existing caller to
/// write `EncodeSession::<_, mp4::Muxer<mp4::Open>>::open(..)`. Pinning `M` in the impl
/// header lets inference resolve it from the impl, so `EncodeSession::open(encoder)` keeps
/// meaning what it always meant.
///
/// Other containers go through [`open_in`](EncodeSession::open_in), which names the muxer
/// explicitly — `EncodeSession::open_in(webm::Muxer::new(), encoder)` reads better than the
/// turbofish would anyway.
impl<E: VideoEncoder> EncodeSession<E, mp4::Muxer<mp4::Open>> {
    /// Register `encoder`'s stream as an MP4 track and begin streaming. Video-only —
    /// see [`open_with_audio`](Self::open_with_audio) for a session with an audio track.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError`] if the muxer rejects the encoder's stream info.
    pub fn open(encoder: E) -> Result<Self, PipelineError> {
        Self::open_in(mp4::Muxer::new(), encoder)
    }

    /// Register `encoder`'s and `audio_encoder`'s streams as tracks (video first, then
    /// audio) and begin streaming.
    ///
    /// Both tracks must be known before this call: [`mp4::Muxer`] is typestate — tracks
    /// can only be added before `begin()`, so there is no way to add an audio track to a
    /// session already opened via [`open`](Self::open)
    /// (`adr/0003-audio-track-and-apm-integration.md` § Context).
    ///
    /// `add_track` requires unique track ids and rejects a duplicate — two independently
    /// constructed encoders both typically report `id: 0` by default (unlike
    /// [`open`](Self::open)'s single-track case, where that default is never a
    /// conflict). This renumbers explicitly, from
    /// [`MuxOpen::FIRST_TRACK_ID`] upward (video first, then audio — `0`/`1` on MP4,
    /// `1`/`2` on `WebM`, which reserves `0`), rather than trusting each encoder's own
    /// default — the same renumbering `tests/screen_mic_av_smoke.rs`
    /// used to do by hand via `StreamInfo::with_id` before migrating onto this
    /// constructor.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError`] if the muxer rejects either encoder's stream info.
    pub fn open_with_audio(
        encoder: E,
        audio_encoder: impl AudioEncoder + 'static,
    ) -> Result<Self, PipelineError> {
        Self::open_in_with_audio(mp4::Muxer::new(), encoder, audio_encoder)
    }
}

impl<E: VideoEncoder, M: MuxOpen> EncodeSession<E, M>
where
    M::Error: Into<ContainerError>,
{
    /// [`open`](Self::open), with a caller-supplied muxer.
    ///
    /// Takes the muxer in its registration phase and consumes it: the session owns the live
    /// muxer from here on. This is how a caller reaches muxer options this facade does not
    /// mirror — `mp4::Muxer::with_fragment_batch` being the motivating one.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError`] if the muxer rejects the encoder's stream info.
    pub fn open_in(mut muxer: M, encoder: E) -> Result<Self, PipelineError> {
        let info = encoder.stream_info().clone();
        // Raise the encoder's own track id only if the container forbids it. Encoders
        // default to `id: 0`, which ISOBMFF accepts and Matroska rejects outright, so
        // leaving it alone would make a WebM session fail on a stream the caller never
        // chose an id for. Clamping upward rather than overwriting keeps an encoder that
        // *did* pick an id in charge of it.
        let info = if info.id() < M::FIRST_TRACK_ID {
            info.with_id(M::FIRST_TRACK_ID)
        } else {
            info
        };
        let track_id = muxer.add_track(info).map_err(Into::into)?;
        Ok(Self {
            encoder,
            muxer: muxer.begin(),
            track_id,
            filters: SmallVec::new(),
            audio: None,
        })
    }

    /// [`open_with_audio`](Self::open_with_audio), with a caller-supplied muxer. See
    /// [`open_in`](Self::open_in).
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError`] if the muxer rejects either encoder's stream info.
    pub fn open_in_with_audio(
        mut muxer: M,
        encoder: E,
        audio_encoder: impl AudioEncoder + 'static,
    ) -> Result<Self, PipelineError> {
        let video_id = M::FIRST_TRACK_ID;
        let audio_id = M::FIRST_TRACK_ID + 1;
        let track_id = muxer
            .add_track(encoder.stream_info().clone().with_id(video_id))
            .map_err(Into::into)?;
        let audio_track_id = muxer
            .add_track(audio_encoder.stream_info().clone().with_id(audio_id))
            .map_err(Into::into)?;
        Ok(Self {
            encoder,
            muxer: muxer.begin(),
            track_id,
            filters: SmallVec::new(),
            audio: Some(AudioTrack {
                encoder: Box::new(audio_encoder),
                track_id: audio_track_id,
                processor: None,
                vad: None,
                vad_scores: VecDeque::new(),
            }),
        })
    }

    /// Attach AEC3 + NS + AGC2 audio enhancement — subsequent
    /// [`write_audio_frame`](Self::write_audio_frame) calls push through `processor`
    /// before reaching the audio encoder. Replaces a previously attached processor, if
    /// any.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError::NoAudioTrack`] if this session was opened via
    /// [`open`](Self::open) (no audio track to attach to) — use
    /// [`open_with_audio`](Self::open_with_audio) instead.
    pub fn attach_audio_processor(
        &mut self,
        processor: AudioProcessor,
    ) -> Result<&mut Self, PipelineError> {
        let audio = self.audio.as_mut().ok_or(PipelineError::NoAudioTrack)?;
        audio.processor = Some(processor);
        Ok(self)
    }

    /// Attach an RNN voice-activity detector — subsequent
    /// [`write_audio_frame`](Self::write_audio_frame) calls score each processed 10ms
    /// block, retrievable via [`poll_vad_score`](Self::poll_vad_score). Replaces a
    /// previously attached detector, if any.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError::NoAudioTrack`] if this session was opened via
    /// [`open`](Self::open) (no audio track to attach to) — use
    /// [`open_with_audio`](Self::open_with_audio) instead.
    pub fn attach_vad(&mut self, vad: VoiceActivityDetector) -> Result<&mut Self, PipelineError> {
        let audio = self.audio.as_mut().ok_or(PipelineError::NoAudioTrack)?;
        audio.vad = Some(vad);
        Ok(self)
    }

    /// Append a filter to the chain (runs after previously pushed filters).
    /// Filters may be pushed at any point before or between `write_frame` calls.
    pub fn push_filter<F: FrameFilter>(&mut self, filter: F) -> &mut Self {
        self.filters.push(Box::new(filter));
        self
    }

    /// Push one frame and drain any packets it produces into the muxer.
    ///
    /// Frames pass through the [`push_filter`](Self::push_filter) chain (if any)
    /// before reaching the encoder. An empty chain costs nothing beyond one
    /// `is_empty()` check. A non-empty chain rejects `Gpu`-backed frames with
    /// [`FilterError::GpuFrameUnsupported`] — v1 filters are CPU-frame-only
    /// (see [ADR-0001](../../adr/0001-frame-filter-hook.md)).
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError`] on filter, encoder, or mux failure.
    pub fn write_frame(&mut self, frame: &VideoFrame) -> Result<(), PipelineError> {
        if self.filters.is_empty() {
            self.encoder.push_frame(frame)?; // unchanged fast path, zero clone
        } else {
            if matches!(frame.storage, VideoFrameStorage::Gpu(_)) {
                return Err(PipelineError::Filter(FilterError::GpuFrameUnsupported));
            }
            // clone: entry point into an owned filter chain — the caller only lent a
            // reference, but VideoFrame::clone() is a Bytes refcount bump (Cpu) or a
            // Copy of a small handle (Gpu, unreachable here), never a pixel memcpy.
            // Paid exactly once per frame, only when a filter chain is attached.
            let mut current = frame.clone();
            for filter in &mut self.filters {
                current = filter.process(current)?;
            }
            self.encoder.push_frame(&current)?;
        }
        self.drain()
    }

    /// Retarget the live CBR bitrate ceiling on the underlying encoder — see
    /// [`mediaway_encoder::VideoEncoder::set_bitrate`]. No session reopen, no dropped
    /// frames; takes effect from the next [`write_frame`](Self::write_frame) call.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError::Encode`] if the underlying encoder was not opened in
    /// CBR mode or cannot retarget bitrate live (`EncodeError::Unsupported`), or on a
    /// backend failure.
    pub fn set_bitrate(&mut self, bitrate_bps: u32) -> Result<(), PipelineError> {
        self.encoder.set_bitrate(bitrate_bps)?;
        Ok(())
    }

    /// Push one microphone/capture-side audio frame and drain any packets it produces
    /// into the muxer.
    ///
    /// With no [`attach_audio_processor`](Self::attach_audio_processor) call, `frame`
    /// goes straight to the audio encoder (the fast path `tests/screen_mic_av_smoke.rs`
    /// now exercises via this method). With a processor attached, `frame` is pushed into it and every
    /// resulting processed 10ms block (zero, one, or several — `AudioProcessor`
    /// re-blocks internally) is scored by [`attach_vad`](Self::attach_vad)'s detector,
    /// if any, then pushed to the audio encoder. See
    /// `adr/0003-audio-track-and-apm-integration.md` § The write path.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError::NoAudioTrack`] if this session was opened via
    /// [`open`](Self::open). Returns [`PipelineError::Apm`] on a transient
    /// `AudioProcessor` failure (the instance keeps working in degraded/passthrough
    /// mode afterward — see that error variant's docs). Otherwise returns
    /// [`PipelineError`] on encoder or mux failure. A [`VoiceActivityDetector`] failure
    /// is **not** propagated here — see [`poll_vad_score`](Self::poll_vad_score).
    pub fn write_audio_frame(&mut self, frame: &AudioFrame) -> Result<(), PipelineError> {
        let Self { audio, muxer, .. } = self;
        let Some(audio) = audio.as_mut() else {
            return Err(PipelineError::NoAudioTrack);
        };

        if let Some(processor) = audio.processor.as_mut() {
            processor.push_capture_frame(frame)?;
            while let Some(block) = processor.poll_processed_frame()? {
                if let Some(vad) = audio.vad.as_mut() {
                    // A disabled VAD's `analyze` errors forever (no honest scalar
                    // passthrough) — that must not block audio encoding, only stop
                    // producing new scores. See `adr/0003-audio-track-and-apm-integration.md`
                    // § Error handling.
                    if let Ok(score) = vad.analyze(&block) {
                        audio.vad_scores.push_back(score);
                    }
                }
                audio.encoder.push_frame(&block)?;
            }
        } else {
            audio.encoder.push_frame(frame)?;
        }
        Self::drain_audio(audio, muxer)
    }

    /// Feed a render-reference (far-end / about-to-be-played) frame to the attached
    /// [`AudioProcessor`], if any — the echo-cancellation reference signal
    /// (`AudioProcessor::push_render_frame`). Only meaningful for a caller that is also
    /// playing audio back (e.g. voice chat); a pure recorder never needs this.
    ///
    /// A no-op when this session has an audio track but no processor attached — there
    /// is nothing to feed the reference into.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError::NoAudioTrack`] if this session was opened via
    /// [`open`](Self::open). Returns [`PipelineError::Apm`] on a transient
    /// `AudioProcessor` failure.
    pub fn write_audio_render_frame(&mut self, frame: &AudioFrame) -> Result<(), PipelineError> {
        let Some(audio) = self.audio.as_mut() else {
            return Err(PipelineError::NoAudioTrack);
        };
        if let Some(processor) = audio.processor.as_mut() {
            processor.push_render_frame(frame)?;
        }
        Ok(())
    }

    /// Pop the next voice-activity score produced by
    /// [`write_audio_frame`](Self::write_audio_frame), if any — one score per processed
    /// 10ms block, in production order. `None` when no [`attach_vad`](Self::attach_vad)
    /// detector is attached, or none is ready yet.
    ///
    /// If the attached [`VoiceActivityDetector`] becomes disabled (a caught backend
    /// panic), this silently stops producing new scores rather than surfacing an error
    /// here — see `adr/0003-audio-track-and-apm-integration.md` § Error handling.
    pub fn poll_vad_score(&mut self) -> Option<f32> {
        self.audio.as_mut()?.vad_scores.pop_front()
    }

    /// Append whatever fMP4 bytes are ready to `out`, returning how many were written.
    ///
    /// The streaming exit. Call it as often as you like during a session — after every
    /// [`write_frame`](Self::write_frame), on a timer, or not at all — and the session's
    /// memory stays bounded by the poll cadence instead of growing with the recording's
    /// length. A one-hour capture that is never polled holds the entire file in RAM until
    /// [`finish`](Self::finish); the same capture polled into a `File` holds a fragment.
    ///
    /// Returns `0` when nothing is ready. That is indistinguishable from "already fully
    /// drained" — a caller that needs the difference must track its own running total
    /// (`adr/0006-encode-session-streaming-bytes.md` § Negative).
    ///
    /// Bytes are appended, never overwritten, so one buffer can be reused across the
    /// whole session: poll, write it out, `clear()`, repeat.
    ///
    /// Polling alone never finishes the stream — the last fragments only appear after
    /// [`finish_into`](Self::finish_into) or [`finish`](Self::finish) flushes the
    /// encoders and muxer.
    pub fn poll_bytes(&mut self, out: &mut Vec<u8>) -> usize {
        self.muxer.poll_bytes(out)
    }

    /// Flush the encoder(s) and muxer, then append the remaining fMP4 bytes to `out` and
    /// return how many were written.
    ///
    /// **This returns what has not been polled yet, not the whole stream.** A session
    /// that was drained with [`poll_bytes`](Self::poll_bytes) along the way ends here
    /// with only its tail; a session that was never polled ends with the complete
    /// recording. Consuming `self` is what makes "write a frame after flushing"
    /// unrepresentable rather than merely wrong
    /// (`adr/0006-encode-session-streaming-bytes.md`).
    ///
    /// **Known gap, inherited from `mediaway-audio-apm`, not fixed here**: a trailing
    /// audio block shorter than 10ms sitting in an attached [`AudioProcessor`]'s
    /// internal buffer is not flushed — `AudioProcessor` has no "flush a partial block"
    /// method today. See `adr/0003-audio-track-and-apm-integration.md` § `finish()`.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError`] on encoder or mux failure.
    pub fn finish_into(mut self, out: &mut Vec<u8>) -> Result<usize, PipelineError> {
        self.encoder.flush()?;
        self.drain()?;
        if let Some(mut audio) = self.audio.take() {
            audio.encoder.flush()?;
            Self::drain_audio(&mut audio, &mut self.muxer)?;
        }
        self.muxer.flush();
        Ok(self.muxer.poll_bytes(out))
    }

    /// Flush the encoder(s) and muxer, returning the fMP4 bytes as one buffer.
    ///
    /// The whole-buffer convenience over [`finish_into`](Self::finish_into), which is the
    /// streaming exit. For a session that was never polled this is the complete
    /// recording — the original and still most common use. For a session drained with
    /// [`poll_bytes`](Self::poll_bytes) it is only the tail; prefer `finish_into` there,
    /// so the tail lands in the same buffer as everything else rather than in a fresh
    /// allocation.
    ///
    /// Carries [`finish_into`](Self::finish_into)'s trailing-partial-audio-block gap.
    ///
    /// # Errors
    ///
    /// Returns [`PipelineError`] on encoder or mux failure.
    pub fn finish(self) -> Result<Vec<u8>, PipelineError> {
        let mut bytes = Vec::new();
        self.finish_into(&mut bytes)?;
        Ok(bytes)
    }

    fn drain(&mut self) -> Result<(), PipelineError> {
        while let Some(mut pkt) = self.encoder.poll_packet()? {
            // A backend whose config record is only known after encoding at least one
            // frame (e.g. `VideoToolbox`, which determines SPS/PPS internally rather
            // than deriving them from open-time config) reports empty `extra_data` from
            // `stream_info()` at `open()` time — by the time `poll_packet` returns a
            // packet, the backend has necessarily finished encoding it, so `stream_info()`
            // now reflects the real, finalized value. A no-op once the muxer already has
            // real `extra_data` (from this call or from `push_packet`'s own in-band
            // Annex-B extraction) or once the moov header is already written.
            if let StreamInfo::Video { extra_data, .. } = self.encoder.stream_info()
                && !extra_data.is_empty()
            {
                // clone: Bytes refcount bump, not a payload copy — set_track_extra_data
                // needs an owned value but this fires at most once in practice (a no-op
                // once the muxer's own extra_data is no longer empty).
                self.muxer
                    .set_track_extra_data(self.track_id, extra_data.clone());
            }
            pkt.stream_id = self.track_id;
            self.muxer.push_packet(&pkt).map_err(Into::into)?;
        }
        Ok(())
    }

    fn drain_audio(audio: &mut AudioTrack, muxer: &mut M::Live) -> Result<(), PipelineError> {
        while let Some(mut pkt) = audio.encoder.poll_packet()? {
            pkt.stream_id = audio.track_id;
            muxer.push_packet(&pkt).map_err(Into::into)?;
        }
        Ok(())
    }
}

#[cfg(test)]
#[path = "session_tests.rs"]
mod tests;