Files
punktfunk/crates/punktfunk-core/src/packet/packetize.rs
T
enricobuehler c5402cb1f7 feat(core+host): LN1 phase-2 — VIDEO_CAP_STREAMED_AU streamed access units
The wire half of sub-frame slice output (latency §7 LN1, planning
design/nvenc-subframe-slice-output.md Phase 2): toward a client that
advertises the new cap bit (0x20), a chunked-poll encoder session ships each
AU's completed FEC blocks while the tail of the frame is still encoding —
the AU's last packet leaves the host as the encode finishes instead of
after it.

Wire semantics (negotiated; zero change for anyone else):
- Non-final blocks ride SENTINEL headers: block_count = 0 (a value no legacy
  sender emits) + frame_bytes = 0 + exactly max_data_per_block data shards,
  so the receiver's shard-offset formula needs no total.
- The final block's headers carry the real frame_bytes/block_count (+
  FLAG_EOF) and RETRO-VALIDATE the whole frame: totals under which a
  received sentinel block is out of range or not full-K kill the frame
  wholesale (no spliced delivery) and the index can't be resurrected.
- Firewall: sentinels are bounded by the negotiated limits (full-K exactly,
  never the last block the limits allow, no total to lie about); the exact
  derived-geometry check runs unchanged on every non-sentinel packet and
  retroactively at pinning. Sentinel opens commit a max_frame_bytes buffer,
  bounded by the existing IN_FLIGHT_BUF_FACTOR budget (amplification test).
- Order-agnostic like legacy: a reversed frame (final block first) opens
  legacy-shaped and still accepts its sentinels against the pinned totals.
- Small/empty streamed AUs degenerate to byte-identical legacy headers.

Host: Packetizer::{begin,push,finish}_streamed seal full-K blocks (data +
parity per block) as chunks arrive; Session::seal_streamed_* share the
pooled-wire + two-lane seal machinery via the new seal_run; the send thread
paces each flush under the frame's existing deadline (pace_sealed split out
of paced_submit) and runs the whole-AU accounting at the last chunk; the
encode pump forwards poll_chunk output as ChunkMsg when the client has the
cap AND the encoder chunks (re-queried per AU — an escalation falls back
seamlessly). Probes never run mid-AU. PUNKTFUNK_STREAMED_AU=0 = host escape
hatch. Client core ORs the cap into Hello (the shared reassembler carries
the support). Sampled first_slice_us vs encode_us PERF log measures the
overlap; the 0xCF stage-field extension stays a follow-up.

Core tests: streamed round-trip (clean/loss/reorder/duplicate, both orders),
sentinel firewall bounds, lying-final wholesale kill + no-resurrect,
open-amplification budget, header-shape pins. Gates still owed before
default-on: security review pass, loss-harness curve, GameStream smoke
(plane untouched structurally), bitrate A/B.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-20 22:53:53 +02:00

519 lines
22 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! Host side: split an access unit into FEC-protected shard packets.
use super::*;
use crate::config::Config;
use crate::error::{PunktfunkError, Result};
use crate::fec::ErasureCoder;
use zerocopy::IntoBytes;
// ---------------------------------------------------------------------------
// Host side: packetization
// ---------------------------------------------------------------------------
/// Splits encoded access units into FEC-protected shard packets. Host-side only.
///
/// Frame numbering: a caller can pass an **explicit** `frame_index` to
/// [`packetize_each`](Self::packetize_each) (the punktfunk/1 encode loop owns the video numbering
/// so the encoder's reference-frame-invalidation bookkeeping stays 1:1 with the wire across
/// encoder rebuilds/resets), or pass `None` to draw from the internal counter (the legacy path —
/// synthetic/spike/ABI sessions where no encoder cares). Speed-test probe filler draws from a
/// **separate** index space ([`alloc_probe_index`](Self::alloc_probe_index)) so a burst never
/// consumes video indexes — see [`crate::quic::VIDEO_CAP_PROBE_SEQ`].
pub struct Packetizer {
next_frame_index: u32,
/// Probe-space frame counter (see [`alloc_probe_index`](Self::alloc_probe_index)).
next_probe_index: u32,
next_seq: u32,
shard_payload: usize,
fec: crate::config::FecConfig,
version: u8,
/// Reusable zero-padded scratch for the frame's final data shard when the frame isn't an
/// exact `shard_payload` multiple (and for the single all-zero shard of an empty frame).
/// Every other data shard is a `shard_payload`-sized slice straight into the frame buffer —
/// blocks are consecutive shard ranges, so only the frame's last shard can be partial.
tail: Vec<u8>,
/// Reusable PER-BLOCK parity buffers for [`ErasureCoder::encode_into`] (plan Phase 1.4):
/// one pool per block index, each growing once to its high-water recovery count, then
/// every frame's parity is generated into them with zero allocations. Per-block (not one
/// shared pool) because the data-first wire order (latency plan T1.3) emits every block's
/// DATA shards before any block's parity — all blocks' parity must stay alive until the
/// frame's second emission pass.
recovery: Vec<Vec<Vec<u8>>>,
/// The peer's per-block `data + recovery` acceptance ceiling, frozen from the **negotiated**
/// config exactly as the far side derives it in [`ReassemblerLimits::from_config`]. Adaptive
/// FEC moves `fec.fec_percent` live ([`set_fec_percent`](Self::set_fec_percent)) but the
/// receiver's ceiling is computed once at session construction and never re-derived, so parity
/// must be clamped against this or a raised percentage puts blocks over the far side's bound —
/// where every packet of the block is dropped wholesale, the frame never completes, and the
/// resulting loss pushes adaptive FEC *higher*. See the `recovery_for` clamp in `packetize_each`.
max_total_shards: usize,
/// The peer's per-frame block ceiling, mirroring [`ReassemblerLimits::from_config`]'s
/// `max_blocks` — the streamed path's bound on how many sentinel blocks it may emit (a
/// streamed AU's size isn't known up front, so this is the only pre-emission guard against
/// producing a frame the receiver must reject).
max_blocks: usize,
}
/// One in-progress **streamed** access unit (design/nvenc-subframe-slice-output.md Phase 2):
/// the caller feeds encoder chunks in as they exist ([`Packetizer::push_streamed`]) and every
/// completed `max_data_per_block × shard_payload` block leaves under SENTINEL headers
/// (`frame_bytes = 0`, `block_count = 0` — "not final yet") before the AU's total size is
/// known; [`Packetizer::finish_streamed`] seals the tail block with the real totals and
/// `FLAG_EOF`. Only ever sent to a peer that advertised
/// [`crate::quic::VIDEO_CAP_STREAMED_AU`].
pub struct StreamedAu {
frame_index: u32,
pts_ns: u64,
user_flags: u32,
/// Bytes not yet sealed into a block. Kept ≤ one block by `push_streamed` (it flushes only
/// while STRICTLY more than a block is buffered, so the final block always has ≥ 1 byte and
/// a sentinel block is never retroactively the frame's last).
pending: Vec<u8>,
/// Sentinel blocks already emitted.
blocks_out: u16,
/// Total AU bytes accumulated so far (`pending` included).
total_bytes: u64,
/// The frame's first packet (block 0, shard 0 — carries `FLAG_SOF`) has been emitted.
opened: bool,
}
impl StreamedAu {
/// The wire frame index this AU is sealed with (the caller's RFI bookkeeping domain).
pub fn frame_index(&self) -> u32 {
self.frame_index
}
}
impl Packetizer {
pub fn new(config: &Config) -> Self {
let max_data = config.fec.max_data_per_block as usize;
let total_data_max = config
.max_frame_bytes
.div_ceil(config.shard_payload.max(1))
.max(1);
Packetizer {
next_frame_index: 0,
next_probe_index: 0,
next_seq: 0,
shard_payload: config.shard_payload,
fec: config.fec,
version: config.phase as u8,
tail: Vec::new(),
recovery: Vec::new(),
// Mirrors `ReassemblerLimits::from_config` — keep the two in step.
max_total_shards: (max_data + config.fec.recovery_for(max_data))
.min(config.fec.scheme.max_total_shards()),
max_blocks: total_data_max.div_ceil(max_data).max(1),
}
}
/// Allocate the next **probe-space** frame index (speed-test filler). A separate counter from
/// the video `frame_index`es so a multi-thousand-AU probe burst never advances the video
/// numbering — the client routes [`FLAG_PROBE`]-flagged shards into its own reassembly window
/// (see [`Reassembler`]), so the two spaces never collide. Only used against clients that
/// advertise [`crate::quic::VIDEO_CAP_PROBE_SEQ`].
pub fn alloc_probe_index(&mut self) -> u32 {
let i = self.next_probe_index;
self.next_probe_index = i.wrapping_add(1);
i
}
/// Live-adjust the FEC recovery percentage (adaptive FEC). Takes effect on the next
/// [`packetize`](Self::packetize); the wire is self-describing (each packet carries its block's
/// data/recovery counts), so the receiver needs no notification. Clamped to ≤ 90.
pub fn set_fec_percent(&mut self, pct: u8) {
self.fec.fec_percent = pct.min(90);
}
/// The current FEC recovery percentage.
pub fn fec_percent(&self) -> u8 {
self.fec.fec_percent
}
/// Packetize one access unit into owned wire packets (header ++ shard payload each).
/// Thin wrapper over [`packetize_each`](Self::packetize_each) — the allocation-free
/// streaming path's reference implementation (tests and the loss harness use this).
pub fn packetize(
&mut self,
frame: &[u8],
pts_ns: u64,
user_flags: u32,
coder: &dyn ErasureCoder,
) -> Result<Vec<Vec<u8>>> {
let mut packets = Vec::new();
self.packetize_each(frame, pts_ns, user_flags, None, coder, |hdr, body| {
let mut pkt = Vec::with_capacity(HEADER_LEN + body.len());
pkt.extend_from_slice(hdr.as_bytes());
pkt.extend_from_slice(body);
packets.push(pkt);
Ok(())
})?;
Ok(packets)
}
/// Packetize one access unit, yielding each packet to `emit` as a `(header, shard bytes)`
/// pair — in exact wire order, which is also the order the session's nonce counter
/// advances. No per-packet allocation happens here, so the caller can write header and
/// shard straight into a pooled wire buffer and seal in place
/// ([`Session::seal_frame`](crate::session::Session::seal_frame)). An `emit` error aborts
/// the frame mid-way (packet numbering has already advanced — callers treat it as fatal).
///
/// Wire order is DATA-FIRST (latency plan T1.3): every block's data shards in block order,
/// then every block's parity shards in block order. In the lossless case a frame completes
/// the moment its last DATA shard arrives, so the completion-gating packet no longer sits
/// behind the parity tail of the paced spread (~fec% of the spread saved). The receiver is
/// order-agnostic (headers are self-describing; the reassembler completes each block on
/// `data + recovery ≥ k`), so this is not a wire-format change. `FLAG_SOF` stays on the
/// first emitted packet (block 0, shard 0); `FLAG_EOF` marks the last emitted packet —
/// the final parity shard, or the final data shard of a FEC-free frame.
///
/// `frame_index`: `Some(i)` stamps the AU with the caller's index — the punktfunk/1 encode
/// loop numbers video AUs itself so the encoder's RFI bookkeeping (LTR marks, DPB timestamps)
/// is 1:1 with what the client sees, surviving encoder rebuilds/resets that restart internal
/// counters. `None` draws from the internal counter (the legacy/self-numbering path). A
/// session must not mix the two styles for the same index space.
pub fn packetize_each(
&mut self,
frame: &[u8],
pts_ns: u64,
user_flags: u32,
frame_index: Option<u32>,
coder: &dyn ErasureCoder,
mut emit: impl FnMut(&PacketHeader, &[u8]) -> Result<()>,
) -> Result<()> {
let payload = self.shard_payload;
let frame_index = frame_index.unwrap_or_else(|| {
let i = self.next_frame_index;
self.next_frame_index = i.wrapping_add(1);
i
});
// At least one (zero-padded) data shard even for an empty frame.
let total_data = frame.len().div_ceil(payload).max(1);
let max_block = self.fec.max_data_per_block as usize;
let block_count = total_data.div_ceil(max_block).max(1);
let frame_bytes = frame.len() as u32;
// Defend the u16 wire fields against silent truncation. `Config::validate`
// already rejects configs that could reach these for valid frame sizes; this is
// the belt-and-suspenders for a frame larger than the negotiated maximum.
if payload > u16::MAX as usize {
return Err(PunktfunkError::InvalidArg("shard_payload exceeds u16"));
}
if block_count > u16::MAX as usize {
return Err(PunktfunkError::Unsupported(
"frame too large: block count exceeds u16",
));
}
// Stage the frame's one possibly-partial shard (the last) in the reusable
// zero-padded scratch; every full shard is referenced in place below.
let full_shards = frame.len() / payload;
self.tail.clear();
self.tail.resize(payload, 0);
let rem = frame.len() % payload;
if rem > 0 {
self.tail[..rem].copy_from_slice(&frame[full_shards * payload..]);
}
let tail = &self.tail;
let shard_at = |s: usize| -> &[u8] {
if s < full_shards {
&frame[s * payload..(s + 1) * payload]
} else {
tail.as_slice()
}
};
// Per-block shard geometry (deterministic — recomputed in both passes).
let block_data_count = |b: usize| ((b + 1) * max_block).min(total_data) - b * max_block;
// Parity for a `k`-shard block: the configured percentage, clamped so the block's wire
// total never exceeds what the peer will accept (see `max_total_shards`). The clamp only
// binds on blocks near `max_data_per_block`; smaller blocks keep the full adaptive range,
// so raising FEC still buys real protection wherever there is headroom. Bound as locals,
// not as a `&self` method: `emit_one` below would otherwise capture all of `self` and
// collide with the `&mut self.recovery[b]` parity borrow.
let (fec, max_total_shards) = (self.fec, self.max_total_shards);
let recovery_for =
move |k: usize| fec.recovery_for(k).min(max_total_shards.saturating_sub(k));
// One parity pool per block, reused across frames (steady-state zero-alloc).
if self.recovery.len() < block_count {
self.recovery.resize_with(block_count, Vec::new);
}
// Total parity across the frame decides where FLAG_EOF lands (the last emitted packet).
let mut total_recovery = 0usize;
for b in 0..block_count {
let k = block_data_count(b);
let m = recovery_for(k);
if k + m > u16::MAX as usize {
return Err(PunktfunkError::Unsupported("block shard count exceeds u16"));
}
total_recovery += m;
}
let mut emit_one =
|next_seq: &mut u32, b: usize, shard_index: usize, body: &[u8], flags: u8| {
let seq = *next_seq;
*next_seq = next_seq.wrapping_add(1);
let k = block_data_count(b);
let hdr = PacketHeader {
pts_ns,
frame_index,
stream_seq: seq,
frame_bytes,
user_flags,
block_index: b as u16,
block_count: block_count as u16,
data_shards: k as u16,
recovery_shards: recovery_for(k) as u16,
shard_index: shard_index as u16,
shard_bytes: payload as u16,
magic: PUNKTFUNK_MAGIC,
version: self.version,
fec_scheme: coder.scheme() as u8,
flags,
};
emit(&hdr, body)
};
let mut next_seq = self.next_seq;
// Pass 1 — per block: generate parity into the block's pool, emit the DATA shards.
for b in 0..block_count {
let first = b * max_block;
let k = block_data_count(b);
// This block's data shards: references into `frame` (plus the staged tail).
let data_shards: Vec<&[u8]> = (first..first + k).map(shard_at).collect();
let recovery_count = recovery_for(k);
coder.encode_into(&data_shards, recovery_count, &mut self.recovery[b])?;
for (shard_index, body) in data_shards.iter().enumerate() {
let mut flags = FLAG_PIC;
if b == 0 && shard_index == 0 {
flags |= FLAG_SOF;
}
if total_recovery == 0 && b + 1 == block_count && shard_index + 1 == k {
flags |= FLAG_EOF;
}
emit_one(&mut next_seq, b, shard_index, body, flags)?;
}
}
// Pass 2 — per block: emit the parity shards (the frame's tail on the wire).
let mut parity_left = total_recovery;
for b in 0..block_count {
let k = block_data_count(b);
let recovery_count = recovery_for(k);
for r in 0..recovery_count {
parity_left -= 1;
let mut flags = FLAG_PIC;
if parity_left == 0 {
flags |= FLAG_EOF;
}
let body: &[u8] = &self.recovery[b][r];
emit_one(&mut next_seq, b, k + r, body, flags)?;
}
}
self.next_seq = next_seq;
Ok(())
}
/// Open a **streamed** access unit (see [`StreamedAu`]). `frame_index` follows the same
/// contract as [`packetize_each`](Self::packetize_each): `Some(i)` = the caller owns the
/// video numbering; `None` draws from the internal counter.
pub fn begin_streamed(
&mut self,
pts_ns: u64,
user_flags: u32,
frame_index: Option<u32>,
) -> StreamedAu {
let frame_index = frame_index.unwrap_or_else(|| {
let i = self.next_frame_index;
self.next_frame_index = i.wrapping_add(1);
i
});
StreamedAu {
frame_index,
pts_ns,
user_flags,
pending: Vec::new(),
blocks_out: 0,
total_bytes: 0,
opened: false,
}
}
/// Feed one encoder chunk into a streamed AU, emitting every block that COMPLETES under
/// sentinel headers (`frame_bytes = 0`, `block_count = 0`, exactly `max_data_per_block`
/// data shards — the receiver's offset formula needs no total). Flushes only while
/// STRICTLY more than one block is buffered, so the frame's last block — whose header must
/// carry the real totals — is never emitted here. Packets reach `emit` in wire order (the
/// caller's nonce order), data then parity per block.
pub fn push_streamed(
&mut self,
au: &mut StreamedAu,
chunk: &[u8],
coder: &dyn ErasureCoder,
mut emit: impl FnMut(&PacketHeader, &[u8]) -> Result<()>,
) -> Result<()> {
au.total_bytes += chunk.len() as u64;
au.pending.extend_from_slice(chunk);
let block_bytes = self.fec.max_data_per_block as usize * self.shard_payload;
while au.pending.len() > block_bytes {
// Room for this sentinel block AND the final block after it, within the peer's
// per-frame block ceiling and the u16 wire field.
if au.blocks_out as usize + 2 > self.max_blocks.min(u16::MAX as usize) {
return Err(PunktfunkError::Unsupported(
"streamed AU exceeds the negotiated max_frame_bytes",
));
}
let sof = !au.opened;
let (bi, pts, uf) = (au.blocks_out, au.pts_ns, au.user_flags);
let fi = au.frame_index;
self.emit_streamed_block(
fi,
pts,
uf,
bi,
&au.pending[..block_bytes],
0,
0,
sof,
false,
coder,
&mut emit,
)?;
au.pending.drain(..block_bytes);
au.blocks_out += 1;
au.opened = true;
}
Ok(())
}
/// Close a streamed AU: seal the final block — its headers carry the REAL
/// `frame_bytes`/`block_count`, which retro-validate the whole frame at the receiver — with
/// `FLAG_EOF` on the last emitted packet. An empty AU degenerates to today's single
/// zero-padded-shard frame (`block_count = 1`, never a sentinel).
pub fn finish_streamed(
&mut self,
au: StreamedAu,
coder: &dyn ErasureCoder,
mut emit: impl FnMut(&PacketHeader, &[u8]) -> Result<()>,
) -> Result<()> {
let frame_bytes = u32::try_from(au.total_bytes)
.map_err(|_| PunktfunkError::Unsupported("streamed AU exceeds u32 bytes"))?;
let block_count = au.blocks_out + 1;
self.emit_streamed_block(
au.frame_index,
au.pts_ns,
au.user_flags,
au.blocks_out,
&au.pending,
frame_bytes,
block_count,
!au.opened,
true,
coder,
&mut emit,
)
}
/// Seal ONE streamed block (data shards + its parity, in wire order). Sentinel blocks pass
/// `frame_bytes = 0` / `block_count = 0`; the final block passes the real totals. `sof`
/// marks the frame's very first packet, `eof` its very last.
#[allow(clippy::too_many_arguments)]
fn emit_streamed_block(
&mut self,
frame_index: u32,
pts_ns: u64,
user_flags: u32,
block_index: u16,
bytes: &[u8],
frame_bytes: u32,
block_count: u16,
sof: bool,
eof: bool,
coder: &dyn ErasureCoder,
emit: &mut impl FnMut(&PacketHeader, &[u8]) -> Result<()>,
) -> Result<()> {
let payload = self.shard_payload;
if payload > u16::MAX as usize {
return Err(PunktfunkError::InvalidArg("shard_payload exceeds u16"));
}
// At least one (zero-padded) data shard even for an empty final block (empty AU).
let k = bytes.len().div_ceil(payload).max(1);
let m = self
.fec
.recovery_for(k)
.min(self.max_total_shards.saturating_sub(k));
if k + m > u16::MAX as usize {
return Err(PunktfunkError::Unsupported("block shard count exceeds u16"));
}
// Stage the one possibly-partial (or empty-frame) shard in the zero-padded scratch.
let full_shards = bytes.len() / payload;
self.tail.clear();
self.tail.resize(payload, 0);
let rem = bytes.len() % payload;
if rem > 0 {
self.tail[..rem].copy_from_slice(&bytes[full_shards * payload..]);
}
let tail = &self.tail;
let shard_at = |s: usize| -> &[u8] {
if s < full_shards {
&bytes[s * payload..(s + 1) * payload]
} else {
tail.as_slice()
}
};
let data_shards: Vec<&[u8]> = (0..k).map(shard_at).collect();
if self.recovery.is_empty() {
self.recovery.push(Vec::new());
}
coder.encode_into(&data_shards, m, &mut self.recovery[0])?;
let mut next_seq = self.next_seq;
let mut emit_one = |next_seq: &mut u32, shard_index: usize, body: &[u8], flags: u8| {
let seq = *next_seq;
*next_seq = next_seq.wrapping_add(1);
let hdr = PacketHeader {
pts_ns,
frame_index,
stream_seq: seq,
frame_bytes,
user_flags,
block_index,
block_count,
data_shards: k as u16,
recovery_shards: m as u16,
shard_index: shard_index as u16,
shard_bytes: payload as u16,
magic: PUNKTFUNK_MAGIC,
version: self.version,
fec_scheme: coder.scheme() as u8,
flags,
};
emit(&hdr, body)
};
for (shard_index, body) in data_shards.iter().enumerate() {
let mut flags = FLAG_PIC;
if sof && shard_index == 0 {
flags |= FLAG_SOF;
}
if eof && m == 0 && shard_index + 1 == k {
flags |= FLAG_EOF;
}
emit_one(&mut next_seq, shard_index, body, flags)?;
}
for r in 0..m {
let mut flags = FLAG_PIC;
if eof && r + 1 == m {
flags |= FLAG_EOF;
}
let body: &[u8] = &self.recovery[0][r];
emit_one(&mut next_seq, k + r, body, flags)?;
}
self.next_seq = next_seq;
Ok(())
}
}