From 30169c312e031876174801f44361317784e5a38e Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Wed, 23 Sep 2026 12:28:06 -0700 Subject: [PATCH] fix(relay): keep only finished groups warm when a track goes idle When the last reader leaves, the origin front parks the track: it copies the cached groups to a warm local track, then drops the source copy, which cancels the upstream subscription and resets the group still in flight. The warm copy had adopted that open group, so its resume boundary pointed mid-group into a now-dead head. A rejoin then subscribed upstream from that frame and the group never reached the returning reader. Adopt only finished groups. The re-splice resumes at the open group's first frame and delivers it whole. Co-Authored-By: Claude Opus 5.5 --- doc/lib/c/index.md | 2 +- rs/moq-net/src/model/group.rs | 5 ++ rs/moq-net/src/model/origin.rs | 9 +++- rs/moq-net/tests/rejoin.rs | 99 ++++++++++++++++++++++++++++++++++ 4 files changed, 112 insertions(+), 3 deletions(-) create mode 100644 rs/moq-net/tests/rejoin.rs diff --git a/doc/lib/c/index.md b/doc/lib/c/index.md index 769e91a51b..1503770533 100644 --- a/doc/lib/c/index.md +++ b/doc/lib/c/index.md @@ -38,7 +38,7 @@ and `target/include/moq.h`. - **Raw decode output.** `moq_video_decoder_output` selects the decoded CPU pixel format (`MOQ_VIDEO_PIXEL_FORMAT_I420` or `_RGBA`) and target size (`width`/`height`, both zero for native; otherwise even and non-zero). Unknown formats and invalid sizes fail `moq_decode_video` before subscribing; accepted requests deliver exactly that layout or fail on the terminal callback. - **Encoded video metadata.** `moq_video_init.hint` is a zero-initialized `moq_video_hint` with `has_*` flags for coded dimensions, bitrate (bits per second), frame rate, and latency preference. Hints seed a video codec track's catalog; detected dimensions take precedence. - **Client config.** A zeroed `moq_client_config` means the defaults for every knob, which is what lets a new one be appended without disturbing callers. Fields cover protocol (`versions`), TLS (`tls_fingerprints`, `tls_roots`, `tls_cert`/`_key`, `tls_host_name`), transport (`bind`, `connect_timeout_us`, the Happy Eyeballs delays, `websocket_enabled`), and tuning (reconnect backoff, `quic_*`). Every duration is in microseconds. A knob whose default isn't zero carries a `has_*` flag, so setting `backoff_timeout_us = 0` needs `has_backoff_timeout = true` to mean "retry forever" rather than "use the default". `moq_client_defaults()` reports what a NULL config dials with. -- **Demand.** A watcher on a published track (`moq_publish_track_demand`, `moq_publish_media_demand`, `moq_encode_video_demand`, `moq_encode_audio_demand`) calls `on_demand` with `MOQ_DEMAND_USED` or `MOQ_DEMAND_UNUSED` right away and again on every change, so an encoder on a battery-powered device runs only while someone is watching. The first call is the current state, so a track that went unused before the watcher existed still reports it. `moq_publish_demand_cancel` stops it; the terminal callback still fires. A container has no single demand and is refused. Demand follows the last real subscriber: an origin that served the track drops its source copy on the unused edge and keeps only the groups it already cached warm for 30 seconds, so the cache linger does not delay the unused edge. +- **Demand.** A watcher on a published track (`moq_publish_track_demand`, `moq_publish_media_demand`, `moq_encode_video_demand`, `moq_encode_audio_demand`) calls `on_demand` with `MOQ_DEMAND_USED` or `MOQ_DEMAND_UNUSED` right away and again on every change, so an encoder on a battery-powered device runs only while someone is watching. The first call is the current state, so a track that went unused before the watcher existed still reports it. `moq_publish_demand_cancel` stops it; the terminal callback still fires. A container has no single demand and is refused. Demand follows the last real subscriber: an origin that served the track drops its source copy on the unused edge and keeps only the finished groups it already cached warm for 30 seconds, so the cache linger does not delay the unused edge. - **Requests.** `moq_publish_dynamic` serves subscriptions to tracks the broadcast never declared: each arrives as a request handle, read its name with `moq_track_request_name`, then `moq_track_request_accept` (a raw track handle), `moq_track_request_video` / `_audio` (the media handle `moq_publish_video` / `_audio` return), or `moq_track_request_abort` with an application code the subscriber sees. Without a live handler an unknown name is refused. `moq_publish_track_dynamic` does the same for fetches of groups a track no longer has cached, delivered as `moq_group_request_*` (`sequence`, `priority`, `frame_start`); `moq_group_request_accept` starts the producer at `frame_start` so written frames keep their group indices. Register it with `moq_track_request_dynamic` before accepting a track that was itself requested by a fetch, so that pending group survives the transition. Both handlers stop with `moq_publish_dynamic_cancel`. - **Everything the bindings can do** ([list](/lib/#what-every-binding-can-do)): media publish and consume with the catalog managed for you, raw pixels and PCM with the codec inside (`moq_encode_video`, `moq_encode_audio`, and the `moq_decode_*` mirrors), raw tracks with timestamps and datagrams, JSON snapshot and stream tracks, group fetch, catalog sections, shared video properties, and stalled hints. The three advertising operations are `moq_origin_create_broadcast` (locally discoverable producer), `moq_publish_announce` / `moq_publish_unannounce` (exact-path advertisement), and `moq_origin_dynamic` (a claim over a path prefix and everything beneath it; `""` for everything). A route is a capability, not an inventory. `moq_origin_announced` takes a literal prefix and an optional relative pattern filter; `moq_announce_update.prefix` stays relative to the origin, while `captures` reports what each wildcard matched when `has_captures` is true. diff --git a/rs/moq-net/src/model/group.rs b/rs/moq-net/src/model/group.rs index d726d391fb..0205ae10b4 100644 --- a/rs/moq-net/src/model/group.rs +++ b/rs/moq-net/src/model/group.rs @@ -803,6 +803,11 @@ impl Producer { self.alive.aborted.load(Ordering::Acquire) } + /// Whether the group was finished: it holds every frame it will ever have. + pub(crate) fn is_finished(&self) -> bool { + self.state.read().fin.is_some() + } + /// The index of the first frame this group still holds, or `None` once it has been /// aborted. Non-zero when the group started later (see [`Self::start_at`]); a reader /// positioned below it is [`Error::Lagged`]. diff --git a/rs/moq-net/src/model/origin.rs b/rs/moq-net/src/model/origin.rs index edd5607203..2d23fb7886 100644 --- a/rs/moq-net/src/model/origin.rs +++ b/rs/moq-net/src/model/origin.rs @@ -1697,12 +1697,17 @@ impl Drop for WarmCopy { } } -/// Cache `source`'s groups on a new local track the origin owns. +/// Cache `source`'s finished groups on a new local track the origin owns. fn warm_copy(source: &track::Consumer) -> Option { let info = source.cached_info()?; let mut track = track::Producer::new(Arc::new(source.broadcast().clone()), source.name(), info); for (group, visible) in source.cached_groups() { - let _ = track.adopt_group(group, visible); + // An open group is left for the re-splice to deliver whole. Dropping the source + // copy resets it mid-transfer, and its dead head would anchor the next takeover + // mid-group, asking upstream for a tail no returning reader can use. + if group.is_finished() { + let _ = track.adopt_group(group, visible); + } } let dynamic = track.dynamic(); Some(WarmCopy { diff --git a/rs/moq-net/tests/rejoin.rs b/rs/moq-net/tests/rejoin.rs new file mode 100644 index 0000000000..5ba6a80b3c --- /dev/null +++ b/rs/moq-net/tests/rejoin.rs @@ -0,0 +1,99 @@ +//! A reader leaving and rejoining a track relayed over the in-memory mock transport. + +mod support; + +use std::time::Duration; + +use moq_net::{Hop, Timestamp, Version, group, track}; +use support::harness::{MockConnectOptions, connect_mock}; + +/// Maximum time any single test may run before being treated as a deadlock. +const TEST_TIMEOUT: Duration = Duration::from_secs(10); + +/// Build an origin producer, spawning its driver on the ambient runtime. +fn produce_origin(hop: u64) -> moq_net::origin::Producer { + let (producer, driver) = moq_net::origin::Producer::new(moq_net::origin::Config::new(Hop::new(hop).unwrap())); + tokio::spawn(support::harness::run(driver)); + producer +} + +async fn read_all(group: &mut group::Consumer) -> Vec> { + let mut frames = Vec::new(); + while let Some(frame) = group.read_frame().await.expect("group aborted") { + frames.push(frame.payload.to_vec()); + } + frames +} + +/// The group in flight when the last reader left comes back whole on rejoin. +/// +/// The relay cancels its idle upstream subscription, resetting that group mid-transfer. +/// Resuming the rejoin at the frame where the reset landed would ask upstream for a tail +/// whose head is gone, so the group would never reach the returning reader. +#[tokio::test] +async fn rejoin_recovers_the_group_reset_on_leave() { + for version in ["moq-lite-05", "moq-lite-06"] { + tokio::time::timeout(TEST_TIMEOUT, async { + let publisher = produce_origin(1); + let relay = produce_origin(2); + + let broadcast = publisher.create_broadcast("bench").unwrap(); + let track = broadcast.create_track("video", None).unwrap(); + broadcast.announce(Default::default()).unwrap(); + + let mut options = MockConnectOptions::new(version.parse::().unwrap()); + options.server_publish = Some(publisher); + options.client_subscribe = Some(relay.clone()); + let _pair = connect_mock(options).await; + + let consumer = relay.consume(); + consumer.routed("bench").await.unwrap(); + let remote = consumer.request_broadcast("bench").await.unwrap(); + + let ts = |ms| Timestamp::from_millis(ms).unwrap(); + let prefs = || track::Subscription::default().with_max_age(Duration::from_secs(10)); + + let mut group = track.append_group().unwrap(); + group.write_frame(ts(0), b"a0".as_ref()).unwrap(); + group.finish().unwrap(); + let mut open = track.append_group().unwrap(); + open.write_frame(ts(100), b"b0".as_ref()).unwrap(); + + let mut sub = remote.track("video").unwrap().subscribe(prefs()).await.unwrap(); + loop { + let mut group = sub.recv_group().await.unwrap().unwrap(); + group.read_frame().await.unwrap().unwrap(); + if group.sequence == 1 { + break; + } + } + + // Leave mid-group; the publisher cuts the open group once demand is gone. + drop(sub); + track.unused().await.unwrap(); + open.write_frame(ts(133), b"b1".as_ref()).unwrap(); + open.finish().unwrap(); + + let mut live = track.append_group().unwrap(); + live.write_frame(ts(5000), b"c0".as_ref()).unwrap(); + live.finish().unwrap(); + + let mut sub = remote.track("video").unwrap().subscribe(prefs()).await.unwrap(); + let mut rejoined = Vec::new(); + loop { + let mut group = sub.recv_group().await.unwrap().unwrap(); + let sequence = group.sequence; + rejoined.push((sequence, read_all(&mut group).await)); + if sequence == 2 { + break; + } + } + assert!( + rejoined.contains(&(1, vec![b"b0".to_vec(), b"b1".to_vec()])), + "{version}: the reset group never came back whole: {rejoined:?}" + ); + }) + .await + .unwrap_or_else(|_| panic!("{version} timed out")); + } +}