diff --git a/doc/concept/standard.md b/doc/concept/standard.md index a6e9b72476..001d5f96b1 100644 --- a/doc/concept/standard.md +++ b/doc/concept/standard.md @@ -33,11 +33,18 @@ to model priority 127, where higher values are served first. A track that never sets a priority is 127 as well, so it goes out as 128 on IETF and 127 on moq-lite. -On drafts 14–19, the Rust publisher serves relative joining `FETCH` requests -with offset zero for `NextObject` subscriptions. The fetch delivers the saved -current-group prefix, and the subscription delivers later objects. Standalone, -absolute joining, and nonzero-offset fetches are refused. Draft-20 uses -subscription fills instead. JavaScript publishing does not yet serve `FETCH`; +The Rust publisher answers a standalone `FETCH` by walking its range one group +at a time, in ascending order, from the cache. A relay fetches each missing +group upstream with a `FETCH` of that one whole group, and an upstream refusal +is the refusal the fetcher sees. A descending range of several groups is +refused. A standalone `FETCH` carries no timestamps, since no `SUBSCRIBE_OK` +declared a timescale for it. + +On drafts 14–19, the Rust publisher also serves relative and absolute joining +`FETCH` requests for `NextObject` subscriptions: the whole groups before the +subscription's group, then that group's saved prefix, while the subscription +delivers later objects. Draft-20 uses subscription fills instead. JavaScript +publishing does not yet serve `FETCH`; Rust and JavaScript subscribers request unfiltered delivery on older drafts because they do not issue joining fetches. Other publishers may replay a cached backlog for that filter; selecting the next group instead would leave static diff --git a/quest/m1/README.md b/quest/m1/README.md index bf58cf9600..eea22d03d4 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -40,7 +40,9 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. - [Publish delay](/quest/m1/publish-delay.md) - js/publish encoders advertise `delay` behind the earliest rendition, like moq-mux - [Data jitter](/quest/m1/data-jitter.md) - JSON and binary tracks with a capture time advertise a detected `delay` and `jitter` +- [Subgroup refusal](/quest/m1/ietf-subgroup-refusal.md) - a non-zero subgroup stream from a moq-transport peer ends that stream, never the session - [Moxygen compatibility](/quest/m1/moxygen/README.md) - one subgroup per group, whole-group FETCH, and one datagram per group, never a full moxygen pass +- [Fetch without SUBSCRIBE](/quest/m1/ietf-fetch-only.md) - a relay fetches from an IETF upstream without subscribing, finished tracks included, with End of Track always reported - [JavaScript FETCH](/quest/m1/js-fetch.md) - generic on-demand group serving and IETF FETCH for browser publishers - [Archive](/quest/m1/archive/README.md) - record selected tracks to any object_store and replay them over FETCH or derived HLS, on the catalog and store the release ships - [Wildcard](/quest/m1/wildcard/README.md) - a relay resolves subscriptions against advertised prefixes, a service claims the prefix it could serve and refuses the rest instead of enumerating broadcasts, and the browser player treats a covering claim as availability diff --git a/quest/m1/ietf-fetch-only.md b/quest/m1/ietf-fetch-only.md new file mode 100644 index 0000000000..c113635ef7 --- /dev/null +++ b/quest/m1/ietf-fetch-only.md @@ -0,0 +1,30 @@ +# [M] Fetch without SUBSCRIBE + +## Goal + +A relay serves a FETCH-only demand for an IETF upstream track without +SUBSCRIBE upstream. A finished upstream track can still be fetched, and a +FETCH_OK that reaches the end of the track says so every time. + +## Plan + +The origin splices a route only once the track's info is known. moq-lite gets +it from TRACK_INFO, so its fetch-only demand never subscribes. IETF has no +TRACK_INFO, so today fetch-only demand still makes the relay SUBSCRIBE upstream. +Two symptoms follow: + +- A finished upstream track refuses the SUBSCRIBE, so its groups cannot be + fetched either. +- The live subscription races the group FETCHes. End of Track is known only + once an upstream FETCH_OK reports it, so a downstream FETCH that runs to the + end can answer before that and leave End of Track unset. moxygen's "FETCH + with large objects" case flakes on this. + +TRACK_STATUS is the likely source of the info. The publisher refuses it today, +so both sides are in scope. Keep the SUBSCRIBE path for real subscription +demand. + +## Related + +- [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line whose FETCH cases this steadies +- [JavaScript FETCH](/quest/m1/js-fetch.md) - the browser publisher answers these fetches diff --git a/quest/m1/ietf-subgroup-refusal.md b/quest/m1/ietf-subgroup-refusal.md new file mode 100644 index 0000000000..06fd256b07 --- /dev/null +++ b/quest/m1/ietf-subgroup-refusal.md @@ -0,0 +1,22 @@ +# [S] Subgroup refusal stays on the stream + +## Goal + +A peer that sends a non-zero subgroup on moq-transport loses that one stream, +never the session. Every other track on the session keeps flowing. + +## Plan + +Against moxygen's `moqtest_server`, a track with two subgroups per group ended +the relay's upstream session, and the server reconnected. Our side refuses the +stream today. Whether the session ends because of how we refuse it (the reset +code, STOP_SENDING, or the alias state it leaves behind) or because the peer +reacts badly to a correct refusal is not known yet. Reproduce it first. If the +peer is at fault, say so on its tracker and keep a regression test for our side. + +A test with an IETF peer that sends a subgroup 1 stream next to a healthy track +is the check. + +## Related + +- [Moxygen compatibility](/quest/m1/moxygen/README.md) - subgroups stay out of scope; only the blast radius is in diff --git a/quest/m1/moxygen/README.md b/quest/m1/moxygen/README.md index 937f5ece5f..f34ba77296 100644 --- a/quest/m1/moxygen/README.md +++ b/quest/m1/moxygen/README.md @@ -30,7 +30,7 @@ Docs stay inline in the change that makes them stale. No new guide. ## Quests -- [Group FETCH](/quest/m1/moxygen/fetch.md) - an IETF FETCH of whole groups is served from cache or fetched upstream, one group at a time +- [Sparse FETCH ranges](/quest/m1/moxygen/fetch-span.md) - a FETCH costs the groups it returns, not the span of its range - [Datagram groups](/quest/m1/moxygen/datagram.md) - an IETF datagram that is one object in a group arrives as a moq-lite datagram group ## Related diff --git a/quest/m1/moxygen/fetch-span.md b/quest/m1/moxygen/fetch-span.md new file mode 100644 index 0000000000..c326968452 --- /dev/null +++ b/quest/m1/moxygen/fetch-span.md @@ -0,0 +1,27 @@ +# [M] Sparse FETCH ranges + +## Goal + +A FETCH's cost follows the groups it can return, not the span of its range. +A track whose sequences are sparse, or whose history is long gone, answers a +wide range without spinning a worker or probing upstream once per missing +sequence. + +## Plan + +The standalone FETCH walk asks `track::Consumer::fetch_group` for every +sequence from start to the newest group, stepping over each miss. With no +fetch handler, every miss resolves at once, so a range over millions of +missing sequences runs without yielding. With a handler, a relay sends one +upstream FETCH per missing sequence. Either way a single request from a peer +buys work proportional to the newest sequence number. + +Seek the next group the track can serve instead of stepping by one. Where a +relay cannot know which groups its upstream still holds, decide what bound +or refusal is honest rather than probing blindly. Benchmark range span and +present-group count as separate axes. + +## Related + +- [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line whose FETCH walk this bounds +- [Fetch without SUBSCRIBE](/quest/m1/ietf-fetch-only.md) - also changes how a relay reaches upstream for fetches diff --git a/quest/m1/moxygen/fetch.md b/quest/m1/moxygen/fetch.md deleted file mode 100644 index f236d1751f..0000000000 --- a/quest/m1/moxygen/fetch.md +++ /dev/null @@ -1,26 +0,0 @@ -# [L] Group FETCH - -## Goal - -An IETF FETCH that asks for whole groups is answered. A group already in -cache is served from there. A miss fetches that group upstream. Subscribers -still FETCH one group at a time. A publisher may be asked for a range, which -the relay walks one group at a time. A publisher refusal is what the -subscriber sees. - -## Plan - -Standalone FETCH and a non-zero joining FETCH are refused today with -"not supported". Walk `track::Consumer::fetch_group`, one group, then the -next. Do not add an archive. - -A joining FETCH is the same walk for the groups it names. A form the walk -cannot express is still an explicit refusal, not a hang. - -The moxygen FETCH cases that ask for whole groups are the check. The rest of -that suite is not. - -## Related - -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line this belongs to -- [JavaScript FETCH](/quest/m1/js-fetch.md) - the browser publisher that fills an upstream miss diff --git a/rs/moq-net/src/ietf/fetch.rs b/rs/moq-net/src/ietf/fetch.rs index 5987d8efbc..57e7d07fe3 100644 --- a/rs/moq-net/src/ietf/fetch.rs +++ b/rs/moq-net/src/ietf/fetch.rs @@ -167,7 +167,8 @@ impl Message for Fetch<'_> { ); let subscriber_priority = subscriber_priority.unwrap_or(128); - let group_order = group_order.unwrap_or(GroupOrder::Descending); + // No preference: the publisher picks the order. + let group_order = group_order.unwrap_or(GroupOrder::Any); Ok(Self { request_id, diff --git a/rs/moq-net/src/ietf/publisher.rs b/rs/moq-net/src/ietf/publisher.rs index 2648a61b8a..4e0101d533 100644 --- a/rs/moq-net/src/ietf/publisher.rs +++ b/rs/moq-net/src/ietf/publisher.rs @@ -44,6 +44,115 @@ enum FillStep { Done, } +/// Where a fetch object sits relative to the one written before it on the same stream. +#[derive(Clone, Copy)] +enum FetchPrior { + /// The first object on the stream. + None, + /// The next object of the same group. + Same, + /// The first object after this earlier group. + Group(u64), +} + +/// One group of a FETCH answer, read out of the cache so a later eviction cannot +/// truncate what FETCH_OK already promised. +struct FetchedGroup { + sequence: u64, + /// The Object ID of the first frame. + first: u64, + frames: Vec, + /// No frame past the last one read exists: the group ended there, or nothing was + /// written past the cap. + complete: bool, +} + +impl FetchPrior { + /// The prior for the next object of a single-group stream, clearing `first`. + fn next(first: &mut bool) -> Self { + match std::mem::take(first) { + true => Self::None, + false => Self::Same, + } + } +} + +impl FetchedGroup { + /// One past the last Object ID read. + fn end(&self) -> u64 { + self.first + self.frames.len() as u64 + } +} + +/// A FETCH range read out of a track. +struct Walked { + groups: Vec, + /// Why the walk stopped before the end of the range, when it did. + stopped: Option, +} + +/// Read `start..end` out of `track` one group at a time, ascending. +/// +/// Each group is a [`track::Consumer::fetch_group`], so a relay fetches a miss upstream. +/// A group that does not exist below the newest one is a hole and is skipped; one at or +/// past it is the end of the track, so the walk stops there. Any other failure answers +/// the whole FETCH. +async fn walk_fetch(track: &track::Consumer, start: Location, end: Location, priority: u8) -> Result { + let mut groups = Vec::new(); + for sequence in start.group..=end.group { + if sequence == end.group && end.object == 0 { + break; + } + + let skip = match sequence == start.group { + true => start.object, + false => 0, + }; + let until = (sequence == end.group).then_some(end.object); + let fetch = group::Fetch { + priority, + frame_start: skip, + }; + + let mut group = match track.fetch_group(sequence, fetch).await { + Ok(group) => group, + Err(Error::NotFound) if track.latest().is_some_and(|latest| sequence < latest) => continue, + Err(err @ Error::NotFound) => { + return Ok(Walked { + groups, + stopped: Some(err), + }); + } + Err(err) => return Err(err), + }; + + // `fetch_group` positions the consumer at `skip`, or refuses a group that no + // longer holds it. + let first = group.index(); + let mut frames = Vec::new(); + let complete = loop { + if let Some(until) = until + && first + frames.len() as u64 >= until + { + break group.frame_count() as u64 <= until; + } + match group.read_frame().await? { + Some(frame) => frames.push(frame), + None => break true, + } + }; + + groups.push(FetchedGroup { + sequence, + first, + frames, + complete, + }); + } + + Ok(Walked { groups, stopped: None }) +} + /// A broadcast whose route table is watched for changes in what we advertise: the /// namespace becoming (un)advertisable, or its path or cost moving. struct Watched { @@ -863,20 +972,13 @@ where stream, sequence, index, - std::mem::take(&mut first), + FetchPrior::next(&mut first), frame.timestamp, timescale, version, ) .await?; - stream.encode(&(frame.payload.len() as u64)).await?; - if frame.payload.is_empty() && matches!(version, Version::Draft14 | Version::Draft15) { - stream.encode(&0u64).await?; - } - if !frame.payload.is_empty() { - let mut payload = frame.payload; - stream.write_all(&mut payload).await?; - } + Self::write_fetch_payload(stream, frame.payload, version).await?; } index += 1; group.keep_alive(); @@ -908,7 +1010,7 @@ where stream, sequence, index, - std::mem::take(&mut first), + FetchPrior::next(&mut first), frame.timestamp, timescale, version, @@ -954,7 +1056,7 @@ where stream: &mut Writer, sequence: u64, object: u64, - first: bool, + prior: FetchPrior, timestamp: Timestamp, timescale: Option, version: Version, @@ -979,10 +1081,10 @@ where return Ok(()); } - let header = match first { + let header = match prior { // The first object must carry its absolute Group and Object IDs. Include the // priority too: "same as the prior object" has no prior to refer to. - true => ietf::FetchObject::Object { + FetchPrior::None => ietf::FetchObject::Object { subgroup: ietf::FetchSubgroup::Zero, group: Some(sequence), object: Some(object), @@ -990,13 +1092,25 @@ where properties, }, // Same group and priority; the Object ID is the prior one plus one. - false => ietf::FetchObject::Object { + FetchPrior::Same => ietf::FetchObject::Object { subgroup: ietf::FetchSubgroup::Zero, group: None, object: None, priority: None, properties, }, + // A later group, walked in ascending order. Draft-18 turned the field into a + // delta, where zero is the very next group. + FetchPrior::Group(prior) => ietf::FetchObject::Object { + subgroup: ietf::FetchSubgroup::Zero, + group: Some(match version { + Version::Draft15 | Version::Draft16 | Version::Draft17 => sequence, + _ => sequence - prior - 1, + }), + object: Some(object), + priority: None, + properties, + }, }; stream.encode(&header).await?; @@ -1004,6 +1118,24 @@ where Ok(()) } + /// Write a whole fetch object's length and payload, after its header. + async fn write_fetch_payload( + stream: &mut Writer, + payload: bytes::Bytes, + version: Version, + ) -> Result<(), Error> { + stream.encode(&(payload.len() as u64)).await?; + // Draft-14 and 15 still carry a Normal status after an empty object. + if payload.is_empty() && matches!(version, Version::Draft14 | Version::Draft15) { + stream.encode(&0u64).await?; + } + if !payload.is_empty() { + let mut payload = payload; + stream.write_all(&mut payload).await?; + } + Ok(()) + } + /// Register a pending subscription until the returned guard drops. fn register_join(&self, request_id: RequestId) -> Join { self.joins.lock().insert(request_id, None); @@ -1013,166 +1145,168 @@ where } } - /// Serve the current-group prefix for a relative joining FETCH with offset zero. + /// Answer a FETCH: a standalone range of the named track, or a joining FETCH's + /// groups up to its subscription's saved start. + /// + /// Both walk the track one group at a time, ascending, and buffer the answer before + /// replying: FETCH_OK names where the response ends, which a range running past the + /// track only learns by reading it, and a refusal can still replace it until then. async fn run_fetch_stream(mut self, mut stream: Stream, msg: ietf::Fetch<'_>) -> Result<(), Error> { - if Filter::is_draft20(self.version) { - return self - .reject_fetch( - stream, - msg.request_id, - &Error::Unsupported, - "joining FETCH removed in draft-20", - ) - .await; - } + let priority = super::priority::from_wire(msg.subscriber_priority); - let subscribe_id = match msg.fetch_type { - FetchType::Standalone { .. } => { - return self - .reject_fetch(stream, msg.request_id, &Error::Unsupported, "not supported") - .await; - } - FetchType::RelativeJoining { - subscriber_request_id, - group_offset, + let (track, start, end, timescale, joined) = match msg.fetch_type { + FetchType::Standalone { + namespace, + track, + start, + end, } => { - if group_offset != 0 { + // An End Object of 0 asks for the whole End Group. + let end = match end.object { + 0 => end.group.checked_add(1).map(|group| Location { group, object: 0 }), + _ => Some(end), + }; + let Some(end) = end.filter(|end| (start.group, start.object) < (end.group, end.object)) else { return self - .reject_fetch(stream, msg.request_id, &Error::Unsupported, "not supported") + .reject_fetch(stream, msg.request_id, &Error::InvalidRange, "empty range") .await; - } - subscriber_request_id + }; + + // The peer must have seen the announcement to name this namespace, so this + // resolves like a SUBSCRIBE does. + let broadcast = match self.serving_origin().await.request_broadcast(&namespace).await { + Ok(broadcast) => broadcast, + Err(err) => return self.reject_fetch(stream, msg.request_id, &err, &err.to_string()).await, + }; + let track = match broadcast.track(&track) { + Ok(track) => track, + Err(err) => return self.reject_fetch(stream, msg.request_id, &err, &err.to_string()).await, + }; + + // No SUBSCRIBE declared a timescale for this request, so its objects go + // out unstamped. + (track, start, end, None, false) } - FetchType::AbsoluteJoining { .. } => { + // Draft-20 replaced joining FETCH with subscription fills. + FetchType::RelativeJoining { .. } | FetchType::AbsoluteJoining { .. } + if Filter::is_draft20(self.version) => + { return self - .reject_fetch(stream, msg.request_id, &Error::Unsupported, "not supported") + .reject_fetch( + stream, + msg.request_id, + &Error::Unsupported, + "joining FETCH removed in draft-20", + ) .await; } + FetchType::RelativeJoining { + subscriber_request_id, .. + } + | FetchType::AbsoluteJoining { + subscriber_request_id, .. + } => { + let (end, cache, timescale) = match self.joined(&mut stream, subscriber_request_id).await? { + Ok(joined) => joined, + Err((err, reason)) => return self.reject_fetch(stream, msg.request_id, &err, reason).await, + }; + let start = match msg.fetch_type { + FetchType::RelativeJoining { group_offset, .. } => end.group.saturating_sub(group_offset), + FetchType::AbsoluteJoining { group_id, .. } if group_id <= end.group => group_id, + _ => { + return self + .reject_fetch( + stream, + msg.request_id, + &Error::InvalidRange, + "joining group past the subscription", + ) + .await; + } + }; + ( + cache, + Location { + group: start, + object: 0, + }, + end, + timescale, + true, + ) + } }; - // Request streams can arrive out of order. Wait on registration, while bounding - // the lifetime of a request whose subscription never arrives or resolves. - let joined = { - let mut pending = false; - let mut deadline = crate::runtime::Deadline::after(&self.runtime, Duration::from_secs(10)); + // The walk writes each group after the one before it, so the order the peer + // reads object IDs in is fixed. One group has no order to get wrong. + let last = match end.object { + 0 => end.group - 1, + _ => end.group, + }; + if msg.group_order == GroupOrder::Descending && start.group != last { + return self + .reject_fetch( + stream, + msg.request_id, + &Error::Unsupported, + "descending FETCH not supported", + ) + .await; + } + + // The subscriber cancelling is the only other way this ends early, and is owed + // nothing. + let walked = { + let mut walk = std::pin::pin!(walk_fetch(&track, start, end, priority)); kio::wait(|waiter| { let mut cx = std::task::Context::from_waker(waiter.waker()); - // The request reader is what the subscriber FINs or resets. The writer - // on a draft-14-16 virtual stream reports closed immediately, which is - // not a cancellation. if stream.reader.poll_closed(&mut cx).is_ready() { - return Poll::Ready(Err(Error::Cancel)); - } - let joins = self.joins.poll(waiter, |joins| match joins.get(&subscribe_id) { - Some(Some(_)) => Poll::Ready(()), - Some(None) => { - pending = true; - Poll::Pending - } - None if pending => Poll::Ready(()), - None => Poll::Pending, - }); - if let Poll::Ready(joins) = joins { - return Poll::Ready(Ok(joins.get(&subscribe_id).cloned().flatten())); + return Poll::Ready(None); } - if deadline.poll(waiter).is_ready() { - return Poll::Ready(if pending { Err(Error::Timeout) } else { Ok(None) }); - } - Poll::Pending + waiter.poll_future(walk.as_mut()).map(Some) }) .await }; - let joined = match joined { - Err(Error::Timeout) => { - return self - .reject_fetch(stream, msg.request_id, &Error::Timeout, "subscription not ready") - .await; - } - result => result?, + let walked = match walked { + Some(Ok(walked)) => walked, + Some(Err(err)) => return self.reject_fetch(stream, msg.request_id, &err, &err.to_string()).await, + None => return Ok(()), }; - let (end, cache, timescale) = match joined { - None => { - return self - .reject_fetch( - stream, - msg.request_id, - &if matches!( - self.version, - Version::Draft14 - | Version::Draft15 | Version::Draft16 - | Version::Draft17 | Version::Draft18 - | Version::Draft19 - ) { - Error::InvalidJoiningRequestId - } else { - Error::NotFound - }, - "no such subscription", - ) - .await; - } - Some(Joined::Unsupported) => { - if matches!(self.version, Version::Draft14 | Version::Draft15 | Version::Draft16) { - self.session.close( - crate::SessionError::ProtocolViolation.to_code(), - "joining FETCH requires Largest Object filter", - ); - return Err(Error::ProtocolViolation); - } - return self - .reject_fetch( - stream, - msg.request_id, - &Error::Unsupported, - "joining filter not supported", - ) - .await; - } - Some(Joined::Empty) => { - return self - .reject_fetch( - stream, - msg.request_id, - &Error::InvalidRange, - "no objects at subscription start", - ) - .await; + + let (end_location, end_of_track) = if joined { + // The subscription starts at the saved Largest Object, so the prefix of that + // group has to be here in full or the two leave a gap. + let whole = walked + .groups + .last() + .is_some_and(|last| last.sequence == end.group && last.end() == end.object); + if !whole { + let (err, reason) = match walked.stopped { + Some(err) => (err, "joining group unavailable"), + None => (Error::Evicted, "joining prefix unavailable"), + }; + return self.reject_fetch(stream, msg.request_id, &err, reason).await; } - Some(Joined::Group { end, cache, timescale }) => (end, cache, timescale), - }; - let priority = super::priority::from_wire(msg.subscriber_priority); - let mut group = match cache - .fetch_group( - end.group, - group::Fetch { - priority, - ..Default::default() - }, - ) - .await - { - Ok(group) => group, - Err(err) => { - return self - .reject_fetch(stream, msg.request_id, &err, "joining group unavailable") - .await; + (end, false) + } else { + let Some(last) = walked.groups.last() else { + // Nothing in range exists: the refusal the walk ran into is the answer. + let err = walked.stopped.unwrap_or(Error::NotFound); + return self.reject_fetch(stream, msg.request_id, &err, &err.to_string()).await; + }; + let delivered = Location { + group: last.sequence, + object: last.end(), + }; + let end_of_track = last.complete && track.final_sequence() == last.sequence.checked_add(1); + // A response that stopped short ends at its last object; one that covered the + // whole range ends where it was asked to. + match walked.stopped.is_some() || end_of_track { + true => (delivered, end_of_track), + false => (end, false), } }; - // Retain the promised prefix before success: a cached group may already have - // evicted its front, and later cache eviction must not truncate this FETCH. - let mut prefix = Vec::new(); - for _ in 0..end.object { - match group.read_frame().await { - Ok(Some(frame)) => prefix.push(frame), - _ => { - return self - .reject_fetch(stream, msg.request_id, &Error::Evicted, "joining prefix unavailable") - .await; - } - } - } - // FETCH_OK on every draft, never REQUEST_OK: section 5.2 allows exactly one FETCH_OK or // REQUEST_ERROR in answer to a FETCH, and REQUEST_OK's own definition lists the other // requests it answers without ever naming this one. @@ -1185,13 +1319,15 @@ where _ => None, }, // Only draft-14 encodes it, and only as the publisher restating the order. - group_order: msg.group_order.any_to_descending(), - end_of_track: false, - end_location: end, + group_order: match msg.group_order { + GroupOrder::Descending => GroupOrder::Descending, + _ => GroupOrder::Ascending, + }, + end_of_track, + end_location, }) .await?; - // The FETCH owns the retained prefix, capped at the subscription snapshot. let uni = self.session.open_uni().await.map_err(Error::from_transport)?; let mut writer = Writer::new(uni, self.version); writer.set_priority(priority); @@ -1201,24 +1337,26 @@ where request_id: msg.request_id, }) .await?; - for (index, frame) in prefix.into_iter().enumerate() { - Self::write_fetch_object( - &mut writer, - end.group, - index as u64, - index == 0, - frame.timestamp, - timescale, - self.version, - ) - .await?; - writer.encode(&(frame.payload.len() as u64)).await?; - if frame.payload.is_empty() && matches!(self.version, Version::Draft14 | Version::Draft15) { - writer.encode(&0u64).await?; - } - if !frame.payload.is_empty() { - let mut payload = frame.payload; - writer.write_all(&mut payload).await?; + let mut prior: Option = None; + for group in walked.groups { + for (object, frame) in (group.first..).zip(group.frames) { + let at = match prior { + None => FetchPrior::None, + Some(prior) if prior == group.sequence => FetchPrior::Same, + Some(prior) => FetchPrior::Group(prior), + }; + Self::write_fetch_object( + &mut writer, + group.sequence, + object, + at, + frame.timestamp, + timescale, + self.version, + ) + .await?; + Self::write_fetch_payload(&mut writer, frame.payload, self.version).await?; + prior = Some(group.sequence); } } writer.close().await?; @@ -1231,6 +1369,76 @@ where Ok(()) } + /// Resolve the subscription a joining FETCH names to its saved start, or the refusal + /// to answer the FETCH with. + async fn joined( + &mut self, + stream: &mut Stream, + subscribe_id: RequestId, + ) -> Result), (Error, &'static str)>, Error> { + // Request streams can arrive out of order. Wait on registration, while bounding + // the lifetime of a request whose subscription never arrives or resolves. + let joined = { + let mut pending = false; + let mut deadline = crate::runtime::Deadline::after(&self.runtime, Duration::from_secs(10)); + kio::wait(|waiter| { + let mut cx = std::task::Context::from_waker(waiter.waker()); + // The request reader is what the subscriber FINs or resets. The writer + // on a draft-14-16 virtual stream reports closed immediately, which is + // not a cancellation. + if stream.reader.poll_closed(&mut cx).is_ready() { + return Poll::Ready(Err(Error::Cancel)); + } + let joins = self.joins.poll(waiter, |joins| match joins.get(&subscribe_id) { + Some(Some(_)) => Poll::Ready(()), + Some(None) => { + pending = true; + Poll::Pending + } + None if pending => Poll::Ready(()), + None => Poll::Pending, + }); + if let Poll::Ready(joins) = joins { + return Poll::Ready(Ok(joins.get(&subscribe_id).cloned().flatten())); + } + if deadline.poll(waiter).is_ready() { + return Poll::Ready(if pending { Err(Error::Timeout) } else { Ok(None) }); + } + Poll::Pending + }) + .await + }; + let refusal = match joined { + Err(Error::Timeout) => (Error::Timeout, "subscription not ready"), + Err(err) => return Err(err), + Ok(Some(Joined::Group { end, cache, timescale })) => return Ok(Ok((end, cache, timescale))), + Ok(None) => ( + match self.version { + Version::Draft14 + | Version::Draft15 + | Version::Draft16 + | Version::Draft17 + | Version::Draft18 + | Version::Draft19 => Error::InvalidJoiningRequestId, + _ => Error::NotFound, + }, + "no such subscription", + ), + Ok(Some(Joined::Unsupported)) => { + if matches!(self.version, Version::Draft14 | Version::Draft15 | Version::Draft16) { + self.session.close( + crate::SessionError::ProtocolViolation.to_code(), + "joining FETCH requires Largest Object filter", + ); + return Err(Error::ProtocolViolation); + } + (Error::Unsupported, "joining filter not supported") + } + Ok(Some(Joined::Empty)) => (Error::InvalidRange, "no objects at subscription start"), + }; + Ok(Err(refusal)) + } + async fn reject_track_status(&self, mut stream: Stream, request_id: RequestId) -> Result<(), Error> { let error_code = request::to_code(&Error::Unsupported, request::Kind::TrackStatus, self.version); if self.version == Version::Draft14 { @@ -3457,6 +3665,341 @@ mod serve_tests { assert!(buf.is_empty()); } + /// Every draft that carries a standalone FETCH. + const FETCH_DRAFTS: [Version; 9] = [ + Version::Draft14, + Version::Draft15, + Version::Draft16, + Version::Draft17, + Version::Draft18, + Version::Draft19, + Version::Draft20, + Version::Draft21, + Version::Draft22, + ]; + + /// Groups `0..count`, each holding `g-0` and `g-1`, skipping `hole`. + fn publish_pairs(h: &mut Serve, count: u64, hole: Option) { + for sequence in (0..count).filter(|sequence| Some(*sequence) != hole) { + let mut group = h.track.create_group(group::Info { sequence }).unwrap(); + for object in 0..2 { + group + .write_frame(timestamp(), format!("{sequence}-{object}").into_bytes()) + .unwrap(); + } + group.finish().unwrap(); + } + } + + /// Run a standalone FETCH of `room/video`, returning what the peer reads back. + async fn standalone_fetch(h: &Serve, start: Location, end: Location, group_order: GroupOrder) -> bytes::Bytes { + let version = h.publisher.version; + let mark = h.log.writes.lock().unwrap().len(); + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + h.publisher + .clone() + .run_fetch_stream( + stream, + ietf::Fetch { + request_id: FETCH_ID, + subscriber_priority: 128, + group_order, + fetch_type: FetchType::Standalone { + namespace: crate::Path::new("room"), + track: "video".into(), + start, + end, + }, + }, + ) + .await + .unwrap(); + bytes::Bytes::from(h.log.writes.lock().unwrap()[mark..].to_vec()) + } + + /// Decode a FETCH_OK and the fetch stream after it, as `(group, object, payload)`. + fn fetch_answer(mut buf: bytes::Bytes, version: Version) -> (ietf::FetchOk, Vec<(u64, u64, String)>) { + assert_eq!( + u64::decode(&mut buf, version).unwrap(), + ietf::FetchOk::ID, + "{version}: not a FETCH_OK" + ); + let ok = ietf::FetchOk::decode(&mut buf, version).unwrap(); + assert_eq!(u64::decode(&mut buf, version).unwrap(), FetchHeader::TYPE); + assert_eq!(FetchHeader::decode(&mut buf, version).unwrap().request_id, FETCH_ID); + + let mut objects = Vec::new(); + let mut prior: Option<(u64, u64)> = None; + while !buf.is_empty() { + let (group, object) = if version == Version::Draft14 { + let group = u64::decode(&mut buf, version).unwrap(); + assert_eq!(u64::decode(&mut buf, version).unwrap(), 0, "subgroup"); + let object = u64::decode(&mut buf, version).unwrap(); + let _priority = u8::decode(&mut buf, version).unwrap(); + let _properties = Vec::::decode(&mut buf, version).unwrap(); + (group, object) + } else { + let ietf::FetchObject::Object { group, object, .. } = + ietf::FetchObject::decode(&mut buf, version).unwrap() + else { + panic!("{version}: unexpected End of Range"); + }; + match (prior, group) { + (None, group) => (group.unwrap(), object.unwrap()), + // No Group ID: the same group, and an absent Object ID Delta is one. + (Some((group, prior)), None) => (group, prior + object.unwrap_or(1)), + (Some((prior, _)), Some(group)) => { + let group = match version { + Version::Draft15 | Version::Draft16 | Version::Draft17 => group, + _ => prior + group + 1, + }; + (group, object.unwrap()) + } + } + }; + let size = u64::decode(&mut buf, version).unwrap() as usize; + let payload = String::from_utf8(buf.split_to(size).to_vec()).unwrap(); + objects.push((group, object, payload)); + prior = Some((group, object)); + } + (ok, objects) + } + + /// The `(group, object, payload)` a range of two-frame groups holds. + fn pairs(groups: impl IntoIterator) -> Vec<(u64, u64, String)> { + groups + .into_iter() + .flat_map(|group| (0..2).map(move |object| (group, object, format!("{group}-{object}")))) + .collect() + } + + /// A range of whole groups is walked one group at a time and answered on one fetch + /// stream, in each draft's own Group ID encoding. + #[tokio::test] + async fn a_standalone_fetch_walks_whole_groups() { + for version in FETCH_DRAFTS { + let mut h = serve(version); + publish_pairs(&mut h, 5, None); + settle().await; + + // An End Object of 0 asks for the whole End Group. + let buf = standalone_fetch( + &h, + Location { group: 1, object: 0 }, + Location { group: 3, object: 0 }, + GroupOrder::Ascending, + ) + .await; + let (ok, objects) = fetch_answer(buf, version); + assert_eq!(ok.end_location, Location { group: 4, object: 0 }, "{version}"); + assert!(!ok.end_of_track, "{version}"); + assert_eq!(objects, pairs(1..=3), "{version}"); + assert!(h.log.resets().is_empty(), "{version}"); + } + } + + /// A range that runs past the end of a finished track stops there, and FETCH_OK + /// says so: it ends one past the last object, with End of Track set. + #[tokio::test] + async fn a_standalone_fetch_past_the_end_reports_the_end_of_track() { + for version in FETCH_DRAFTS { + let mut h = serve(version); + publish_pairs(&mut h, 5, None); + h.track.finish().unwrap(); + settle().await; + + let buf = standalone_fetch( + &h, + Location { group: 3, object: 1 }, + Location { group: 9, object: 0 }, + GroupOrder::Any, + ) + .await; + let (ok, objects) = fetch_answer(buf, version); + assert_eq!(ok.end_location, Location { group: 4, object: 2 }, "{version}"); + assert!(ok.end_of_track, "{version}"); + assert_eq!(objects, pairs(3..=4)[1..].to_vec(), "{version}"); + } + } + + /// A missing group below the newest one is a hole the walk steps over. + #[tokio::test] + async fn a_standalone_fetch_skips_a_missing_group() { + for version in [Version::Draft16, Version::Draft18] { + let mut h = serve(version); + publish_pairs(&mut h, 4, Some(2)); + settle().await; + + let buf = standalone_fetch( + &h, + Location { group: 1, object: 0 }, + Location { group: 3, object: 1 }, + GroupOrder::Ascending, + ) + .await; + let (ok, objects) = fetch_answer(buf, version); + assert_eq!(ok.end_location, Location { group: 3, object: 1 }, "{version}"); + let mut expected = pairs([1]); + expected.push((3, 0, "3-0".into())); + assert_eq!(objects, expected, "{version}"); + } + } + + /// Read a REQUEST_ERROR (or draft-14 FETCH_ERROR) off a refused FETCH, as its code. + fn fetch_refusal(mut buf: bytes::Bytes, version: Version) -> u64 { + let id = u64::decode(&mut buf, version).unwrap(); + let code = match version { + Version::Draft14 => { + assert_eq!(id, ietf::FetchError::ID); + ietf::FetchError::decode(&mut buf, version).unwrap().error_code + } + _ => { + assert_eq!(id, ietf::RequestError::ID); + ietf::RequestError::decode(&mut buf, version).unwrap().error_code + } + }; + assert!(buf.is_empty(), "{version}: a refusal opens no fetch stream"); + code + } + + /// Nothing in range exists, so the walk's own refusal is the answer. + #[tokio::test] + async fn a_standalone_fetch_past_the_end_is_refused() { + for version in FETCH_DRAFTS { + let mut h = serve(version); + publish_pairs(&mut h, 2, None); + settle().await; + + let buf = standalone_fetch( + &h, + Location { group: 5, object: 0 }, + Location { group: 6, object: 0 }, + GroupOrder::Ascending, + ) + .await; + assert_eq!(fetch_refusal(buf, version), does_not_exist(version), "{version}"); + } + } + + /// The walk only runs forwards, so a descending range of several groups is refused. + #[tokio::test] + async fn a_descending_fetch_of_several_groups_is_refused() { + for version in [Version::Draft14, Version::Draft16, Version::Draft18] { + let mut h = serve(version); + publish_pairs(&mut h, 3, None); + settle().await; + + let buf = standalone_fetch( + &h, + Location { group: 0, object: 0 }, + Location { group: 1, object: 0 }, + GroupOrder::Descending, + ) + .await; + assert_eq!(fetch_refusal(buf, version), 0x3, "{version}"); + } + } + + /// A joining FETCH with an offset walks the whole groups before the subscription's + /// group, then that group's prefix, and ends exactly where the subscription starts. + #[tokio::test] + async fn a_joining_fetch_with_an_offset_walks_back_whole_groups() { + for version in JOINING_DRAFTS { + let mut h = serve(version); + publish_pairs(&mut h, 4, None); + let mut group = h.track.create_group(group::Info { sequence: 4 }).unwrap(); + group.write_frame(timestamp(), b"4-0".as_slice()).unwrap(); + settle().await; + + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + let mut serving = std::pin::pin!( + h.publisher + .clone() + .run_subscribe_stream(stream, subscribe(Filter::NextObject, None)) + ); + registered(&h, serving.as_mut()).await; + + for (fetch_type, first) in [ + ( + FetchType::RelativeJoining { + subscriber_request_id: RequestId(REQUEST_ID), + group_offset: 2, + }, + 2, + ), + ( + FetchType::AbsoluteJoining { + subscriber_request_id: RequestId(REQUEST_ID), + group_id: 1, + }, + 1, + ), + ] { + let mark = h.log.writes.lock().unwrap().len(); + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + h.publisher + .clone() + .run_fetch_stream( + stream, + ietf::Fetch { + request_id: FETCH_ID, + subscriber_priority: 128, + group_order: GroupOrder::Ascending, + fetch_type, + }, + ) + .await + .unwrap(); + let buf = bytes::Bytes::from(h.log.writes.lock().unwrap()[mark..].to_vec()); + + let (ok, objects) = fetch_answer(buf, version); + assert_eq!(ok.end_location, Location { group: 4, object: 1 }, "{version}"); + assert!(!ok.end_of_track, "{version}"); + let mut expected = pairs(first..=3); + expected.push((4, 0, "4-0".into())); + assert_eq!(objects, expected, "{version}"); + } + } + } + + /// An absolute joining FETCH starting past the subscription's group names an empty range. + #[tokio::test] + async fn an_absolute_joining_fetch_past_the_subscription_is_refused() { + let version = Version::Draft16; + let mut h = serve(version); + publish_pairs(&mut h, 3, None); + settle().await; + + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + let mut serving = std::pin::pin!( + h.publisher + .clone() + .run_subscribe_stream(stream, subscribe(Filter::NextObject, None)) + ); + registered(&h, serving.as_mut()).await; + + let mark = h.log.writes.lock().unwrap().len(); + let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); + h.publisher + .clone() + .run_fetch_stream( + stream, + ietf::Fetch { + request_id: FETCH_ID, + subscriber_priority: 128, + group_order: GroupOrder::Ascending, + fetch_type: FetchType::AbsoluteJoining { + subscriber_request_id: RequestId(REQUEST_ID), + group_id: 5, + }, + }, + ) + .await + .unwrap(); + let buf = bytes::Bytes::from(h.log.writes.lock().unwrap()[mark..].to_vec()); + assert_eq!(fetch_refusal(buf, version), invalid_range(version)); + } + /// A fill against an empty track has an empty range: no fetch stream is owed. #[tokio::test] async fn an_empty_track_opens_no_fill_stream() { @@ -4896,8 +5439,8 @@ mod tests { (writes, h.log.resets()) } - /// Send a FETCH we don't implement, returning the same pair. - async fn fetch_unsupported(version: Version, fetch_type: FetchType<'_>) -> (Vec, Vec) { + /// Send a FETCH we refuse, returning the same pair. + async fn fetch_refused(version: Version, fetch_type: FetchType<'_>) -> (Vec, Vec) { let h = harness(version); let stream = Stream::open(&mut h.session.clone(), version).await.unwrap(); @@ -4941,8 +5484,8 @@ mod tests { /// Every FETCH we refuse goes out through its own error encoder, so it needs the same /// finish: a reset there loses the rejection the same way. #[tokio::test] - async fn unsupported_fetch_is_refused_without_resetting_the_stream() { - let unsupported = || { + async fn a_refused_fetch_does_not_reset_the_stream() { + let refused = || { [ ( "standalone", @@ -4954,25 +5497,20 @@ mod tests { }, ), ( - "relative joining with an offset", - FetchType::RelativeJoining { - subscriber_request_id: RequestId(3), - group_offset: 1, - }, - ), - ( - "absolute joining", - FetchType::AbsoluteJoining { - subscriber_request_id: RequestId(3), - group_id: 7, + "empty range", + FetchType::Standalone { + namespace: crate::Path::new("nothing/here"), + track: "video".into(), + start: Location { group: 2, object: 0 }, + end: Location { group: 1, object: 0 }, }, ), ] }; for version in [Version::Draft17, Version::Draft18, Version::Draft19, Version::Draft20] { - for (label, fetch_type) in unsupported() { - let (writes, resets) = fetch_unsupported(version, fetch_type).await; + for (label, fetch_type) in refused() { + let (writes, resets) = fetch_refused(version, fetch_type).await; assert!(!writes.is_empty(), "{version} {label}: nothing was sent"); assert_eq!( diff --git a/rs/moq-net/src/ietf/subscriber.rs b/rs/moq-net/src/ietf/subscriber.rs index b2cea32c26..0f18d9bea2 100644 --- a/rs/moq-net/src/ietf/subscriber.rs +++ b/rs/moq-net/src/ietf/subscriber.rs @@ -145,6 +145,9 @@ struct State { // Joining FETCH request ids, mapped to the SUBSCRIBE they name. fetches: HashMap, + // Group FETCH request ids, filling a cache miss. + group_fetches: HashMap>, + // Track aliases chosen by the remote publisher. aliases: TrackAliases, @@ -285,6 +288,48 @@ impl Fill { } } +/// A standalone FETCH of one whole group, from FETCH_OK to its fetch stream. +enum GroupFetch { + /// Waiting on FETCH_OK, which the fetch stream can overtake. + Pending, + /// Accepted into the track cache, waiting for the fetch stream to write it. + Ready { + producer: group::Producer, + timescale: Option, + }, + /// A fetch stream is writing the group. + Receiving, + /// The group is written or failed. + Done, +} + +/// A group FETCH's entry in [`State::group_fetches`], removed when its request ends. +struct GroupFetchEntry { + state: Lock, + fetch_id: RequestId, + slot: kio::Producer, +} + +impl GroupFetchEntry { + fn new(state: &Lock, fetch_id: RequestId, slot: kio::Producer) -> Self { + state.lock().group_fetches.insert(fetch_id, slot.clone()); + Self { + state: state.clone(), + fetch_id, + slot, + } + } +} + +impl Drop for GroupFetchEntry { + fn drop(&mut self) { + self.state.lock().group_fetches.remove(&self.fetch_id); + // A fetch stream that overtook a refused or unserved FETCH_OK holds its own clone + // of the slot, so only closing it wakes that stream. + let _ = self.slot.close(); + } +} + /// A pre-draft-20 joining FETCH, sent as its own request after SUBSCRIBE. #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum JoiningFetch { @@ -1760,6 +1805,9 @@ where true => request.resolving_start(), false => request, }; + // Serves cache misses with a group FETCH. Registered before accepting, so a miss + // queued meanwhile waits for it rather than failing for want of a handler. + let dynamic = request.dynamic(); let mut track = request.accept(info); if live { let _ = track.start_at(largest.map(|largest| largest.group)); @@ -1806,9 +1854,12 @@ where enum End { Unused, Done(Result), + Fetch(group::Request), } let mut fetch_done = fetching.is_none(); + // Group FETCHes for cache misses, cancelled with the subscription. + let mut group_fetches = TaskSet::owned(); let cancelled = { let mut done = std::pin::pin!(Self::read_publish_done(&mut stream.reader, self.version)); loop { @@ -1819,6 +1870,11 @@ where { fetch_done = true; } + // An error is the track closing, which the arms below report. + if let Poll::Ready(Ok(request)) = dynamic.poll_requested_group(waiter) { + return Poll::Ready(End::Fetch(request)); + } + let _ = group_fetches.poll(waiter); if track.poll_unused(waiter).is_ready() { return Poll::Ready(End::Unused); } @@ -1827,6 +1883,15 @@ where .await; match end { + End::Fetch(request) => { + let fetch = self.clone().run_group_fetch( + broadcast_path.to_owned(), + track_name.clone(), + request, + timescale, + ); + group_fetches.push(fetch); + } End::Unused => match track.abort_unused(Error::Cancel) { Ok(()) => { tracing::info!(broadcast = %self.origin.absolute(&broadcast_path), track = %track_name, "subscribe cancelled"); @@ -2317,6 +2382,9 @@ enum Ended { Track, } +/// Object status: no object at or past this location in the group exists (draft-14, 15). +const END_OF_GROUP: u64 = 0x3; + /// Object status: no object at or past this location exists (every implemented draft). const END_OF_TRACK: u64 = 0x4; @@ -2528,6 +2596,11 @@ where let _: u64 = stream.decode().await?; let header: ietf::FetchHeader = stream.decode().await?; + let group_fetch = self.state.lock().group_fetches.get(&header.request_id).cloned(); + if let Some(slot) = group_fetch { + return self.recv_group_fetch(stream, slot).await; + } + let (subscribe_id, fill, joining, largest, _counted) = { let state = self.state.lock(); // A draft-20 fill is named by the SUBSCRIBE's request id. A pre-draft-20 joining @@ -2743,40 +2816,286 @@ where } } - // The properties carry the frame's presentation timestamp (the Timestamp Object - // Property) in the units the track declared. A track that declared none opted - // out, so its frames are stamped on arrival instead. - let timestamp = match (object.properties, timescale) { - (Some(properties), Some(timescale)) => { - let mut properties = bytes::Bytes::from(properties); - ietf::decode_object_time(&mut properties, timescale, self.version)? - } - _ => None, - }; - let timestamp = timestamp.unwrap_or_else(|| crate::Timestamp::from(self.runtime.now())); - - // A fetch object has no status field from draft-16 on; a zero length is simply - // an empty object. Draft-14 and 15 still encode Normal (0) after a zero length. - let size: u64 = stream.decode().await?; - if size == 0 && matches!(self.version, Version::Draft14 | Version::Draft15) { - let status: u64 = stream.decode().await?; - if status != 0 { - return Err(Error::Unsupported); + let (_, next, producer) = head.as_mut().expect("the head was created above"); + if !self + .recv_fetch_payload(stream, producer, object.properties, timescale) + .await? + { + return Err(Error::Unsupported); + } + *next += 1; + } + + Ok(()) + } + + /// Read one fetch object's length and payload into `producer`, after its header. + /// + /// Returns `false` for a draft-14 or 15 end-of-group or end-of-track marker, which + /// is a status rather than a frame. + async fn recv_fetch_payload( + &self, + stream: &mut Reader, + producer: &mut group::Producer, + properties: Option>, + timescale: Option, + ) -> Result { + // The properties carry the frame's presentation timestamp (the Timestamp Object + // Property) in the units the track declared. A track that declared none opted + // out, so its frames are stamped on arrival instead. + let timestamp = match (properties, timescale) { + (Some(properties), Some(timescale)) => { + let mut properties = bytes::Bytes::from(properties); + ietf::decode_object_time(&mut properties, timescale, self.version)? + } + _ => None, + }; + let timestamp = timestamp.unwrap_or_else(|| crate::Timestamp::from(self.runtime.now())); + + // A fetch object has no status field from draft-16 on; a zero length is simply + // an empty object. Draft-14 and 15 still encode Normal (0) after a zero length. + let size: u64 = stream.decode().await?; + if size == 0 && matches!(self.version, Version::Draft14 | Version::Draft15) { + match stream.decode::().await? { + 0 => {} + END_OF_GROUP | END_OF_TRACK => return Ok(false), + _ => return Err(Error::Unsupported), + } + } + + // `create_frame_owned` is the allocation chokepoint and rejects an oversized `size` + // before allocating, so no pre-check is needed. + let mut frame = producer.create_frame_owned(frame::Info { size, timestamp })?; + if let Err(err) = std::future::poll_fn(|cx| stream.poll_read_frame(cx, &mut frame)).await { + let _ = frame.abort(err.clone()); + return Err(err); + } + frame.finish()?; + Ok(true) + } + + /// Fetch one whole group from the publisher to fill a cache miss, with a standalone + /// FETCH for this subscription's track. + /// + /// The group is accepted only once FETCH_OK arrives, so a refusal reaches every + /// waiting [`track::Consumer::fetch_group`] as the publisher's own error. The objects + /// arrive on a fetch stream, which [`Self::recv_fill`] routes into the group. + async fn run_group_fetch( + self, + broadcast: PathOwned, + name: String, + request: group::Request, + timescale: Option, + ) { + let sequence = request.sequence(); + if self.going_away.is_set() { + request.reject(Error::GoingAway); + return; + } + + let fetch_id = match self.control.next_request_id(&self.runtime).await { + Ok(id) => id, + Err(err) => return request.reject(err), + }; + + // Registered before the FETCH goes out, since its fetch stream can overtake FETCH_OK. + let slot = kio::Producer::new(GroupFetch::Pending); + let _registered = GroupFetchEntry::new(&self.state, fetch_id, slot.clone()); + + let mut stream = match Stream::open(&mut self.session.clone(), self.version).await { + Ok(stream) => stream, + Err(err) => return request.reject(err), + }; + + let res = async { + stream.writer.encode(&ietf::Fetch::ID).await?; + stream + .writer + .encode(&ietf::Fetch { + request_id: fetch_id, + subscriber_priority: super::priority::to_wire(request.priority()), + group_order: GroupOrder::Ascending, + // An End Object of 0 is the whole End Group. + fetch_type: FetchType::Standalone { + namespace: broadcast.clone(), + track: name.as_str().into(), + start: ietf::Location { + group: sequence, + object: 0, + }, + end: ietf::Location { + group: sequence, + object: 0, + }, + }, + }) + .await?; + self.read_group_fetch_response(&mut stream).await + } + .await; + + let ok = match res { + Ok(ok) => ok, + Err(err) => { + tracing::debug!(%err, group = sequence, "group fetch refused"); + request.reject(err); + let _ = stream.writer.close().await; + return; + } + }; + + // The publisher knows where the track ends, which a range FETCH downstream needs. + if ok.end_of_track { + let end = ok.end_location; + request.finish_track_at(end.group + u64::from(end.object > 0)); + } + + // An empty answer opens no fetch stream at all. + if (ok.end_location.group, ok.end_location.object) <= (sequence, 0) { + request.reject(Error::NotFound); + let _ = stream.writer.close().await; + return; + } + + // Only a track nothing has subscribed to yet takes this info, as SUBSCRIBE_OK would + // have set it. + let info = track::Info::default() + .with_timescale(Timescale::MICRO) + .with_max_age(self.origin.default_max_age()); + let producer = match request.accept(info) { + Ok(producer) => producer, + // Already served by a concurrent fetch, or the track closed. + Err(err) => { + tracing::debug!(%err, group = sequence, "group fetch not served"); + let _ = stream.writer.close().await; + return; + } + }; + if let Ok(mut state) = slot.write() { + *state = GroupFetch::Ready { producer, timescale }; + } + + // Hold the request open until its fetch stream is done: closing our side first is + // what a draft-14-16 adapter reads as cancelling the FETCH. + let _ = slot + .wait(|state| match &**state { + GroupFetch::Done => Poll::Ready(()), + _ => Poll::Pending, + }) + .await; + let _ = stream.writer.close().await; + } + + /// Read the answer to a group FETCH: FETCH_OK, or the publisher's refusal as an error. + async fn read_group_fetch_response(&self, stream: &mut Stream) -> Result { + let type_id: u64 = stream.reader.decode().await?; + let size: u16 = stream.reader.decode().await?; + let mut data = stream.reader.read_exact(size as usize).await?; + + match type_id { + ietf::FetchOk::ID => Ok(ietf::FetchOk::decode_msg(&mut data, self.version)?), + ietf::FetchError::ID if self.version == Version::Draft14 => { + let msg = ietf::FetchError::decode_msg(&mut data, self.version)?; + Err(request::from_code(msg.error_code, request::Kind::Fetch, self.version)) + } + ietf::RequestError::ID => { + let msg = ietf::RequestError::decode_msg(&mut data, self.version)?; + Err(request::from_code(msg.error_code, request::Kind::Fetch, self.version)) + } + _ => Err(Error::UnexpectedMessage), + } + } + + /// Write a group FETCH's objects into the group it accepted. + async fn recv_group_fetch( + &mut self, + stream: &mut Reader, + slot: kio::Producer, + ) -> Result<(), Error> { + // FETCH_OK can trail its own fetch stream. Taking the group in the same step is + // what refuses a second stream for one request. + let taken = kio::wait(|waiter| { + match slot.poll(waiter, |state| match &**state { + GroupFetch::Pending => Poll::Pending, + _ => Poll::Ready(()), + }) { + Poll::Ready(Ok(mut state)) => { + Poll::Ready(match std::mem::replace(&mut *state, GroupFetch::Receiving) { + GroupFetch::Ready { producer, timescale } => Ok((producer, timescale)), + other => { + *state = other; + Err(Error::Unsupported) + } + }) } + Poll::Ready(Err(_)) => Poll::Ready(Err(Error::Dropped)), + Poll::Pending => Poll::Pending, } + }) + .await; + let (producer, timescale) = taken?; - let (_, next, producer) = head.as_mut().expect("the head was created above"); + let sequence = producer.info().sequence; + let mut producer = crate::recv::Group::new(producer); + let res = match self + .recv_group_fetch_objects(stream, &mut producer, sequence, timescale) + .await + { + Ok(()) => producer.finish(), + Err(err) => { + let _ = producer.abort(err.clone()); + Err(err) + } + }; - // `create_frame_owned` is the allocation chokepoint and rejects an oversized `size` - // before allocating, so no pre-check is needed. - let mut frame = producer.create_frame_owned(frame::Info { size, timestamp })?; - if let Err(err) = std::future::poll_fn(|cx| stream.poll_read_frame(cx, &mut frame)).await { - let _ = frame.abort(err.clone()); - return Err(err); + if let Ok(mut state) = slot.write() { + *state = GroupFetch::Done; + } + res + } + + /// Decode one whole group's objects: all in `sequence`, numbered from 0 with no gaps. + async fn recv_group_fetch_objects( + &self, + stream: &mut Reader, + producer: &mut group::Producer, + sequence: u64, + timescale: Option, + ) -> Result<(), Error> { + let mut next = 0u64; + let mut prior_group = None; + let mut ended = false; + while let Some(object) = decode_fetch_object(stream, self.version).await? { + if ended { + tracing::warn!(sequence, "a group fetch continued past its end marker"); + return Err(Error::ProtocolViolation); + } + if !object.subgroup_ok { + tracing::warn!("subgroup ID is not supported, dropping group fetch"); + return Err(Error::Unsupported); } - frame.finish()?; - *next += 1; + let group = resolve_fetch_group(self.version, prior_group, object.group)?; + if let Some(group) = group { + prior_group = Some(group); + } + let id = match (object.group.is_some(), object.object) { + (true, Some(id)) => Some(id), + (false, None | Some(1)) => Some(next), + _ => None, + }; + if group.is_some_and(|group| group != sequence) || id != Some(next) { + tracing::warn!(sequence, next, group = ?group, object = ?id, "a group fetch must be one whole group"); + return Err(Error::Unsupported); + } + + match self + .recv_fetch_payload(stream, producer, object.properties, timescale) + .await? + { + true => next += 1, + false => ended = true, + } } Ok(()) @@ -2965,7 +3284,7 @@ impl GroupIngest { let frame = group.create_frame_owned(frame::Info { size: 0, timestamp })?; frame.finish()?; self.phase = IngestPhase::Delta; - } else if status == 3 && !self.has_end { + } else if status == END_OF_GROUP && !self.has_end { self.phase = IngestPhase::Finished(Ended::Group); } else if status == END_OF_TRACK { // Defined on every implemented draft, whether or not the header marks @@ -6437,6 +6756,59 @@ mod stitch_tests { assert_eq!(frames[0].1, b"g7-0"); assert_eq!(frames[1].1, b"g7-1"); } + + /// A draft-14/15 end-of-group marker ends a group fetch, so an object after it is a + /// violation rather than another frame. + #[tokio::test] + async fn a_group_fetch_refuses_an_object_past_its_end_marker() { + const DRAFT: Version = Version::Draft15; + let header = |object: u64| ietf::FetchObject::Object { + subgroup: ietf::FetchSubgroup::Zero, + group: Some(SEQUENCE), + object: Some(object), + priority: Some(0), + properties: None, + }; + let mut buf = bytes::BytesMut::new(); + header(0).encode(&mut buf, DRAFT).unwrap(); + 1u64.encode(&mut buf, DRAFT).unwrap(); + buf.put_slice(b"a"); + header(1).encode(&mut buf, DRAFT).unwrap(); + 0u64.encode(&mut buf, DRAFT).unwrap(); + END_OF_GROUP.encode(&mut buf, DRAFT).unwrap(); + header(1).encode(&mut buf, DRAFT).unwrap(); + 1u64.encode(&mut buf, DRAFT).unwrap(); + buf.put_slice(b"b"); + + let mut session = ScriptedSession::per_stream_eof(vec![buf.to_vec()]); + let tasks = TaskSet::new(); + let subscriber = Subscriber::new( + crate::time::Clock::tokio(), + session.clone(), + crate::origin::Config::new(crate::Hop::new(1).unwrap()).produce(), + Control::new(None, false), + None, + peer::PeerSetup::default(), + crate::Hop::new(1).unwrap(), + None, + DRAFT, + tasks.0.clone(), + Default::default(), + ); + let (_, recv) = session.open_bi().await.unwrap(); + let mut stream = Reader::new(recv, DRAFT); + + let track = track::Producer::new( + std::sync::Arc::new(crate::broadcast::Info::default()), + "video", + track::Info::default(), + ); + let mut group = track.create_group(group::Info { sequence: SEQUENCE }).unwrap(); + let res = subscriber + .recv_group_fetch_objects(&mut stream, &mut group, SEQUENCE, None) + .await; + assert!(matches!(res, Err(Error::ProtocolViolation)), "{res:?}"); + } } /// A scripted peer answering SUBSCRIBE then FETCH, so the join is spelled on the wire. diff --git a/rs/moq-net/src/model/resume.rs b/rs/moq-net/src/model/resume.rs index 0d9b8ba038..0bf34326c7 100644 --- a/rs/moq-net/src/model/resume.rs +++ b/rs/moq-net/src/model/resume.rs @@ -690,6 +690,13 @@ impl Consumer { self.state.read().latest() } + /// The newest segment's declared exclusive end, where fetches are routed. + pub(crate) fn final_sequence(&self) -> Option { + // Copied out: the segment's track takes its own lock. + let track = self.state.read().segments.last().map(|segment| segment.track.clone())?; + track.final_sequence() + } + /// One past the newest position across the segments: where a route taking this /// logical track over would resume. pub(crate) fn resume_position(&self) -> Option { @@ -796,11 +803,13 @@ impl kio::Pollable for Fetching { type Output = Result; fn poll(&self, waiter: &kio::Waiter) -> Poll { - if let Some(group) = (Consumer { + if let Some(mut group) = (Consumer { state: self.state.clone(), }) .cached_group(self.sequence, self.options.frame_start) { + // Sitting where the caller asked, as a fetch from the segment's own track would. + group.start_at(self.options.frame_start); return Poll::Ready(Ok(group)); } @@ -2803,6 +2812,33 @@ mod test { assert_eq!(read(&mut group), b"b4"); } + /// A cached copy is handed back sitting at the requested frame, the same as a fetch + /// the segment's own track answers. + #[tokio::test] + async fn a_cached_fetch_starts_at_the_requested_frame() { + let (track_a, consumer_a) = track_pair("a"); + let mut producer = Producer::new(); + producer.switch(&consumer_a, None).unwrap(); + + let mut group = track_a.create_group(group::Info { sequence: 0 }).unwrap(); + for payload in ["f0", "f1", "f2"] { + group + .write_frame(crate::Timestamp::ZERO, payload.as_bytes().to_vec()) + .unwrap(); + } + group.finish().unwrap(); + + let mut group = producer + .consume() + .fetch_group(0, group::Fetch::default().with_frame_start(1)) + .now_or_never() + .expect("cached fetch should resolve") + .unwrap(); + assert_eq!(group.index(), 1); + let frame = group.read_frame().now_or_never().unwrap().unwrap().unwrap(); + assert_eq!(frame.payload.as_ref(), b"f1"); + } + #[tokio::test] async fn fetch_waits_for_first_segment() { let (mut track_a, consumer_a) = track_pair("a"); diff --git a/rs/moq-net/src/model/track.rs b/rs/moq-net/src/model/track.rs index 105ce7c636..75a7bb0de3 100644 --- a/rs/moq-net/src/model/track.rs +++ b/rs/moq-net/src/model/track.rs @@ -2744,6 +2744,16 @@ impl Consumer { } } + /// The declared exclusive final sequence, or `None` while the track is open ended. + /// + /// A spliced track answers for its newest segment, which is where fetches go. + pub(crate) fn final_sequence(&self) -> Option { + match &self.inner { + ConsumerKind::Plain(state) => state.read().final_sequence, + ConsumerKind::Spliced(resume) => resume.final_sequence(), + } + } + /// The frame-precise point a replacement route should resume from: one past the /// last frame this copy produced. `None` if it produced nothing. /// @@ -2938,6 +2948,16 @@ impl group::Request { res } + /// Declare the track's exclusive final sequence, as the publisher answering this + /// fetch reported it. A no-op once the track declared one, or holds a later group. + pub(crate) fn finish_track_at(&self, final_sequence: u64) { + if let Ok(mut state) = TrackState::modify(&self.state) + && state.final_sequence.is_none() + { + let _ = state.set_final(final_sequence); + } + } + /// Reject the fetch, resolving every joined [`Consumer::fetch_group`] with `err`. pub fn reject(mut self, err: Error) { self.done = true; diff --git a/rs/moq-tokio/tests/broadcast.rs b/rs/moq-tokio/tests/broadcast.rs index 23a16c1fde..3e1cd2a472 100644 --- a/rs/moq-tokio/tests/broadcast.rs +++ b/rs/moq-tokio/tests/broadcast.rs @@ -376,6 +376,130 @@ async fn broadcast_moq_lite_05_fetch_webtransport() { lite05_fetch_roundtrip("https").await; } +/// A cache miss over moq-transport is a standalone FETCH of the one group, served from +/// the publisher's cache, and a group the publisher lacks comes back as the publisher's +/// own refusal. +async fn transport_fetch_roundtrip(version: &str) { + let pub_origin = moq_tokio::origin::spawn(); + let broadcast = pub_origin.create_broadcast("test").expect("failed to create broadcast"); + broadcast + .announce(Default::default()) + .expect("failed to announce broadcast"); + let track = broadcast.create_track("video", None).expect("failed to create track"); + for sequence in 0..3u64 { + let mut group = track.append_group().expect("failed to append group"); + for frame in 0..2 { + group + .write_frame( + moq_net::Timestamp::ZERO, + bytes::Bytes::from(format!("{sequence}-{frame}")), + ) + .expect("failed to write frame"); + } + group.finish().expect("failed to finish group"); + } + + let mut server_config = moq_tokio::listen::Config::default(); + server_config.bind = Some("[::]:0".parse().unwrap()); + server_config.tls.generate = vec!["localhost".into()]; + server_config.version = vec![version.parse().unwrap()]; + let server = server_config.init(Default::default()).expect("failed to init server"); + let mut server = server.listen().await.expect("failed to listen"); + let addr = server.local_addr().expect("failed to get local addr"); + + let sub_origin = moq_tokio::origin::spawn(); + let sub_consumer = sub_origin.consume(); + let mut announcements = sub_consumer.announced(); + + let mut client_config = moq_tokio::connect::Config::default(); + client_config.tls.insecure = Some(true); + client_config.version = vec![version.parse().unwrap()]; + let client = client_config.init(Default::default()).expect("failed to init client"); + let url: url::Url = format!("moqt://localhost:{}", addr.port()).parse().unwrap(); + + let server_handle = tokio::spawn(async move { + let request = server.accept().await.expect("no incoming connection"); + let session = request.with_publisher(&pub_origin).ok().await?; + let _broadcast = broadcast; + let _track = track; + let _ = session.closed().await; + Ok::<_, anyhow::Error>(()) + }); + + let client = client.with_subscriber(sub_origin); + let (_client, connection) = tokio::time::timeout(TIMEOUT, connect_once(client, url)) + .await + .expect("client connect timed out") + .expect("client connect failed"); + + tokio::time::timeout(TIMEOUT, announcements.next()) + .await + .expect("announce timed out") + .expect("origin closed"); + let bc = tokio::time::timeout(TIMEOUT, sub_consumer.request_broadcast("test")) + .await + .expect("request timed out") + .expect("announced broadcast resolves"); + + let mut group = tokio::time::timeout(TIMEOUT, async { bc.track("video").unwrap().fetch_group(1, None).await }) + .await + .expect("fetch timed out") + .expect("fetch failed"); + assert_eq!(group.sequence, 1); + for frame in 0..2 { + let payload = tokio::time::timeout(TIMEOUT, group.read_frame()) + .await + .expect("read timed out") + .expect("read failed") + .expect("group ended early"); + assert_eq!(payload.payload, bytes::Bytes::from(format!("1-{frame}"))); + } + let end = tokio::time::timeout(TIMEOUT, group.read_frame()) + .await + .expect("read timed out") + .expect("read failed"); + assert!(end.is_none(), "the group ends after its frames"); + + let refused = tokio::time::timeout(TIMEOUT, async { bc.track("video").unwrap().fetch_group(7, None).await }) + .await + .expect("fetch timed out"); + assert!( + matches!(refused, Err(moq_net::Error::NotFound)), + "expected the publisher's refusal, got {:?}", + refused.err() + ); + + drop(connection); + server_handle + .await + .expect("server task panicked") + .expect("server task failed"); +} + +#[tracing_test::traced_test] +#[tokio::test] +async fn fetch_moq_transport_14() { + transport_fetch_roundtrip("moq-transport-14").await; +} + +#[tracing_test::traced_test] +#[tokio::test] +async fn fetch_moq_transport_16() { + transport_fetch_roundtrip("moq-transport-16").await; +} + +#[tracing_test::traced_test] +#[tokio::test] +async fn fetch_moq_transport_18() { + transport_fetch_roundtrip("moq-transport-18").await; +} + +#[tracing_test::traced_test] +#[tokio::test] +async fn fetch_moq_transport_20() { + transport_fetch_roundtrip("moq-transport-20").await; +} + /// A fetch must be served while a live subscription is active on the same track. /// The relay subscribes starting at the latest group, so an older group isn't /// cached and the fetch has to issue a wire FETCH concurrently with the