From 4240b76182f0e5556115f27b70499abcdaaeab41 Mon Sep 17 00:00:00 2001 From: enricobuehler Date: Fri, 31 Jul 2026 16:30:45 +0200 Subject: [PATCH] feat(android): slices feed the decoder as they arrive, not when the AU closes The decode loop used to hold every access unit until its last wire packet landed. With slice-progressive delivery (frame_parts) each newly contiguous piece now goes straight into MediaCodec under BUFFER_FLAG_PARTIAL_FRAME, the closing piece drops the flag, and the decoder chews the front of the frame while its tail is still on the wire. Feeding partial input means owning its failure modes without a codec flush: a broken sequence (dropped piece, orphan, oversize) closes the dead AU with an empty non-partial buffer at its own pts, arms the re-anchor freeze so the concealed output never reaches glass, and requests a recovery keyframe - the same machinery ordinary loss already rides. Per-AU accounting keeps its units: the RFI gap detector notes an AU once, and the HUD stamps, host/network split and the phase-lock arrival sensor ride only the completing delivery. The opt-in is decoder truth and loop truth: every decoder this device would use must pass FEATURE_PartialFrame, and only the async loop may see parts - the legacy sync loop feeds whole AUs and stays that way. Smoke-checked on-glass (NP3 vs the Windows host): byte-identical degenerate path, presenter profile unchanged. Co-Authored-By: Claude Fable 5 --- .../kotlin/io/unom/punktfunk/HostConnect.kt | 6 +- .../io/unom/punktfunk/kit/NativeBridge.kt | 4 + .../io/unom/punktfunk/kit/VideoDecoders.kt | 25 +++ .../android/native/src/decode/async_loop.rs | 143 ++++++++++++++++-- clients/android/native/src/session/connect.rs | 8 +- 5 files changed, 167 insertions(+), 19 deletions(-) diff --git a/clients/android/app/src/main/kotlin/io/unom/punktfunk/HostConnect.kt b/clients/android/app/src/main/kotlin/io/unom/punktfunk/HostConnect.kt index 9e96414a..31e89b29 100644 --- a/clients/android/app/src/main/kotlin/io/unom/punktfunk/HostConnect.kt +++ b/clients/android/app/src/main/kotlin/io/unom/punktfunk/HostConnect.kt @@ -49,7 +49,11 @@ suspend fun connectToHost( host, port, w, h, hz, identity.certPem, identity.privateKeyPem, pinHex, settings.bitrateKbps, settings.compositor, gamepadPref, - hdrEnabled, VideoDecoders.multiSliceTolerant(), settings.audioChannels, + hdrEnabled, VideoDecoders.multiSliceTolerant(), + // Slice-progressive delivery: decoder truth AND the async decode loop — the legacy + // sync loop feeds whole AUs only, so parts must never arrive when it is selected. + settings.lowLatencyMode && VideoDecoders.partialFrameCapable(), + settings.audioChannels, // What this device can decode (H.264|HEVC always, AV1 when a real decoder exists) + // the user's soft codec preference — the host resolves the emitted codec from both. VideoDecoders.decodableCodecBits(), settings.preferredCodec(), timeoutMs, diff --git a/clients/android/kit/src/main/kotlin/io/unom/punktfunk/kit/NativeBridge.kt b/clients/android/kit/src/main/kotlin/io/unom/punktfunk/kit/NativeBridge.kt index 5a3c815b..ed03a18c 100644 --- a/clients/android/kit/src/main/kotlin/io/unom/punktfunk/kit/NativeBridge.kt +++ b/clients/android/kit/src/main/kotlin/io/unom/punktfunk/kit/NativeBridge.kt @@ -51,6 +51,10 @@ object NativeBridge { * ([VideoDecoders.multiSliceTolerant]) — advertises `VIDEO_CAP_MULTI_SLICE`; false keeps * the host at single-slice frames (the safe pre-0.17 wire shape). */ multiSliceOk: Boolean, + /** Every decoder this device would use accepts partial-frame input + * ([VideoDecoders.partialFrameCapable]) — opts into slice-progressive delivery (the + * decode loop then feeds slices with `BUFFER_FLAG_PARTIAL_FRAME` as they arrive). */ + framePartsOk: Boolean, audioChannels: Int, /** `quic::CODEC_*` bitfield of codecs this device decodes ([VideoDecoders.decodableCodecBits]); * `0` falls back to H.264|HEVC. The host resolves the emitted codec from this ∩ its GPU. */ diff --git a/clients/android/kit/src/main/kotlin/io/unom/punktfunk/kit/VideoDecoders.kt b/clients/android/kit/src/main/kotlin/io/unom/punktfunk/kit/VideoDecoders.kt index b2dfcce2..bbf591c6 100644 --- a/clients/android/kit/src/main/kotlin/io/unom/punktfunk/kit/VideoDecoders.kt +++ b/clients/android/kit/src/main/kotlin/io/unom/punktfunk/kit/VideoDecoders.kt @@ -75,6 +75,31 @@ object VideoDecoders { } } + /** + * True when EVERY decoder this device would stream with supports partial-frame input + * (MediaCodec `FEATURE_PartialFrame`): the client may then opt into slice-progressive + * delivery and feed slices with `BUFFER_FLAG_PARTIAL_FRAME` ahead of the AU's tail. + * Codec selection happens after connect, so all advertised mimes must qualify — the same + * conservative shape as [multiSliceTolerant] (an unprobeable pick disqualifies). + */ + fun partialFrameCapable(): Boolean { + val mimes = buildList { + add("video/avc") + add("video/hevc") + if (decodableCodecBits() and 4 != 0) add("video/av01") + } + val infos = runCatching { MediaCodecList(MediaCodecList.REGULAR_CODECS).codecInfos } + .getOrNull() ?: return false + return mimes.all { mime -> + val pick = pickDecoder(mime)?.name ?: return@all false + val info = infos.firstOrNull { it.name == pick } ?: return@all false + runCatching { + info.getCapabilitiesForType(mime) + .isFeatureSupported(CodecCapabilities.FEATURE_PartialFrame) + }.getOrDefault(false) + } + } + fun pickDecoder(mime: String): DecoderChoice? { if (mime.isEmpty()) return null val infos = runCatching { MediaCodecList(MediaCodecList.REGULAR_CODECS).codecInfos } diff --git a/clients/android/native/src/decode/async_loop.rs b/clients/android/native/src/decode/async_loop.rs index 39c8d496..fc4dcac7 100644 --- a/clients/android/native/src/decode/async_loop.rs +++ b/clients/android/native/src/decode/async_loop.rs @@ -278,6 +278,8 @@ pub(super) fn run_async( let mut discarded: u64 = 0; // AUs larger than the codec input buffer, dropped whole (see `feed`/`feed_ready`). let mut oversized_dropped: u64 = 0; + // Slice-progressive continuity ledger (see `PartFeed`). + let mut part_open: Option = None; // Freeze-until-reanchor gate (see the sync loop for the rationale). Armed on a frame-index gap // (the feeder's Au verdict), a parked-AU overflow drop, a dropped-count climb, or a recoverable // codec error; `recovery_flags` carries each AU's user_flags from `dispatch_event` (feed) to @@ -358,6 +360,8 @@ pub(super) fn run_async( &mut free_inputs, &mut fed, &mut oversized_dropped, + &mut part_open, + &mut gate, ); let had_output = !ready.is_empty(); let rendered_before = rendered; @@ -580,11 +584,15 @@ fn feeder_loop( // invalidation request so an RFI-capable host recovers with a cheap clean P-frame // instead of a full IDR (the frames_dropped keyframe path is the backstop). The gap // verdict rides the Au event so the decode loop arms its freeze gate on the same signal. - let gap = client.note_frame_index(frame.frame_index); + // Slice-progressive parts repeat their AU's index — note it once, on the + // AU's first piece (or a whole delivery), so the RFI gap detector keeps + // counting AUs. + let au_first = frame.part.is_none_or(|p| p.first); + let gap = au_first && client.note_frame_index(frame.frame_index); // Park the receipt stamp (keyed by the pts the codec echoes) whenever the `decode` // stage is consumed: the HUD, or the ABR decode signal (`measure_decode`). The // HUD-only `received` point + host/network split stay gated on the overlay. - if stats.enabled() || measure_decode { + if (stats.enabled() || measure_decode) && frame.complete { // Core reassembly-completion stamp (ABI v9), NOT the pull instant: stamping // here would fold the hand-off queue wait into the network latency figure // (a client-side standing backlog masquerading as network). 0 = older core. @@ -607,7 +615,10 @@ fn feeder_loop( let lat_ns = received_ns + clock_offset - frame.pts_ns as i128; let lat_us = (lat_ns > 0 && lat_ns < 10_000_000_000) .then_some((lat_ns / 1000) as u64); - stats.note_received(frame.data.len(), lat_us, clock_offset != 0); + // On a parts stream the completing delivery carries only the AU's + // suffix — its offset restores the full AU byte count for bitrate. + let au_len = frame.part.map_or(0, |p| p.offset as usize) + frame.data.len(); + stats.note_received(au_len, lat_us, clock_offset != 0); if let Some(hostnet_us) = lat_us { pending_split.push_back((frame.pts_ns, hostnet_us)); if pending_split.len() > PENDING_SPLIT_CAP { @@ -669,20 +680,27 @@ fn dispatch_event( if gap { gate.arm(Instant::now()); } - recovery_flags.push_back((f.pts_ns / 1000, f.flags)); - if recovery_flags.len() > IN_FLIGHT_CAP { - recovery_flags.pop_front(); + // One entry per AU (parts share the pts): the completing delivery carries it. + if f.complete { + recovery_flags.push_back((f.pts_ns / 1000, f.flags)); + if recovery_flags.len() > IN_FLIGHT_CAP { + recovery_flags.pop_front(); + } } // Phase-lock v3 sensor: the ARRIVAL stamp (reassembly completion, realtime) — the // phase the host actually controls. The latch-based v2 sensor measured downstream // of the decoder pipeline, which absorbed the host's actuation (on-glass 07-31). - arrival_stamps.push(if f.received_ns > 0 { - f.received_ns as i128 - } else { - now_realtime_ns() - }); - if arrival_stamps.len() > 256 { - arrival_stamps.remove(0); + // Phase sensor: AU completion is the arrival the host's hold actually moves — + // prefix parts would smear the phase toward the first slice's landing. + if f.complete { + arrival_stamps.push(if f.received_ns > 0 { + f.received_ns as i128 + } else { + now_realtime_ns() + }); + if arrival_stamps.len() > 256 { + arrival_stamps.remove(0); + } } pending_aus.push_back(f); if pending_aus.len() > FRAME_PARK_CAP { @@ -715,10 +733,35 @@ fn dispatch_event( false } +/// `AMEDIACODEC_BUFFER_FLAG_PARTIAL_FRAME` (NDK ≥ 26, gated by the Kotlin +/// `FEATURE_PartialFrame` probe): this input buffer is a PIECE of an AU — the codec assembles +/// pieces until a buffer WITHOUT the flag closes the AU. +const BUFFER_FLAG_PARTIAL_FRAME: u32 = 8; + +/// The slice-progressive feed's open access unit: parts already queued into the codec under +/// [`BUFFER_FLAG_PARTIAL_FRAME`], awaiting the rest. Loop-local — a codec rebuild tears the +/// whole loop down, so the state can never outlive the codec instance it fed. +pub(super) struct PartFeed { + index: u32, + /// The AU byte offset the next part must carry — a mismatch means the hand-off dropped a + /// piece (memory cap / jump-to-live clear) and the AU is unrecoverable. + expected: usize, + /// The dead-close pts: an abandoned AU is CLOSED with an empty non-PARTIAL buffer at its + /// own pts — the codec then emits (concealed garbage) at that pts, which the reanchor + /// freeze gate withholds from glass while the keyframe request recovers the chain. No + /// mid-stream codec flush needed. + pts_us: u64, +} + /// Queue as many parked AUs as there are free input buffer slots (async mode: the indices come from /// `InputAvailable` callbacks, not a dequeue). Each AU is copied into its codec input buffer and /// submitted; an AU larger than the buffer is DROPPED (+ a recovery keyframe requested) — a /// truncated AU is corrupt input the decoder chews on silently, poisoning the reference chain. +/// +/// Slice-progressive deliveries ([`Frame::part`]) feed as they arrive: every piece rides +/// [`BUFFER_FLAG_PARTIAL_FRAME`] except the AU's last, all at the AU's pts. `part_open` is the +/// continuity ledger — any break (gap, orphan, oversize) abandons the AU per [`PartFeed::pts_us`]'s +/// close contract and re-syncs at the next `first`. fn feed_ready( codec: &MediaCodec, client: &NativeClient, @@ -726,11 +769,49 @@ fn feed_ready( free_inputs: &mut VecDeque, fed: &mut u64, oversized_dropped: &mut u64, + part_open: &mut Option, + gate: &mut ReanchorGate, ) { while !pending_aus.is_empty() && !free_inputs.is_empty() { let idx = free_inputs.pop_front().unwrap(); let frame = pending_aus.pop_front().unwrap(); let pts_us = frame.pts_ns / 1000; + let (first, last, offset) = match frame.part { + None => (true, true, 0usize), + Some(p) => (p.first, p.last, p.offset as usize), + }; + // Continuity ledger. `continues` = this piece extends the open AU exactly; + // anything else with an AU open means that AU died mid-flight and must be closed + // (empty non-PARTIAL buffer at ITS pts) before this frame may touch the codec. + let continues = part_open + .as_ref() + .is_some_and(|o| frame.frame_index == o.index && offset == o.expected && !first); + if !continues { + if let Some(o) = part_open.take() { + // Spend THIS slot on the close; the current frame re-queues for the next one. + if let Err(e) = codec.queue_input_buffer_by_index(idx, 0, 0, o.pts_us, 0) { + log::warn!("decode: close of abandoned partial AU {}: {e}", o.index); + } + log::warn!( + "decode: partial AU {} abandoned mid-feed — closed empty, requesting keyframe", + o.index + ); + // The close makes the codec emit concealed garbage at the dead pts — freeze it + // off the glass until the recovery keyframe re-anchors. + gate.arm(Instant::now()); + let _ = client.request_keyframe(); + pending_aus.push_front(frame); + continue; + } + // No AU open: an orphan non-first piece lost its head upstream — discard and + // re-sync at the next `first` (the recovery request rides the same loss). + if !first { + free_inputs.push_front(idx); + gate.arm(Instant::now()); + let _ = client.request_keyframe(); + continue; + } + } let Some(dst) = codec.input_buffer(idx) else { log::warn!("decode: input_buffer({idx}) returned None — dropping AU"); continue; @@ -747,6 +828,16 @@ fn feed_ready( *oversized_dropped ); let _ = client.request_keyframe(); + if frame.part.is_some() { + gate.arm(Instant::now()); + // Pieces already queued can't be unqueued: poison the ledger so the next + // delivery mismatches and takes the close-empty path above. + *part_open = Some(PartFeed { + index: frame.frame_index, + expected: usize::MAX, + pts_us, + }); + } continue; } let n = au.len(); @@ -755,10 +846,32 @@ fn feed_ready( unsafe { std::ptr::copy_nonoverlapping(au.as_ptr(), dst.as_mut_ptr().cast::(), n); } - if let Err(e) = codec.queue_input_buffer_by_index(idx, 0, n, pts_us, 0) { + let flags = if last { 0 } else { BUFFER_FLAG_PARTIAL_FRAME }; + if let Err(e) = codec.queue_input_buffer_by_index(idx, 0, n, pts_us, flags) { log::warn!("decode: queue_input_buffer_by_index: {e}"); + if frame.part.is_some() && !last { + // The piece never reached the codec — same unrecoverable-AU shape as oversize. + *part_open = Some(PartFeed { + index: frame.frame_index, + expected: usize::MAX, + pts_us, + }); + } } else { - *fed += 1; + // `fed` counts ACCESS UNITS toward the HUD's fed/decoded balance — the closing + // piece (or a whole AU) bumps it. + if last { + *fed += 1; + } + *part_open = if last { + None + } else { + Some(PartFeed { + index: frame.frame_index, + expected: offset + n, + pts_us, + }) + }; } } } diff --git a/clients/android/native/src/session/connect.rs b/clients/android/native/src/session/connect.rs index eefbe48e..cffd8e48 100644 --- a/clients/android/native/src/session/connect.rs +++ b/clients/android/native/src/session/connect.rs @@ -116,6 +116,7 @@ pub extern "system" fn Java_io_unom_punktfunk_kit_NativeBridge_nativeConnect<'lo gamepad_pref: jint, hdr_enabled: jboolean, multi_slice_ok: jboolean, + frame_parts_ok: jboolean, audio_channels: jint, video_codecs: jint, preferred_codec: jint, @@ -225,9 +226,10 @@ pub extern "system" fn Java_io_unom_punktfunk_kit_NativeBridge_nativeConnect<'lo // No non-video caps: this client does not render the host cursor locally (no shape/state // planes in the jni surface), so advertising CLIENT_CAP_CURSOR would stream cursor-less. 0, - // Slice-progressive delivery: off until the decode loop feeds MediaCodec with - // BUFFER_FLAG_PARTIAL_FRAME (P2d flips this behind a FEATURE_PartialFrame probe). - false, + // Slice-progressive delivery, by decoder truth (Kotlin probes FEATURE_PartialFrame on + // every decoder this device would use): AU prefixes then arrive as `Frame::part` + // pieces and the decode loop feeds them with BUFFER_FLAG_PARTIAL_FRAME. + frame_parts_ok != 0, launch, // a store-qualified library id to boot into a game, or None for the desktop device_name, // Kotlin's Build.MODEL — the host's approval-list / trust-store label pin, // Some → Crypto on host-fp mismatch