diff --git a/quest/m1/moxygen/README.md b/quest/m1/moxygen/README.md index bd040aad9a..8fae11bd13 100644 --- a/quest/m1/moxygen/README.md +++ b/quest/m1/moxygen/README.md @@ -45,7 +45,6 @@ Docs stay inline in the change that makes them stale. No new guide. ## Required - [Sparse FETCH ranges](/quest/m1/moxygen/fetch-span.md) - a FETCH costs the groups it returns, not the span of its range -- [Group fetch fill](/quest/m1/moxygen/fetch-fill.md) - a cache fill from an IETF upstream is complete or refused, validated against a concrete end signal, and asks from the frame the reader wants ## Related diff --git a/quest/m1/moxygen/fetch-fill.md b/quest/m1/moxygen/fetch-fill.md deleted file mode 100644 index 95b73655cf..0000000000 --- a/quest/m1/moxygen/fetch-fill.md +++ /dev/null @@ -1,58 +0,0 @@ -# [M] Group fetch fill - -## Goal - -A relay filling a cache miss from an IETF upstream caches exactly what the -upstream promised, or nothing. A short or malformed fetch stream fails the -group loudly instead of landing in the cache as a complete one, and the -upstream FETCH asks for the frames the downstream reader actually wants. - -## Plan - -[#4276](https://github.com/moq-dev/moq/pull/4276) added the fill: each cache -miss becomes a standalone FETCH of one group (`run_group_fetch` in -`rs/moq-net/src/ietf/subscriber.rs`), and `recv_group_fetch_objects` writes -the fetch stream into the accepted group. Three of Codex's review findings -merged unanswered, and the maintainer ruled that each blocks the line: - -- **Truncation is cached as complete** - ([r4113942736](https://github.com/moq-dev/moq/pull/4276#discussion_r4113942736)). - FETCH_OK names an `end_location`, but it never reaches the decoder. If the - peer cleanly ends the stream early, the group is finished and cached short, - and later readers see a normal end. Check the exclusive received boundary - (last object ID + 1) against FETCH_OK's `end_location`, which is exclusive, - before finishing the producer, and abort the group otherwise. That only - works when `end_location` names a concrete object: for a whole-group request our own `run_fetch_stream` answers with - the requested boundary (`(group + 1, 0)`), so a stream holding only object - 0 looks like a valid one-object group. The line's accepted declines keep - that requested end, so for this case find a wire signal that marks a - complete group (such as an End of Group status object) and require it; if - the drafts we speak offer none, ask the maintainer rather than guess. -- **A first object with no IDs is accepted** - ([r4113942737](https://github.com/moq-dev/moq/pull/4276#discussion_r4113942737)). - The `(false, None | Some(1))` arm treats omitted Group and Object IDs as - "same group, next object", which only means something when a prior object - exists. On the first object `prior_group` is `None` and `next` is 0, so an - anonymous object is cached under the requested group. The first object must - carry explicit IDs resolving to the requested group and start object; refuse - it as a protocol violation otherwise. The same holds for every field that - inherits from a predecessor (draft-15+ `FetchSubgroup::Prior`, an omitted - publisher priority): on the first object each must be explicit. -- **The upstream FETCH ignores the requested frame offset** - ([r4113942733](https://github.com/moq-dev/moq/pull/4276#discussion_r4113942733)). - `group::Request::frame_start()` carries the downstream reader's start, but - the upstream FETCH always asks from object 0, the decoder numbers from 0, - and the handler never calls `start_at` on the producer `accept` returns, - though `frame_start`'s docs require it. An upstream that evicted the - prefix but holds the suffix refuses a request it could have answered. Ask - from `frame_start`, `start_at` it, and number the fill from it, so the - request, the decoder, and the producer agree on where the group starts. - -Each fix gets a regression test in the existing subscriber test module that -fails without it. Keep the whole-group-only scope of the line: this is about -filling faithfully, not about new FETCH shapes. - -## Related - -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line this blocks -- [Fetch without SUBSCRIBE](/quest/m1/ietf-fetch-only.md) - also changes how a relay reaches upstream for fetches diff --git a/rs/moq-net/src/ietf/subscriber.rs b/rs/moq-net/src/ietf/subscriber.rs index f7d852e90a..cc97454318 100644 --- a/rs/moq-net/src/ietf/subscriber.rs +++ b/rs/moq-net/src/ietf/subscriber.rs @@ -288,7 +288,7 @@ impl Fill { } } -/// A standalone FETCH of one whole group, from FETCH_OK to its fetch stream. +/// A standalone FETCH of one group through its end, from FETCH_OK to its fetch stream. enum GroupFetch { /// Waiting on FETCH_OK, which the fetch stream can overtake. Pending, @@ -296,6 +296,10 @@ enum GroupFetch { Ready { producer: group::Producer, timescale: Option, + /// The first Object ID requested, which the producer starts at. + start: u64, + /// The exclusive last Object ID, when FETCH_OK's End Location falls inside the group. + end: Option, }, /// A fetch stream is writing the group. Receiving, @@ -2812,7 +2816,8 @@ where head: &mut Option<(u64, u64, crate::recv::Group)>, ) -> Result<(), Error> { let mut prior_group = None; - while let Some(object) = decode_fetch_object(stream, self.version).await? { + let mut first = true; + while let Some(object) = decode_fetch_object(stream, self.version, std::mem::take(&mut first)).await? { if !object.subgroup_ok { tracing::warn!("subgroup ID is not supported, dropping fill"); return Err(Error::Unsupported); @@ -2929,8 +2934,9 @@ where Ok(true) } - /// Fetch one whole group from the publisher to fill a cache miss, with a standalone - /// FETCH for this subscription's track. + /// Fetch a group from the publisher to fill a cache miss, with a standalone FETCH for + /// this subscription's track. It asks from the frame the reader wants through the end + /// of the group, so a publisher that evicted the prefix can still answer. /// /// 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 @@ -2943,6 +2949,7 @@ where timescale: Option, ) { let sequence = request.sequence(); + let start = request.frame_start(); if self.going_away.is_set() { request.reject(Error::GoingAway); return; @@ -2976,7 +2983,7 @@ where track: name.as_str().into(), start: ietf::Location { group: sequence, - object: 0, + object: start, }, end: ietf::Location { group: sequence, @@ -3006,18 +3013,22 @@ where } // An empty answer opens no fetch stream at all. - if (ok.end_location.group, ok.end_location.object) <= (sequence, 0) { + let end = ok.end_location; + if (end.group, end.object) <= (sequence, start) { request.reject(Error::NotFound); let _ = stream.writer.close().await; return; } + // An exclusive End Location inside the group is its Largest Object plus one, so + // the stream owes every object up to it. One past the group covers it whole. + let end = (end.group == sequence).then_some(end.object); // 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) { + let mut producer = match request.accept(info) { Ok(producer) => producer, // Already served by a concurrent fetch, or the track closed. Err(err) => { @@ -3026,8 +3037,19 @@ where return; } }; + // The objects keep the IDs they have in the group rather than restarting at 0. + if let Err(err) = producer.start_at(start) { + let _ = producer.abort(err); + let _ = stream.writer.close().await; + return; + } if let Ok(mut state) = slot.write() { - *state = GroupFetch::Ready { producer, timescale }; + *state = GroupFetch::Ready { + producer, + timescale, + start, + end, + }; } // Hold the request open until its fetch stream is done: closing our side first is @@ -3076,7 +3098,12 @@ where }) { Poll::Ready(Ok(mut state)) => { Poll::Ready(match std::mem::replace(&mut *state, GroupFetch::Receiving) { - GroupFetch::Ready { producer, timescale } => Ok((producer, timescale)), + GroupFetch::Ready { + producer, + timescale, + start, + end, + } => Ok((producer, timescale, start, end)), other => { *state = other; Err(Error::Unsupported) @@ -3088,12 +3115,11 @@ where } }) .await; - let (producer, timescale) = taken?; + let (producer, timescale, start, end) = taken?; - 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) + .recv_group_fetch_objects(stream, &mut producer, start, end, timescale) .await { Ok(()) => producer.finish(), @@ -3109,18 +3135,22 @@ where res } - /// Decode one whole group's objects: all in `sequence`, numbered from 0 with no gaps. + /// Decode one group's objects: all in the producer's group, numbered from `start` with + /// no gaps, and through `end` (exclusive) when FETCH_OK named one inside the group. async fn recv_group_fetch_objects( &self, stream: &mut Reader, producer: &mut group::Producer, - sequence: u64, + start: u64, + end: Option, timescale: Option, ) -> Result<(), Error> { - let mut next = 0u64; + let sequence = producer.info().sequence; + let mut next = start; let mut prior_group = None; let mut ended = false; - while let Some(object) = decode_fetch_object(stream, self.version).await? { + let mut first = true; + while let Some(object) = decode_fetch_object(stream, self.version, std::mem::take(&mut first)).await? { if ended { tracing::warn!(sequence, "a group fetch continued past its end marker"); return Err(Error::ProtocolViolation); @@ -3139,8 +3169,14 @@ where (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"); + // Another group, or an ID before the next one, is outside what was asked for. + if group.is_some_and(|group| group != sequence) || id.is_some_and(|id| id < next) { + tracing::warn!(sequence, next, group = ?group, object = ?id, "a group fetch answered outside its range"); + return Err(Error::ProtocolViolation); + } + // A skipped ID is an object that does not exist, a hole the model cannot hold. + if id != Some(next) { + tracing::warn!(sequence, next, object = ?id, "a group fetch skipped an object"); return Err(Error::Unsupported); } @@ -3151,6 +3187,22 @@ where true => next += 1, false => ended = true, } + if end.is_some_and(|end| next > end) { + tracing::warn!(sequence, ?end, "a group fetch continued past FETCH_OK's End Location"); + return Err(Error::ProtocolViolation); + } + } + + // A clean FIN short of the promised end would otherwise cache a truncated group as + // a complete one. + if end.is_some_and(|end| next < end) { + tracing::warn!( + sequence, + next, + ?end, + "a group fetch ended before FETCH_OK's End Location" + ); + return Err(Error::ProtocolViolation); } Ok(()) @@ -3166,9 +3218,13 @@ struct FetchedObject { } /// Decode the next fetch object header, or `None` at stream end. +/// +/// The `first` object on a stream has no prior to inherit from, so a field it leaves to +/// the prior object is a protocol violation. async fn decode_fetch_object( stream: &mut Reader, version: Version, + first: bool, ) -> Result, Error> { if version == Version::Draft14 { let Some(group) = stream.decode_maybe::().await? else { @@ -3196,17 +3252,27 @@ async fn decode_fetch_object( subgroup, group, object, + priority, properties, - .. - }) => Some(FetchedObject { - group, - object, - subgroup_ok: matches!( - subgroup, - ietf::FetchSubgroup::Zero | ietf::FetchSubgroup::Prior | ietf::FetchSubgroup::Explicit(0) - ), - properties, - }), + }) => { + let inherits = group.is_none() + || object.is_none() + || priority.is_none() + || matches!(subgroup, ietf::FetchSubgroup::Prior | ietf::FetchSubgroup::PriorPlusOne); + if first && inherits { + tracing::warn!("the first fetch object refers to a prior object"); + return Err(Error::ProtocolViolation); + } + Some(FetchedObject { + group, + object, + subgroup_ok: matches!( + subgroup, + ietf::FetchSubgroup::Zero | ietf::FetchSubgroup::Prior | ietf::FetchSubgroup::Explicit(0) + ), + properties, + }) + } }) } @@ -6896,35 +6962,153 @@ mod stitch_tests { 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; + let mut run = GroupFetchRun::new(DRAFT, buf.to_vec()).await; + let res = run.recv_objects(0, None).await; assert!(matches!(res, Err(Error::ProtocolViolation)), "{res:?}"); } + + /// The first object on a fetch stream has no prior object, so leaving any field to + /// the prior one is a violation, not "the requested group, from its start". + #[tokio::test] + async fn a_group_fetch_refuses_a_first_object_that_inherits() { + use ietf::FetchSubgroup::{Prior, Zero}; + let header = |group, object, subgroup, priority| ietf::FetchObject::Object { + subgroup, + group, + object, + priority, + properties: None, + }; + let headers = [ + // Both IDs omitted: "the requested group, from its start" is exactly the bug. + header(None, None, Zero, Some(0)), + header(Some(SEQUENCE), None, Zero, Some(0)), + header(Some(SEQUENCE), Some(0), Prior, Some(0)), + header(Some(SEQUENCE), Some(0), Zero, None), + ]; + + for version in [Version::Draft15, VERSION] { + for header in &headers { + let mut buf = bytes::BytesMut::new(); + header.encode(&mut buf, version).unwrap(); + 1u64.encode(&mut buf, version).unwrap(); + buf.put_slice(b"a"); + + let mut run = GroupFetchRun::new(version, buf.to_vec()).await; + let res = run.recv_objects(0, None).await; + assert!( + matches!(res, Err(Error::ProtocolViolation)), + "{version} {header:?}: {res:?}" + ); + } + } + } + + /// FETCH_OK's End Location inside the group promises every object before it. A stream + /// that FINs short of it, or runs past it, fails the group instead of caching it. + #[tokio::test] + async fn a_group_fetch_must_reach_its_end_location() { + const END: u64 = 3; + for (count, complete) in [(2, false), (3, true), (4, false)] { + let payloads: Vec<&[u8]> = [b"a", b"b", b"c", b"d"][..count].iter().map(|p| &p[..]).collect(); + let mut run = GroupFetchRun::new(VERSION, group_fetch_objects(SEQUENCE, 0, &payloads)).await; + + let group = run.track.create_group(group::Info { sequence: SEQUENCE }).unwrap(); + let mut consumer = group.consume(); + let slot = kio::Producer::new(GroupFetch::Ready { + producer: group, + timescale: None, + start: 0, + end: Some(END), + }); + let res = run.subscriber.recv_group_fetch(&mut run.stream, slot).await; + assert_eq!(res.is_ok(), complete, "{count} objects: {res:?}"); + + let mut read = 0; + let end = loop { + match consumer.read_frame().await { + Ok(Some(_)) => read += 1, + Ok(None) => break Ok(read), + Err(err) => break Err(err), + } + }; + match complete { + true => assert_eq!(end.expect("a complete group"), END), + false => assert!(end.is_err(), "{count} objects: the group must fail, not end"), + } + } + } + + /// A group fetch's objects after the FETCH_HEADER: the first one names the group and + /// `start`, and every later one is the next object. + fn group_fetch_objects(sequence: u64, start: u64, payloads: &[&[u8]]) -> Vec { + let mut buf = bytes::BytesMut::new(); + for (index, payload) in payloads.iter().enumerate() { + let first = index == 0; + ietf::FetchObject::Object { + subgroup: ietf::FetchSubgroup::Zero, + group: first.then_some(sequence), + object: first.then_some(start), + priority: first.then_some(0), + properties: None, + } + .encode(&mut buf, VERSION) + .unwrap(); + (payload.len() as u64).encode(&mut buf, VERSION).unwrap(); + buf.put_slice(payload); + } + buf.to_vec() + } + + /// A subscriber reading one scripted group fetch stream, already past its header. + struct GroupFetchRun { + subscriber: Subscriber, + stream: Reader<::RecvStream, Version>, + track: track::Producer, + _tasks: (Tasks, TaskSet), + } + + impl GroupFetchRun { + async fn new(version: Version, objects: Vec) -> Self { + let mut session = ScriptedSession::per_stream_eof(vec![objects]); + 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, + version, + tasks.0.clone(), + Default::default(), + ); + let (_, recv) = session.open_bi().await.unwrap(); + let track = track::Producer::new( + std::sync::Arc::new(crate::broadcast::Info::default()), + "video", + track::Info::default(), + ); + + Self { + subscriber, + stream: Reader::new(recv, version), + track, + _tasks: tasks, + } + } + + /// Decode the stream into a fresh group numbered from `start`. + async fn recv_objects(&mut self, start: u64, end: Option) -> Result<(), Error> { + let mut group = self.track.create_group(group::Info { sequence: SEQUENCE }).unwrap(); + group.start_at(start).unwrap(); + self.subscriber + .recv_group_fetch_objects(&mut self.stream, &mut group, start, end, None) + .await + } + } } /// A scripted peer answering SUBSCRIBE then FETCH, so the join is spelled on the wire. @@ -6935,6 +7119,7 @@ mod joining_fetch_tests { coding::Encode as _, lite::test_transport::ScriptedSession, model::ProduceTest, + transport::poll::Session as _, util::{TaskSet, Tasks}, }; @@ -7235,4 +7420,122 @@ mod joining_fetch_tests { let writes = log.writes.lock().unwrap().clone(); assert!(writes.is_empty(), "a refused join must not write SUBSCRIBE"); } + + /// A cache miss for a group's tail asks upstream from the frame the reader wants and + /// numbers what arrives from there, so a publisher that evicted the prefix can answer. + #[tokio::test(start_paused = true)] + async fn a_group_fetch_asks_from_the_wanted_frame() { + const VERSION: Version = Version::Draft20; + const GROUP: u64 = 4; + const START: u64 = 2; + + let ok = message_bytes( + ietf::FetchOk::ID, + &ietf::FetchOk { + request_id: None, + group_order: GroupOrder::Ascending, + end_of_track: false, + end_location: ietf::Location { + group: GROUP + 1, + object: 0, + }, + }, + VERSION, + ); + + // The fetch stream answering our first request id, from object START on. + let mut objects = bytes::BytesMut::new(); + ietf::FetchHeader::TYPE.encode(&mut objects, VERSION).unwrap(); + ietf::FetchHeader { + request_id: RequestId(1), + } + .encode(&mut objects, VERSION) + .unwrap(); + for (index, payload) in [b"c", b"d"].iter().enumerate() { + let first = index == 0; + ietf::FetchObject::Object { + subgroup: ietf::FetchSubgroup::Zero, + group: first.then_some(GROUP), + object: first.then_some(START), + priority: first.then_some(0), + properties: None, + } + .encode(&mut objects, VERSION) + .unwrap(); + 1u64.encode(&mut objects, VERSION).unwrap(); + objects.extend_from_slice(&payload[..]); + } + + let session = ScriptedSession::per_stream_eof(vec![ok, objects.to_vec()]); + let (tasks, _task_set) = crate::util::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, + VERSION, + tasks, + Default::default(), + ); + + let track = track::Producer::new( + std::sync::Arc::new(crate::broadcast::Info::default()), + "video", + track::Info::default(), + ); + let dynamic = track.dynamic(); + let consumer = track.consume(); + let mut fetch = std::pin::pin!(consumer.fetch_group(GROUP, group::Fetch::default().with_frame_start(START))); + assert!(futures::poll!(fetch.as_mut()).is_pending()); + let request = dynamic.requested_group().await.expect("no group requested"); + + let serving = tokio::spawn(subscriber.clone().run_group_fetch( + Path::new("broadcast").to_owned(), + "video".into(), + request, + None, + )); + settle().await; + + // The fetch stream, as the peer would open it. + let (_, recv) = session.clone().open_bi().await.unwrap(); + let mut stream = Reader::new(recv, VERSION); + subscriber.clone().recv_fill(&mut stream).await.expect("group fetch"); + serving.await.expect("run_group_fetch"); + + let mut group = fetch.await.expect("fetched"); + assert_eq!(group.index(), START, "the group starts where the reader wanted"); + let mut payloads = Vec::new(); + while let Some(frame) = group.read_frame().await.expect("complete") { + payloads.push(frame.payload.to_vec()); + } + assert_eq!(payloads, [b"c".to_vec(), b"d".to_vec()]); + + let messages = decode_messages(&session.log, VERSION); + let fetch = messages.iter().find(|(id, _)| *id == ietf::Fetch::ID).expect("FETCH"); + let mut body = fetch.1.clone(); + let msg = ietf::Fetch::decode_msg(&mut body, VERSION).unwrap(); + let FetchType::Standalone { start, end, .. } = msg.fetch_type else { + panic!("a group fetch is standalone: {:?}", msg.fetch_type); + }; + assert_eq!( + start, + ietf::Location { + group: GROUP, + object: START + } + ); + assert_eq!( + end, + ietf::Location { + group: GROUP, + object: 0 + }, + "through the end of the group" + ); + } }