From dcbf857bc694f88cd3437776ce7a2f87d160adb5 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 5 Oct 2026 15:33:25 -0700 Subject: [PATCH 1/4] quest: claim unannounce-demand-release From 6a0b3f7458774ff9adb866a785e9adfb7597de58 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 5 Oct 2026 16:28:24 -0700 Subject: [PATCH 2/4] fix(moq-net): release a retracted broadcast's track demand when its last reader leaves Since #4741, a front that ends on a retraction dropped its tracks' resume producers, but the logical track's state (kept allocated by weak handles) still held the serving copy, so the source track stayed subscribed until its publisher ended it. A live producer never does. The front now holds the tracks it ended in flight until their last reader leaves (or nothing owns the origin), then lets go of the copy, so the source goes unused as an unread track parks. Co-Authored-By: Claude Opus 5.5 --- doc/concept/moq-lite.md | 3 +- quest/m0/broadcast-epoch/README.md | 4 - .../unannounce-demand-release.md | 78 ----------------- rs/moq-net/src/model/origin.rs | 84 +++++++++++++++---- rs/moq-net/src/model/resume.rs | 34 ++++++-- rs/moq-net/src/util.rs | 6 ++ rs/moq-net/tests/unannounce_release.rs | 50 +++++++++++ rs/moq-tokio/tests/unannounce.rs | 47 +++++++++++ 8 files changed, 202 insertions(+), 104 deletions(-) delete mode 100644 quest/m0/broadcast-epoch/unannounce-demand-release.md create mode 100644 rs/moq-net/tests/unannounce_release.rs create mode 100644 rs/moq-tokio/tests/unannounce.rs diff --git a/doc/concept/moq-lite.md b/doc/concept/moq-lite.md index 6f906db54f..bf3817736c 100644 --- a/doc/concept/moq-lite.md +++ b/doc/concept/moq-lite.md @@ -115,7 +115,8 @@ can be neither discovered nor requested. A broadcast published locally competes with remote routes to its path on cost like any other route, winning only a tie. Retracting a route (an unannounce, or the peer's `ANNOUNCE_END`) stops new requests from resolving through it but leaves subscriptions already -in flight alone: each track runs to its own end or failure. On moq-lite 05 and +in flight alone: each track runs to its own end or failure, or until its last +subscriber leaves, which cancels it upstream. On moq-lite 05 and newer, a clean end requires `SUBSCRIBE_END` before the publisher's FIN. A FIN without that declaration fails the subscription with `ProtocolViolation`; older moq-lite versions use FIN alone. moq-transport requires `PUBLISH_DONE` before FIN. diff --git a/quest/m0/broadcast-epoch/README.md b/quest/m0/broadcast-epoch/README.md index 7f968003ee..1898c1dbb5 100644 --- a/quest/m0/broadcast-epoch/README.md +++ b/quest/m0/broadcast-epoch/README.md @@ -55,9 +55,6 @@ Decided: [Stats epochs](/quest/m0/broadcast-epoch/stats-epoch.md), which also gates the release (decided 2026-10-04): a restarted stats node under a reused name stalls its viewers the same way. -- [Retracted demand release](/quest/m0/broadcast-epoch/unannounce-demand-release.md) also - gates the release (decided 2026-10-05): a regression from #4741 on main that - `release` lacks. - Decided in the 2026-10-05 audit: the m1 quests gating this line (stats epochs, the bounded stats aggregate it requires, and retracted demand release) moved under it, and the OBS half of GStreamer and OBS moved to m1 @@ -80,7 +77,6 @@ This README owns: ## Required -- [Retracted demand release](/quest/m0/broadcast-epoch/unannounce-demand-release.md) - a retracted broadcast's track demand is released when its last subscriber leaves, as before #4741 - [Origin](/quest/m0/broadcast-epoch/origin.md) - moq-net publish mints an epoch, consumers follow the newest live one, and bare requests resolve to it on every version - [Apps](/quest/m0/broadcast-epoch/apps.md) - moq-cli, the browser publish and watch components, and demo/web publish under epochs and play bare names - [Gateways](/quest/m0/broadcast-epoch/gateways.md) - RTMP, SRT, and WHIP ingest mint an epoch per incoming connection, so an encoder reconnect is a clean takeover diff --git a/quest/m0/broadcast-epoch/unannounce-demand-release.md b/quest/m0/broadcast-epoch/unannounce-demand-release.md deleted file mode 100644 index 8a5b7002f5..0000000000 --- a/quest/m0/broadcast-epoch/unannounce-demand-release.md +++ /dev/null @@ -1,78 +0,0 @@ -# [S] A retracted broadcast still releases its tracks' demand - -## Goal - -Once a broadcast is unannounced and the retraction has settled, dropping the -last subscriber of one of its tracks resolves `track.demand().unused()` on the -publisher, as it does before the retraction and on `release`. Today the demand -stays used forever while the producer lives. - -In a relay this keeps every demand-driven emit loop (stats, overlays) and its -upstream subscription running after an unannounce, for nobody. moq.pro's -billing-meter test `a_reannounced_node_is_read_once` fails on `main` because of -it. `release` doesn't have the regression, so this lands before the next -release cut. - -## Plan - -Regression from [#4741](https://github.com/moq-dev/moq/pull/4741): the repro -below passes on its parent (16b1fe2da) and on `release` (3492aebf4), and fails -at 0382d309f and on `main` 2704e10e2. -The unsettled case passes, so the hold comes from the origin processing the -retraction while the reader is still subscribed. - -Suspect, unverified: a front's `Action::End` in `rs/moq-net/src/model/origin.rs` -keeps a used track's copy in flight ("readers follow the copy they read to its -end"). The copy then outlives its readers and keeps the source track -subscribed until the publisher ends it, which a live producer never does. A -copy kept in flight after a retraction should still be released when its last -reader leaves. - -Land the repro as the regression test (moq-tokio, paused time; `control` and -`unannounced` pass on `main`, `unannounced_settled` hits the 600 s virtual -timeout): - -```rust -use std::time::Duration; - -async fn case(unannounce: bool, settle: bool) { - let origin = moq_tokio::origin::spawn(); - let broadcast = origin.publish("a/b", moq_net::origin::Route::default()).unwrap(); - let track = broadcast.create_track("t", None).unwrap(); - let consumer = origin.consume().request_broadcast("a/b").await.unwrap(); - let sub = consumer.track("t").unwrap().subscribe(None).await.unwrap(); - track.demand().used().await.unwrap(); - if unannounce { - broadcast.unannounce(); - if settle { - tokio::time::sleep(Duration::from_secs(1)).await; - } - } - drop(sub); - drop(consumer); - let demand = track.demand(); - let unused = demand.unused(); - tokio::time::timeout(Duration::from_secs(600), unused).await.expect("released").unwrap(); -} - -#[tokio::test(start_paused = true)] -async fn control() { case(false, false).await } -#[tokio::test(start_paused = true)] -async fn unannounced() { case(true, false).await } -#[tokio::test(start_paused = true)] -async fn unannounced_settled() { case(true, true).await } -``` - -Also cover a remote source, the case a relay hits (`rs/moq-net/tests`, the mock -harness, every version): publisher origin P announces the broadcast, relay R -pulls it over `connect_mock`, and a reader subscribes to the track on R. P -unannounces and the retraction settles. Drop R's reader, then assert P's -`track.demand().unused()` resolves, which shows R's upstream SUBSCRIBE was -cancelled. If -[Request linger](/quest/m1/request-linger.md) lands first, the release waits -out its linger, not forever. - -## Related - -- [Request linger](/quest/m1/request-linger.md) - bounds how long an abandoned upstream request may outlive its readers -- [Broadcast epochs](/quest/m0/broadcast-epoch/README.md) - gates the release #4741 ships in diff --git a/rs/moq-net/src/model/origin.rs b/rs/moq-net/src/model/origin.rs index 00cd050446..ba9bade4f9 100644 --- a/rs/moq-net/src/model/origin.rs +++ b/rs/moq-net/src/model/origin.rs @@ -2300,10 +2300,34 @@ impl TrackIo { } } -/// Drives one front: feeds the world's events to a [`Front`] and performs the -/// actions it returns, until the front ends. The decisions live in the machine; -/// this only waits and executes, so nothing here decides anything twice. -async fn run_front(task: FrontTask) { +/// Drives one front until it ends, then holds the tracks it left in flight until +/// their last reader leaves, or until nothing owns the origin: a reader never keeps +/// the driver from finishing. +async fn run_front(task: FrontTask, origin: TasksWeak) { + let mut in_flight = serve_front(task).await; + // Each keeps its copy for the readers still on their way, then lets go as an unread + // track parks, so the copy never keeps its source subscribed for nobody. + kio::wait(|waiter| { + if origin.poll_orphaned(waiter).is_ready() { + return Poll::Ready(()); + } + in_flight.retain(|io| { + io.weak.poll_unused(waiter); + io.weak.is_used() + }); + match in_flight.is_empty() { + true => Poll::Ready(()), + false => Poll::Pending, + } + }) + .await; +} + +/// Feeds the world's events to a [`Front`] and performs the actions it returns, +/// until the front ends; returns the tracks it ended in flight. The decisions live +/// in the machine; this only waits and executes, so nothing here decides anything +/// twice. +async fn serve_front(task: FrontTask) -> Vec { let FrontTask { shared, broadcast, @@ -2567,6 +2591,7 @@ async fn run_front(task: FrontTask) { // subscriptions already in flight): their readers follow the copy // they read to its end, since no front is left to replace it. broadcast.close(); + let mut in_flight = Vec::new(); for (_, mut io) in tracks.drain() { let used = io.weak.is_used(); // A reader still waiting on its source's answer is in flight @@ -2584,10 +2609,12 @@ async fn run_front(task: FrontTask) { true => Ok(()), false => Err(err.clone()), }); + continue; } - // Dropping `io.routes` concludes the rest: readers follow the copy. + io.routes.conclude(); + in_flight.push(io); } - return; + return in_flight; } } } @@ -4300,15 +4327,18 @@ impl Consumer { // Released before the push: a set whose handles are gone drops the task, // and the `Watch` it carries unregisters under this same lock. drop(state); - self.tasks.push(run_front(FrontTask { - shared: self.shared.clone(), - broadcast, - path: absolute, - horizon: self.horizon, - watch, - request, - timers: self.timers.clone(), - })); + self.tasks.push(run_front( + FrontTask { + shared: self.shared.clone(), + broadcast, + path: absolute, + horizon: self.horizon, + watch, + request, + timers: self.timers.clone(), + }, + self.tasks.clone(), + )); kio::Pending::new(Requesting::queued(consumer).with_path(requested).with_stats(scope)) } @@ -8298,6 +8328,30 @@ mod tests { drop(consumer); } + /// A reader still on a retracted broadcast's track keeps its front draining, but + /// never keeps the driver from finishing once nothing owns the origin. + #[tokio::test(start_paused = true)] + async fn a_retracted_reader_never_holds_the_driver() { + let (producer, driver) = Producer::new(Config::new(origin(1))); + let run = tokio::spawn(crate::time::run(driver)); + let consumer = producer.consume(); + let broadcast = producer.publish("room/alice", Route::default()).unwrap(); + let track = broadcast.create_track("video", None).unwrap(); + let resolved = consumer.request_broadcast("room/alice").await.expect("resolves"); + let sub = resolved.track("video").unwrap().subscribe(None).await.unwrap(); + + broadcast.unannounce(); + settle(|| resolved.is_closed()).await; + drop(broadcast); + drop(producer); + tokio::time::timeout(Duration::from_secs(5), run) + .await + .expect("driver must finish once the producers are gone") + .unwrap(); + drop(sub); + drop(track); + } + #[test] fn watch_wakes_only_for_covering_changes() { let producer = origin(1).produce(); diff --git a/rs/moq-net/src/model/resume.rs b/rs/moq-net/src/model/resume.rs index 048f292920..641dd60d8c 100644 --- a/rs/moq-net/src/model/resume.rs +++ b/rs/moq-net/src/model/resume.rs @@ -31,10 +31,14 @@ const MAX_DELIVERED: usize = 1024; struct Route { /// Bumped whenever `copy` changes, so a reader can tell a replacement apart. generation: u64, - /// The serving route's copy; `None` while nobody reads the track, or between routes. + /// The serving route's copy; `None` while nobody reads the track, between routes, or + /// once the front let go. copy: Option, /// How the track ended, once the front decided. end: Option>, + /// No front is left to replace the copy, though one may still hold it for readers + /// on their way: readers follow it to its end. + concluded: bool, } /// Every reader of a logical track, so a new copy subscribes them all as it is served: @@ -43,8 +47,9 @@ type Readers = Arc>>>>; /// The front's side of a logical track: which copy serves it, and how it ends. /// -/// Dropping it concludes the track: readers follow the last copy to its end, since no -/// front is left to replace it. +/// Dropping it concludes the track and lets go of the serving copy: readers already on +/// it follow it to its end, and nothing else keeps its route subscribed, however long +/// the logical track's state stays allocated. pub(crate) struct Producer { state: kio::Producer, readers: Readers, @@ -93,6 +98,14 @@ impl Producer { } } + /// No front will replace the serving copy: readers follow it to its end, while this + /// still holds it for the readers on their way. + pub(crate) fn conclude(&self) { + if let Ok(mut route) = self.state.write() { + route.concluded = true; + } + } + pub(crate) fn consume(&self) -> Consumer { Consumer { state: self.state.consume(), @@ -101,6 +114,15 @@ impl Producer { } } +impl Drop for Producer { + fn drop(&mut self) { + // Same generation: a reader on the copy keeps it as the serving one. + if let Ok(mut route) = self.state.write() { + route.copy = None; + } + } +} + /// A reader's handle on a logical track's [`Producer`]. #[derive(Clone)] pub(crate) struct Consumer { @@ -161,7 +183,7 @@ impl Consumer { /// Poll for the route state to move past `generation`, or the track to end. fn poll_changed(&self, generation: u64, waiter: &kio::Waiter) -> Poll<()> { match self.state.poll(waiter, |route| { - match route.generation != generation || route.end.is_some() { + match route.generation != generation || route.end.is_some() || route.concluded { true => Poll::Ready(()), false => Poll::Pending, } @@ -176,7 +198,7 @@ impl Consumer { /// decision. fn end(&self) -> Option>> { let route = self.state.read(); - match (&route.end, route.is_closed()) { + match (&route.end, route.concluded || route.is_closed()) { (Some(end), _) => Some(Some(end.clone())), (None, true) => Some(None), (None, false) => None, @@ -820,7 +842,7 @@ impl Recover { route.generation, route.copy.clone(), route.end.clone(), - route.is_closed(), + route.concluded || route.is_closed(), ) }; // Registered before looking, so a route change from here on wakes the reader. diff --git a/rs/moq-net/src/util.rs b/rs/moq-net/src/util.rs index 988a3069ed..c1eee58489 100644 --- a/rs/moq-net/src/util.rs +++ b/rs/moq-net/src/util.rs @@ -128,6 +128,12 @@ impl TasksWeak { state.queued.push_back(task.maybe_boxed()); } } + + /// Poll for every owning handle being gone, after which the set finishes as soon + /// as its children do. + pub fn poll_orphaned(&self, waiter: &kio::Waiter) -> Poll<()> { + self.alive.poll_closed(waiter) + } } /// A dynamic set of child tasks polled by its parent driver, built on diff --git a/rs/moq-net/tests/unannounce_release.rs b/rs/moq-net/tests/unannounce_release.rs new file mode 100644 index 0000000000..b420fa3bd0 --- /dev/null +++ b/rs/moq-net/tests/unannounce_release.rs @@ -0,0 +1,50 @@ +//! A relay reading a broadcast that its publisher unannounced still cancels its upstream +//! subscription once its last reader leaves, so the publisher's demand goes unused. + +mod support; + +use std::time::Duration; + +use moq_net::{Hop, Version}; +use support::harness::{MockConnectOptions, connect_mock}; + +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 release(version: Version) { + let publisher = produce_origin(1); + let relay = produce_origin(2); + let mut options = MockConnectOptions::new(version); + options.server_publish = Some(publisher.consume()); + options.client_subscribe = Some(relay.clone()); + let _pair = connect_mock(options).await; + + let broadcast = publisher.publish("a/b", moq_net::origin::Route::default()).unwrap(); + let track = broadcast.create_track("t", None).unwrap(); + + let consumer = relay.consume(); + consumer.routed("a/b").await.expect("routed"); + let remote = consumer.request_broadcast("a/b").await.expect("resolves"); + let sub = remote.track("t").unwrap().subscribe(None).await.expect("subscribes"); + track.demand().used().await.unwrap(); + + broadcast.unannounce(); + tokio::time::sleep(Duration::from_secs(1)).await; + + drop(sub); + drop(remote); + tokio::time::timeout(Duration::from_secs(600), track.demand().unused()) + .await + .unwrap_or_else(|_| panic!("{version}: the relay kept its upstream subscription")) + .unwrap(); +} + +#[tokio::test(start_paused = true)] +async fn unannounced_remote_track_releases_demand() { + for name in Version::names() { + release(name.parse().unwrap()).await; + } +} diff --git a/rs/moq-tokio/tests/unannounce.rs b/rs/moq-tokio/tests/unannounce.rs new file mode 100644 index 0000000000..9e27bae341 --- /dev/null +++ b/rs/moq-tokio/tests/unannounce.rs @@ -0,0 +1,47 @@ +//! Integration test: an unannounced broadcast still releases a track's demand +//! once its last subscriber leaves. + +use std::time::Duration; + +use moq_tokio::moq_net; + +/// Subscribe to a published track, optionally unannounce (and let the origin +/// process the retraction), then drop the subscriber and wait for the +/// publisher's demand to go unused. +async fn release(unannounce: bool, settle: bool) { + let origin = moq_tokio::origin::spawn(); + let broadcast = origin.publish("a/b", moq_net::origin::Route::default()).unwrap(); + let track = broadcast.create_track("t", None).unwrap(); + let consumer = origin.consume().request_broadcast("a/b").await.unwrap(); + let sub = consumer.track("t").unwrap().subscribe(None).await.unwrap(); + track.demand().used().await.unwrap(); + + if unannounce { + broadcast.unannounce(); + if settle { + tokio::time::sleep(Duration::from_secs(1)).await; + } + } + + drop(sub); + drop(consumer); + tokio::time::timeout(Duration::from_secs(600), track.demand().unused()) + .await + .expect("demand released") + .unwrap(); +} + +#[tokio::test(start_paused = true)] +async fn announced_track_releases_demand() { + release(false, false).await +} + +#[tokio::test(start_paused = true)] +async fn unannounced_track_releases_demand() { + release(true, false).await +} + +#[tokio::test(start_paused = true)] +async fn settled_unannounce_releases_demand() { + release(true, true).await +} From d417e6e82f7a7c3283ec21d9117e41d63685b5a5 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 5 Oct 2026 18:05:08 -0700 Subject: [PATCH 3/4] fix(moq-net): keep a retracted copy for pending readers past the origin's end Clearing the serving copy in resume::Producer's Drop broke a subscriber that asked before the retraction but was not polled yet: once the origin was orphaned, the tail dropped its tracks and the subscriber read Error::Dropped. Copy-clearing moves into an explicit resume::Producer::release, called only when a track's last reader leaves (the tail) or when the front ends with a track unread and not parked yet (Action::End). The orphan exit just drops. Co-Authored-By: Claude Opus 5.5 --- rs/moq-net/src/model/origin.rs | 69 ++++++++++++++++++++++++++++++++-- rs/moq-net/src/model/resume.rs | 24 ++++++------ 2 files changed, 78 insertions(+), 15 deletions(-) diff --git a/rs/moq-net/src/model/origin.rs b/rs/moq-net/src/model/origin.rs index ba9bade4f9..8f56afcce1 100644 --- a/rs/moq-net/src/model/origin.rs +++ b/rs/moq-net/src/model/origin.rs @@ -2311,10 +2311,13 @@ async fn run_front(task: FrontTask, origin: TasksWeak) { if origin.poll_orphaned(waiter).is_ready() { return Poll::Ready(()); } - in_flight.retain(|io| { + for io in std::mem::take(&mut in_flight) { io.weak.poll_unused(waiter); - io.weak.is_used() - }); + match io.weak.is_used() { + true => in_flight.push(io), + false => io.routes.release(), + } + } match in_flight.is_empty() { true => Poll::Ready(()), false => Poll::Pending, @@ -2609,6 +2612,11 @@ async fn serve_front(task: FrontTask) -> Vec { true => Ok(()), false => Err(err.clone()), }); + // An unread track the front has not parked yet still holds + // its copy, which would keep its source subscribed. + if !used { + io.routes.release(); + } continue; } io.routes.conclude(); @@ -8352,6 +8360,61 @@ mod tests { drop(track); } + /// A subscriber that asked before the retraction but was not polled yet still finds + /// the copy after the origin is orphaned and its driver finishes. + #[tokio::test(start_paused = true)] + async fn an_orphaned_front_keeps_the_copy_for_a_pending_subscriber() { + let (producer, driver) = Producer::new(Config::new(origin(1))); + let run = tokio::spawn(crate::time::run(driver)); + let consumer = producer.consume(); + let broadcast = producer.publish("room/alice", Route::default()).unwrap(); + let mut track = broadcast.create_track("video", None).unwrap(); + let resolved = consumer.request_broadcast("room/alice").await.expect("resolves"); + let sub = resolved.track("video").unwrap().subscribe(None).await.unwrap(); + let pending = resolved.track("video").unwrap().subscribe(None); + + broadcast.unannounce(); + settle(|| resolved.is_closed()).await; + drop(broadcast); + drop(producer); + tokio::time::timeout(Duration::from_secs(5), run) + .await + .expect("driver must finish once the producers are gone") + .unwrap(); + track.write_frame(crate::Timestamp::ZERO, b"late".as_ref()).unwrap(); + let mut late = pending.await.expect("subscribes"); + let mut group = late.recv_group().await.expect("recv").expect("the source's group"); + assert_eq!(&group.read_frame().await.unwrap().unwrap().payload[..], b"late"); + drop(sub); + } + + /// A track whose last reader leaves as its broadcast ends, before the front parks it, + /// lets go of its copy: the source is not kept subscribed for nobody. + #[tokio::test(start_paused = true)] + async fn a_track_unread_as_its_broadcast_ends_releases_its_source() { + let producer = origin(1).produce(); + let consumer = producer.consume(); + let broadcast = producer.publish("room/alice", Route::default()).unwrap(); + let track = broadcast.create_track("video", None).unwrap(); + let resolved = consumer.request_broadcast("room/alice").await.expect("resolves"); + let sub = resolved.track("video").unwrap().subscribe(None).await.unwrap(); + tokio::time::timeout(Duration::from_secs(1), track.demand().used()) + .await + .expect("the front subscribed the source") + .unwrap(); + + // The source closing wakes the front before the reader leaving does, so it ends + // with the track unread and its copy not parked yet. + drop(sub); + drop(broadcast); + settle(|| resolved.is_closed()).await; + + tokio::time::timeout(Duration::from_secs(5), track.demand().unused()) + .await + .expect("the ended copy must not keep the source subscribed") + .unwrap(); + } + #[test] fn watch_wakes_only_for_covering_changes() { let producer = origin(1).produce(); diff --git a/rs/moq-net/src/model/resume.rs b/rs/moq-net/src/model/resume.rs index 641dd60d8c..4ea64199d3 100644 --- a/rs/moq-net/src/model/resume.rs +++ b/rs/moq-net/src/model/resume.rs @@ -47,9 +47,10 @@ type Readers = Arc>>>>; /// The front's side of a logical track: which copy serves it, and how it ends. /// -/// Dropping it concludes the track and lets go of the serving copy: readers already on -/// it follow it to its end, and nothing else keeps its route subscribed, however long -/// the logical track's state stays allocated. +/// Dropping it concludes the track: readers follow the last copy to its end, since no +/// front is left to replace it. [`Producer::release`] also lets go of that copy, so it +/// stops keeping its route subscribed however long the logical track's state stays +/// allocated. pub(crate) struct Producer { state: kio::Producer, readers: Readers, @@ -106,6 +107,14 @@ impl Producer { } } + /// Let go of the serving copy once no front will replace it and nobody reads the + /// track. Same generation: a reader on the copy keeps it as the serving one. + pub(crate) fn release(self) { + if let Ok(mut route) = self.state.write() { + route.copy = None; + } + } + pub(crate) fn consume(&self) -> Consumer { Consumer { state: self.state.consume(), @@ -114,15 +123,6 @@ impl Producer { } } -impl Drop for Producer { - fn drop(&mut self) { - // Same generation: a reader on the copy keeps it as the serving one. - if let Ok(mut route) = self.state.write() { - route.copy = None; - } - } -} - /// A reader's handle on a logical track's [`Producer`]. #[derive(Clone)] pub(crate) struct Consumer { From e53983fd1ca8f15d4b078c767d6bdf1208e30b4f Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 5 Oct 2026 18:57:51 -0700 Subject: [PATCH 4/4] docs(moq-net): note the tail's lookup and orphan invariants Co-Authored-By: Claude Opus 5.5 --- rs/moq-net/src/model/origin.rs | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/rs/moq-net/src/model/origin.rs b/rs/moq-net/src/model/origin.rs index 8f56afcce1..62896865d3 100644 --- a/rs/moq-net/src/model/origin.rs +++ b/rs/moq-net/src/model/origin.rs @@ -2306,8 +2306,12 @@ impl TrackIo { async fn run_front(task: FrontTask, origin: TasksWeak) { let mut in_flight = serve_front(task).await; // Each keeps its copy for the readers still on their way, then lets go as an unread - // track parks, so the copy never keeps its source subscribed for nobody. + // track parks, so the copy never keeps its source subscribed for nobody. A plain + // `is_used` suffices: the front closed its broadcast, which refuses every lookup, so + // no reader can arrive once the last one left. kio::wait(|waiter| { + // Dropped without a release, so a reader on its way keeps the copy; it stays held + // only while the logical track's state does, for an origin nobody owns. if origin.poll_orphaned(waiter).is_ready() { return Poll::Ready(()); } @@ -8403,8 +8407,8 @@ mod tests { .expect("the front subscribed the source") .unwrap(); - // The source closing wakes the front before the reader leaving does, so it ends - // with the track unread and its copy not parked yet. + // The front looks at a closed source before demand edges, so it ends with the + // track unread and its copy not parked yet. drop(sub); drop(broadcast); settle(|| resolved.is_closed()).await;