diff --git a/doc/bin/cli.md b/doc/bin/cli.md index 43f26fdbcd..1d109da204 100644 --- a/doc/bin/cli.md +++ b/doc/bin/cli.md @@ -182,8 +182,8 @@ zero-based `frame` and padded standard base64. `` is the literal track name. `/fetch` splits its path on the last `/`, so the two agree only for names without one. Fetch only dials `--connect`, and refuses a listener or cluster flag. It gives up after 30 seconds, as `/fetch` -does, and exits non-zero when the broadcast or group is not found, the relay -refuses, or the deadline passes. +does, and exits non-zero when the broadcast or group is not found (before +writing anything), the relay refuses, or the deadline passes. ## Multiple stages diff --git a/doc/bin/relay/http.md b/doc/bin/relay/http.md index 28cad8de5c..80460efb30 100644 --- a/doc/bin/relay/http.md +++ b/doc/bin/relay/http.md @@ -13,7 +13,7 @@ operational ones that must stay private. | Endpoint | Returns | | --- | --- | | `GET /announced/` | Broadcasts announced under the prefix. | -| `GET /fetch//?group=N` | One group from the cache, the latest by default. Useful for catch-up and debugging. | +| `GET /fetch//?group=N` | One group from the cache, the latest by default, or `404` if the track has no such group. Useful for catch-up and debugging. | | `GET /certificate.sha256` | The fingerprint of the first configured TLS certificate, for pinning a self-signed dev certificate. | | `GET /health` | `200 ok`, unauthenticated, for load balancers. | diff --git a/quest/m1/README.md b/quest/m1/README.md index fee11fb135..87ac03950a 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -17,7 +17,6 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. ## Quests -- [Missing fetch group](/quest/m1/fetch-missing-group.md) - HTTP /fetch answers 404 and `moq fetch` fails cleanly for a group the track lacks - [JS fetch answer](/quest/m1/js-fetch-answer.md) - js/net's lite fetch settles on the publisher's answer, and a JS publisher's miss resets with NotFound - [libmoq hidden opt-in](/quest/m1/libmoq-hidden.md) - `moq_origin_announced` takes a `hidden` flag so C callers can list `.`-named broadcasts - [Announce compression](/quest/m1/announce-compression.md) - a lite-07 announce reuses the path head and hop-chain tail of a live announcement on its stream instead of resending them diff --git a/quest/m1/fetch-missing-group.md b/quest/m1/fetch-missing-group.md deleted file mode 100644 index e144de69aa..0000000000 --- a/quest/m1/fetch-missing-group.md +++ /dev/null @@ -1,18 +0,0 @@ -# [S] Fetching a missing group fails before any output - -## Goal - -HTTP `/fetch//?group=N` answers 404 for a group the track -does not have, and `moq fetch --group N` exits with a clean "not found" -before writing anything. Today the lookup hands back a group consumer that -only fails on its first frame read, so the endpoint answers 200 with a -cut-off body. - -## Plan - -- Resolve whether the group exists before the response starts, in the - shared lookup the relay and the CLI both use. Where the check belongs (the - relay helper or a `moq-net` track consumer method) is the implementer's - call; propose it in the PR. -- Test: a missing group is a 404 over HTTP and a non-zero `moq fetch` exit - with no stdout; an existing group is byte-identical to today. diff --git a/rs/moq-cli/src/fetch.rs b/rs/moq-cli/src/fetch.rs index 11cfc90258..fad7b7862b 100644 --- a/rs/moq-cli/src/fetch.rs +++ b/rs/moq-cli/src/fetch.rs @@ -86,7 +86,6 @@ async fn fetch( None => format!("no group of `{}` found", args.track), })?; - // A missing sequence can resolve to a group that fails on its first read. let sequence = group.sequence; let mut index = 0; while let Some(frame) = group @@ -209,13 +208,13 @@ mod tests { (result, out) } - /// The relay's HTTP `/fetch` body for the same group. - async fn curl(&self, query: &str) -> Vec { + /// The relay's HTTP `/fetch` status and body for the same group. + async fn curl(&self, query: &str) -> (u16, Vec) { let response = reqwest::get(format!("http://{}/fetch/demo/data{query}", self.http)) .await .expect("HTTP fetch"); - assert_eq!(response.status(), 200); - response.bytes().await.expect("HTTP body").to_vec() + let status = response.status().as_u16(); + (status, response.bytes().await.expect("HTTP body").to_vec()) } } @@ -254,7 +253,7 @@ mod tests { let (result, out) = fixture.fetch(&["data", "--group", "1"], TIMEOUT).await; result.expect("fetch"); assert_eq!(out, frames(1).concat()); - assert_eq!(out, fixture.curl("?group=1").await); + assert_eq!(fixture.curl("?group=1").await, (200, out)); } /// No `--group` reads the newest group, as `/fetch` does by default. @@ -266,9 +265,11 @@ mod tests { let (result, out) = fixture.fetch(&["data"], TIMEOUT).await; result.expect("fetch"); assert_eq!(out, frames(2).concat()); - assert_eq!(out, fixture.curl("").await); + assert_eq!(fixture.curl("").await, (200, out)); } + /// A missing sequence fails the lookup itself, before any output, as `/fetch` + /// answers 404 rather than starting a body. #[tokio::test] async fn a_missing_sequence_fails() { let _env = EnvGuard::clear(ENV); @@ -276,9 +277,9 @@ mod tests { let (result, out) = fixture.fetch(&["data", "--group", "99"], TIMEOUT).await; let err = result.expect_err("group 99 does not exist"); - let err = format!("{err:#}"); - assert!(err.contains("group 99") && err.contains("not found"), "{err}"); + assert_eq!(err.to_string(), "group 99 of `data` not found", "{err:#}"); assert!(out.is_empty()); + assert_eq!(fixture.curl("?group=99").await, (404, Vec::new())); } /// A track with no group never resolves "newest", so the deadline ends it. diff --git a/rs/moq-net/src/coding/reader.rs b/rs/moq-net/src/coding/reader.rs index 2005bd11ae..12bf76bf5f 100644 --- a/rs/moq-net/src/coding/reader.rs +++ b/rs/moq-net/src/coding/reader.rs @@ -227,8 +227,8 @@ impl Reader { Poll::Ready(Ok(())) } - /// Poll for whether data is available in the buffer or stream. - fn poll_has_more(&mut self, cx: &mut Context<'_>) -> Poll> { + /// Poll for whether data is available in the buffer or stream: `false` once it finishes. + pub(crate) fn poll_has_more(&mut self, cx: &mut Context<'_>) -> Poll> { if !self.buffer.is_empty() { return Poll::Ready(Ok(true)); } diff --git a/rs/moq-net/src/lite/subscriber.rs b/rs/moq-net/src/lite/subscriber.rs index 46b90b0d5c..67a79c0821 100644 --- a/rs/moq-net/src/lite/subscriber.rs +++ b/rs/moq-net/src/lite/subscriber.rs @@ -3503,6 +3503,12 @@ enum FetchRunState { stream: Stream, frame_start: u64, }, + /// Flushed; waiting for the publisher to answer before accepting. + Answer { + request: Option, + stream: Stream, + frame_start: u64, + }, Ingest { stream: Stream, producer: group::Producer, @@ -3613,7 +3619,33 @@ impl kio::Task for FetchServeRun { else { unreachable!() }; + self.state = FetchRunState::Answer { + request, + stream, + frame_start, + }; + } + FetchRunState::Answer { stream, .. } => { + // Lite has no FETCH_OK: a publisher without the group resets the + // stream instead. Accepting before the first byte (or a FIN, for an + // empty group) would resolve every joined `fetch_group` to a group + // that only fails on its first read, so wait for the answer. + let answered = ready!(stream.reader.poll_has_more(&mut cx)); + let FetchRunState::Answer { + request, + stream, + frame_start, + } = std::mem::replace(&mut self.state, FetchRunState::Done) + else { + unreachable!() + }; let request = request.expect("request pending"); + if let Err(err) = answered { + tracing::debug!(track = %self.serve.name, group = self.group, %err, "fetch refused"); + stream.writer.abort(&err); + request.reject(err); + return Poll::Ready(()); + } // Make the group available (resolving the downstream fetch) and fill // it. The track::Info only takes effect if the track isn't accepted yet diff --git a/rs/moq-net/src/model/group.rs b/rs/moq-net/src/model/group.rs index 0205ae10b4..e9f3a4bbda 100644 --- a/rs/moq-net/src/model/group.rs +++ b/rs/moq-net/src/model/group.rs @@ -1804,8 +1804,10 @@ impl Fetch { /// The handler fulfills it by calling [`Self::accept`], which inserts the group /// into the track cache (resolving every [`track::Consumer::fetch_group`] that joined the /// attempt) and returns a [`Producer`] to fill. A relay typically opens a wire -/// FETCH, reads FETCH_OK, then accepts. The request carries its own producer handle, -/// so it works the same whether or not the track has been accepted yet. +/// FETCH and waits for the publisher to answer before accepting, so a group the +/// publisher lacks is rejected rather than accepted and then aborted. The request +/// carries its own producer handle, so it works the same whether or not the track +/// has been accepted yet. pub struct Request { pub(crate) state: kio::Producer, pub(crate) fetch: kio::Shared, diff --git a/rs/moq-net/src/model/track.rs b/rs/moq-net/src/model/track.rs index 6cfda585dc..8ac10f45b9 100644 --- a/rs/moq-net/src/model/track.rs +++ b/rs/moq-net/src/model/track.rs @@ -2609,8 +2609,9 @@ impl Consumer { /// or `group::Fetch::default()`. /// /// The returned future resolves to [`Error::NotFound`] when the group can never be served - /// (past the final sequence, or no [`Dynamic`] on the track), or the track's abort error - /// if it's already closed. Concurrent fetches for the same sequence coalesce onto one + /// (past the final sequence, or no [`Dynamic`] on the track), the handler's rejection + /// (a relay's upstream miss is [`StreamError::NotFound`](crate::StreamError::NotFound)), + /// or the track's abort error if it's already closed. Concurrent fetches for the same sequence coalesce onto one /// handler request. pub fn fetch_group(&self, sequence: u64, options: impl Into>) -> kio::Pending { let options = options.into().unwrap_or_default(); diff --git a/rs/moq-relay/src/fetch.rs b/rs/moq-relay/src/fetch.rs index 94dfb063ca..5cc958e8bd 100644 --- a/rs/moq-relay/src/fetch.rs +++ b/rs/moq-relay/src/fetch.rs @@ -3,8 +3,9 @@ /// The newest needs a live subscription to learn its sequence, since a fetch can /// only retrieve a sequence already known. Once known it is fetched rather than /// read off the subscription, so an evicted group is retrieved from upstream -/// instead of waited on forever. Returns [`moq_net::Error::NotFound`] when the -/// group can never be served, including a track that ends before any group. +/// instead of waited on forever. Fails before returning when the group can never +/// be served: [`moq_net::Error::NotFound`] locally, including a track that ends +/// before any group, or [`moq_net::StreamError::NotFound`] from upstream. pub async fn fetch_group( track: &moq_net::track::Consumer, sequence: Option, diff --git a/rs/moq-relay/src/web.rs b/rs/moq-relay/src/web.rs index f948c49fc0..9f4bb55e50 100644 --- a/rs/moq-relay/src/web.rs +++ b/rs/moq-relay/src/web.rs @@ -926,7 +926,10 @@ async fn serve_fetch( }; let group = match async { crate::fetch_group(&broadcast.track(&track)?, sequence).await }.await { Ok(group) => group, - Err(moq_net::Error::NotFound) => return Err(StatusCode::NOT_FOUND), + // A miss upstream arrives as the stream reset that refused the FETCH. + Err(moq_net::Error::NotFound | moq_net::Error::Stream(moq_net::StreamError::NotFound)) => { + return Err(StatusCode::NOT_FOUND); + } Err(_) => return Err(StatusCode::INTERNAL_SERVER_ERROR), };