Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion doc/concept/moq-lite.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
4 changes: 0 additions & 4 deletions quest/m0/broadcast-epoch/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
78 changes: 0 additions & 78 deletions quest/m0/broadcast-epoch/unannounce-demand-release.md

This file was deleted.

151 changes: 136 additions & 15 deletions rs/moq-net/src/model/origin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2300,10 +2300,41 @@ 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. 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(());
}
for io in std::mem::take(&mut in_flight) {
io.weak.poll_unused(waiter);
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,
}
})
.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<TrackIo> {
let FrontTask {
shared,
broadcast,
Expand Down Expand Up @@ -2567,6 +2598,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
Expand All @@ -2584,10 +2616,17 @@ async fn run_front(task: FrontTask) {
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;
}
// Dropping `io.routes` concludes the rest: readers follow the copy.
io.routes.conclude();
in_flight.push(io);
}
return;
return in_flight;
}
}
}
Expand Down Expand Up @@ -4300,15 +4339,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))
}

Expand Down Expand Up @@ -8298,6 +8340,85 @@ 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);
}

/// 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 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;

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();
Expand Down
32 changes: 27 additions & 5 deletions rs/moq-net/src/model/resume.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<track::Consumer>,
/// How the track ended, once the front decided.
end: Option<Result<()>>,
/// 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:
Expand All @@ -44,7 +48,9 @@ type Readers = Arc<Mutex<Vec<Weak<Mutex<Reader>>>>>;
/// 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.
/// 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<Route>,
readers: Readers,
Expand Down Expand Up @@ -93,6 +99,22 @@ 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;
}
}

/// 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(),
Expand Down Expand Up @@ -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,
}
Expand All @@ -176,7 +198,7 @@ impl Consumer {
/// decision.
fn end(&self) -> Option<Option<Result<()>>> {
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,
Expand Down Expand Up @@ -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.
Expand Down
6 changes: 6 additions & 0 deletions rs/moq-net/src/util.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading