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
4 changes: 2 additions & 2 deletions doc/bin/cli.md
Original file line number Diff line number Diff line change
Expand Up @@ -182,8 +182,8 @@ zero-based `frame` and padded standard base64.
`<track>` 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

Expand Down
2 changes: 1 addition & 1 deletion doc/bin/relay/http.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ operational ones that must stay private.
| Endpoint | Returns |
| --- | --- |
| `GET /announced/<prefix>` | Broadcasts announced under the prefix. |
| `GET /fetch/<broadcast>/<track>?group=N` | One group from the cache, the latest by default. Useful for catch-up and debugging. |
| `GET /fetch/<broadcast>/<track>?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. |

Expand Down
1 change: 0 additions & 1 deletion quest/m1/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
18 changes: 0 additions & 18 deletions quest/m1/fetch-missing-group.md

This file was deleted.

19 changes: 10 additions & 9 deletions rs/moq-cli/src/fetch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<u8> {
/// The relay's HTTP `/fetch` status and body for the same group.
async fn curl(&self, query: &str) -> (u16, Vec<u8>) {
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())
}
}

Expand Down Expand Up @@ -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.
Expand All @@ -266,19 +265,21 @@ 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);
let fixture = Fixture::new().await;

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.
Expand Down
4 changes: 2 additions & 2 deletions rs/moq-net/src/coding/reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -227,8 +227,8 @@ impl<S: crate::transport::poll::RecvStream, V: StreamCodes> Reader<S, V> {
Poll::Ready(Ok(()))
}

/// Poll for whether data is available in the buffer or stream.
fn poll_has_more(&mut self, cx: &mut Context<'_>) -> Poll<Result<bool, Error>> {
/// 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<Result<bool, Error>> {
if !self.buffer.is_empty() {
return Poll::Ready(Ok(true));
}
Expand Down
32 changes: 32 additions & 0 deletions rs/moq-net/src/lite/subscriber.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3503,6 +3503,12 @@ enum FetchRunState<S: crate::transport::poll::Session> {
stream: Stream<S, Version>,
frame_start: u64,
},
/// Flushed; waiting for the publisher to answer before accepting.
Answer {
request: Option<group::Request>,
stream: Stream<S, Version>,
frame_start: u64,
},
Ingest {
stream: Stream<S, Version>,
producer: group::Producer,
Expand Down Expand Up @@ -3613,7 +3619,33 @@ impl<S: crate::transport::poll::Session> kio::Task for FetchServeRun<S> {
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));
Comment on lines +3629 to +3633

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Gate TypeScript fetches on the publisher answer

For a TypeScript/browser consumer fetching a missing lite group, js/net/src/lite/subscriber.ts:706-767 still creates and returns a group mirror immediately after sending FETCH; its background #runFetchResponse only closes that group when the reset arrives. Consequently, fetchGroup() continues to resolve and the failure appears on the first frame read, which is exactly the behavior this patch fixes only in Rust. Apply the equivalent first-byte/FIN gate and missing-group regression test in js/net so the APIs remain mirrored.

AGENTS.md reference: AGENTS.md:L94-L98

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Leaving the JS gate out of this PR. It is the Rust side: lite fetch_group waits for the publisher's first byte or FIN, and HTTP /fetch maps an upstream miss to 404. The browser gap is the separate quest landed in #4175 (quest/m1/js-fetch-answer.md), not this diff.

(Written by Grok 4.7)

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
Expand Down
6 changes: 4 additions & 2 deletions rs/moq-net/src/model/group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<track::TrackState>,
pub(crate) fetch: kio::Shared<track::FetchState>,
Expand Down
5 changes: 3 additions & 2 deletions rs/moq-net/src/model/track.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Option<group::Fetch>>) -> kio::Pending<Fetching> {
let options = options.into().unwrap_or_default();
Expand Down
5 changes: 3 additions & 2 deletions rs/moq-relay/src/fetch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<u64>,
Expand Down
5 changes: 4 additions & 1 deletion rs/moq-relay/src/web.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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),
};

Expand Down
Loading