From c650bc92b938bb102ca3a6a47a99bd27f9ce0446 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 10:47:49 -0700 Subject: [PATCH 1/5] quest(moxygen): claim datagram From 78d6eefebcd0a0f72cc2be3cf372f559acea2d19 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 11:37:22 -0700 Subject: [PATCH 2/5] feat(moq-net): carry datagrams over moq-transport as OBJECT_DATAGRAM The IETF session never read or wrote QUIC datagrams, so a datagram track delivered nothing over moq-transport. Copy the moq-lite path: the publisher sends each model datagram as an OBJECT_DATAGRAM at object 0 that ends its group, and the subscriber inserts one back at its Group ID so a relay keeps the sequence. An object past 0 or a non-Normal status is dropped. Completes quest/m1/moxygen/datagram.md. Co-Authored-By: Claude Opus 5.5 --- doc/concept/standard.md | 5 + doc/lib/rs/moq-net.md | 2 +- go/wrapper/README.md | 2 +- kt/README.md | 2 +- py/moq-rs/README.md | 2 +- quest/m1/moxygen/README.md | 1 - quest/m1/moxygen/datagram.md | 23 --- rs/moq-net/src/fuzz.rs | 3 +- rs/moq-net/src/ietf/datagram.rs | 268 ++++++++++++++++++++++++++++++ rs/moq-net/src/ietf/mod.rs | 2 + rs/moq-net/src/ietf/publisher.rs | 58 ++++++- rs/moq-net/src/ietf/session.rs | 25 +++ rs/moq-net/src/ietf/subscriber.rs | 116 +++++++++++++ rs/moq-net/src/model/datagram.rs | 6 +- rs/moq-net/src/model/track.rs | 5 +- rs/moq-net/tests/datagram.rs | 102 +++++++----- swift/README.md | 2 +- 17 files changed, 546 insertions(+), 78 deletions(-) delete mode 100644 quest/m1/moxygen/datagram.md create mode 100644 rs/moq-net/src/ietf/datagram.rs diff --git a/doc/concept/standard.md b/doc/concept/standard.md index f2c7e9929e..8ca241e0fa 100644 --- a/doc/concept/standard.md +++ b/doc/concept/standard.md @@ -41,6 +41,11 @@ 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 tracks waiting for a group that never arrives. +A moq-lite datagram is a single-frame group, so on moq-transport it travels +as an `OBJECT_DATAGRAM` at object 0 whose Group ID is the sequence, and a relay +forwards it without renumbering. A datagram carrying any other Object ID, or a +status other than Normal, is dropped. + Several project drafts extend the IETF wire without breaking it, since `SETUP` ignores unknown parameters: [cluster](/draft/moq-cluster) routing hop lists, [solicit](/draft/moq-solicit) to make announcements opt-in, diff --git a/doc/lib/rs/moq-net.md b/doc/lib/rs/moq-net.md index bd775dbe4b..e9ee8b442d 100644 --- a/doc/lib/rs/moq-net.md +++ b/doc/lib/rs/moq-net.md @@ -21,7 +21,7 @@ above ([hang](/lib/rs/hang)); relays and CDNs implement only this. - **Tracks** carry groups with a priority, a retention window, and a timescale. Subscribers set their own priority and max age and can change them live. - **Groups** are written frame by frame and delivered on independent streams. Old groups are cached for fetch-by-sequence; stale groups are skipped per the subscriber's budget. - **Track ends**: `finish()` ends a track at its live edge, while `finish_at(n)` declares the exclusive end ahead of it and still accepts the groups below. A subscriber awaits it with `finished()`. A remote track ends only once every group below its end has arrived or was dropped; one reset before its header arrived is skipped after the subscription's max age on moq-lite (one second without one), or after one second on IETF. -- **Datagrams** send a single small frame unreliably on moq-lite 05+. +- **Datagrams** send a single small frame unreliably on moq-lite 05+ and moq-transport. - **Routes** record the relay hops and a cost, which is what the relay [cluster](/bin/relay/cluster) routes on. A hop of 0 marks the chain anonymous: `Route::is_anonymous()` is true, and that route ranks below every fully identified one. `Route::source()` says where a delivered route entered: `Source::Local`, or `Source::Peer(hop)` when a handle marked `origin::Producer::peer()` announced it. `origin::Consumer::local()` sees only the local ones. - **Stats** counters per broadcast and session, drained by [`moq-stats`](https://docs.rs/moq-stats). diff --git a/go/wrapper/README.md b/go/wrapper/README.md index cf241671b5..ccc1f77ea9 100644 --- a/go/wrapper/README.md +++ b/go/wrapper/README.md @@ -110,7 +110,7 @@ Raw tracks support best-effort datagrams alongside groups: `TrackProducer.Append sends one `Frame` and returns its sequence number, while `TrackConsumer.RecvDatagram` and `TrackConsumer.Datagrams` receive them in arrival order. Payloads are capped at 1200 bytes. Datagram delivery requires a datagram-capable transport and lite-05 or -newer moq-lite; IETF moq-transport, pre-lite-05, WebSocket, and TCP paths do not +newer moq-lite, or moq-transport; pre-lite-05, WebSocket, and TCP paths do not deliver them, and there is no stream fallback. ## Versioning diff --git a/kt/README.md b/kt/README.md index 62c509e960..6990712969 100644 --- a/kt/README.md +++ b/kt/README.md @@ -52,7 +52,7 @@ The `dev.moq` package is intentionally thin: Kotlin has extension functions, so - **Fetched media**: `fetchMediaGroup(...).frames()` streams the decoded frames of one retained group, then completes. - **Duration extensions** (`Durations.kt`): the FFI carries microseconds as integers, so `stats.rtt`, `backoff.initial`, `frame.timestamp`, and their siblings read back as a `kotlin.time.Duration`. - **`logLevel(...)`**: configures native Rust tracing without importing the raw bindings package. -- **Raw datagrams**: `TrackProducer.appendDatagram(Frame(payload, timestampUs))` sends one best-effort frame and returns its sequence; `TrackConsumer.recvDatagram()` and `datagrams()` receive them. Payloads are capped at 1200 bytes, require a datagram-capable transport plus lite-05 or newer moq-lite, and have no stream fallback. +- **Raw datagrams**: `TrackProducer.appendDatagram(Frame(payload, timestampUs))` sends one best-effort frame and returns its sequence; `TrackConsumer.recvDatagram()` and `datagrams()` receive them. Payloads are capped at 1200 bytes, require a datagram-capable transport plus lite-05 or newer moq-lite or moq-transport, and have no stream fallback. - **`MoqException.isShutdown`** (`Errors.kt`): true for the graceful `Cancelled`/`Closed` cases. ## Versioning diff --git a/py/moq-rs/README.md b/py/moq-rs/README.md index 0902031a5a..90281ccc6d 100644 --- a/py/moq-rs/README.md +++ b/py/moq-rs/README.md @@ -212,7 +212,7 @@ Every handle whose cleanup is `cancel()` is an async context manager, so exiting - **`Catalog`**. `.audio: dict[str, Audio]`, `.video: dict[str, Video]`, `.display`, `.rotation`, `.flip`. - **`Frame`**. `.payload: bytes`, `.timestamp_us: int`. The unit of every write and every raw read. - **`MediaFrame`**. `.payload: bytes`, `.timestamp_us: int`, `.keyframe: bool`. Returned by media subscriptions. `keyframe` marks a group start or video keyframe; for audio it is true only at a group start. -- **`Datagram`**. `.sequence: int`, `.timestamp_us: int`, `.payload: bytes`. Delivered only on datagram-capable transports and lite-05 or newer moq-lite. +- **`Datagram`**. `.sequence: int`, `.timestamp_us: int`, `.payload: bytes`. Delivered only on datagram-capable transports with lite-05 or newer moq-lite, or moq-transport. - **`Audio`**. `.codec`, `.sample_rate`, `.channel_count`, `.bitrate`, `.description`. - **`Video`**. `.codec`, `.coded: Dimensions`, `.display_aspect`, `.bitrate`, `.stalled`, `.framerate`, `.description`. A true `.stalled` recommends temporarily avoiding the rendition without making it unavailable. - **`Subscription`**. Subscriber delivery preferences: priority, staleness, and optional group range. diff --git a/quest/m1/moxygen/README.md b/quest/m1/moxygen/README.md index d7d7e72434..097e85216a 100644 --- a/quest/m1/moxygen/README.md +++ b/quest/m1/moxygen/README.md @@ -32,7 +32,6 @@ Docs stay inline in the change that makes them stale. No new guide. - [Default track priority](/quest/m1/moxygen/priority.md) - an unset track priority is the midpoint on moq-lite and on IETF, not the least urgent value - [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 -- [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/datagram.md b/quest/m1/moxygen/datagram.md deleted file mode 100644 index c794b5e57d..0000000000 --- a/quest/m1/moxygen/datagram.md +++ /dev/null @@ -1,23 +0,0 @@ -# [M] Datagram groups - -## Goal - -An IETF datagram that carries one object for a group is delivered as a -moq-lite datagram: one single-frame group, sequence preserved. Several -objects in one datagram group stay unsupported. - -## Plan - -moq-lite already does this. A datagram is a subscribe id, a group sequence, a -timestamp, and a payload. `insert_datagram` keeps the sequence so a relay -does not renumber it. - -The IETF session does not read or write QUIC datagrams, so a datagram -subscribe delivers nothing. Copy the lite path onto that session. Do not -invent a second object list inside the group. - -The moxygen cases with one object per group are the ones this can pass. - -## Related - -- [Moxygen compatibility](/quest/m1/moxygen/README.md) - the line this belongs to diff --git a/rs/moq-net/src/fuzz.rs b/rs/moq-net/src/fuzz.rs index 5fc9f793b6..66cb16c1f0 100644 --- a/rs/moq-net/src/fuzz.rs +++ b/rs/moq-net/src/fuzz.rs @@ -61,7 +61,7 @@ const IETF_VERSIONS: &[ietf::Version] = &[ const LITE_KINDS: u8 = 21; /// How many types [`ietf_wire`] dispatches over. -const IETF_KINDS: u8 = 39; +const IETF_KINDS: u8 = 40; /// Split the two selector bytes off the input: a version and a type. fn select(data: &[u8], versions: usize) -> Option<(usize, u8, &[u8])> { @@ -318,6 +318,7 @@ pub fn ietf_wire(data: &[u8]) -> bool { 36 => roundtrip::(rest, version, stable), 37 => roundtrip::(rest, version, stable), 38 => roundtrip::(rest, version, stable), + 39 => roundtrip::(rest, version, stable), _ => unreachable!("kind is taken modulo IETF_KINDS"), } } diff --git a/rs/moq-net/src/ietf/datagram.rs b/rs/moq-net/src/ietf/datagram.rs new file mode 100644 index 0000000000..04815eb680 --- /dev/null +++ b/rs/moq-net/src/ietf/datagram.rs @@ -0,0 +1,268 @@ +//! OBJECT_DATAGRAM: one Object carried in a QUIC datagram (draft-14 section 10.3.1 through +//! draft-20 section 11.3.1). +//! +//! The model counterpart is [`crate::Datagram`], a single-frame group, so only an Object at +//! ID 0 maps onto it. The Type is a set of flags on every draft; draft-14 lacks the +//! DEFAULT_PRIORITY bit and a status with an omitted Object ID. + +use bytes::{Buf, BufMut, Bytes}; + +use crate::coding::{Decode, DecodeError, Encode, EncodeError}; + +use super::Version; + +/// The bits of an OBJECT_DATAGRAM Type. +mod flag { + pub const PROPERTIES: u64 = 0x01; + pub const END_OF_GROUP: u64 = 0x02; + pub const ZERO_OBJECT_ID: u64 = 0x04; + pub const DEFAULT_PRIORITY: u64 = 0x08; + pub const STATUS: u64 = 0x20; + /// Every defined bit. Anything else, including the reserved 0x10, is invalid. + pub const ALL: u64 = PROPERTIES | END_OF_GROUP | ZERO_OBJECT_ID | DEFAULT_PRIORITY | STATUS; +} + +/// What follows an OBJECT_DATAGRAM's header. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum DatagramBody { + /// The Object Payload, delimited by the datagram boundary. + Payload(Bytes), + /// The Object Status of an Object without a payload. + Status(u64), +} + +/// A decoded OBJECT_DATAGRAM. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ObjectDatagram { + pub track_alias: u64, + pub group_id: u64, + /// The Object ID, or `None` when the ZERO_OBJECT_ID bit omits it (Object 0). + pub object_id: Option, + /// The Publisher Priority, or `None` to inherit the subscription's (draft-15+). + pub publisher_priority: Option, + /// No Object past this one exists in the group. + pub end_of_group: bool, + /// The Object Properties block without its length prefix, which carries the Timestamp. + pub properties: Option>, + pub body: DatagramBody, +} + +impl ObjectDatagram { + /// Whether `kind` is a Type this draft defines. + fn valid(kind: u64, version: Version) -> bool { + match version { + Version::Draft14 => kind <= 0x07 || kind == 0x20 || kind == 0x21, + // A status cannot also end the group. + _ => kind & !flag::ALL == 0 && !(kind & flag::STATUS != 0 && kind & flag::END_OF_GROUP != 0), + } + } +} + +impl Encode for ObjectDatagram { + fn encode(&self, w: &mut W, version: Version) -> Result<(), EncodeError> { + let mut kind = 0; + if self.properties.is_some() { + kind |= flag::PROPERTIES; + } + if self.end_of_group { + kind |= flag::END_OF_GROUP; + } + if self.object_id.is_none() { + kind |= flag::ZERO_OBJECT_ID; + } + if self.publisher_priority.is_none() { + kind |= flag::DEFAULT_PRIORITY; + } + if matches!(self.body, DatagramBody::Status(_)) { + kind |= flag::STATUS; + } + if !Self::valid(kind, version) { + return Err(EncodeError::InvalidState); + } + + kind.encode(w, version)?; + self.track_alias.encode(w, version)?; + self.group_id.encode(w, version)?; + if let Some(object_id) = self.object_id { + object_id.encode(w, version)?; + } + if let Some(priority) = self.publisher_priority { + priority.encode(w, version)?; + } + if let Some(properties) = &self.properties { + // A present but empty block is a protocol violation for the peer. + if properties.is_empty() { + return Err(EncodeError::InvalidState); + } + properties.encode(w, version)?; + } + + match &self.body { + DatagramBody::Status(status) => status.encode(w, version)?, + DatagramBody::Payload(payload) => { + // Runs to the datagram boundary: written raw, no length prefix. + if w.remaining_mut() < payload.len() { + return Err(EncodeError::Short); + } + w.put_slice(payload); + } + } + Ok(()) + } +} + +impl Decode for ObjectDatagram { + fn decode(r: &mut R, version: Version) -> Result { + let kind = u64::decode(r, version)?; + if !Self::valid(kind, version) { + return Err(DecodeError::InvalidValue); + } + + let track_alias = u64::decode(r, version)?; + let group_id = u64::decode(r, version)?; + let object_id = match kind & flag::ZERO_OBJECT_ID != 0 { + true => None, + false => Some(u64::decode(r, version)?), + }; + let publisher_priority = match kind & flag::DEFAULT_PRIORITY != 0 { + true => None, + false => Some(u8::decode(r, version)?), + }; + let properties = match kind & flag::PROPERTIES != 0 { + true => { + let properties = Vec::::decode(r, version)?; + if properties.is_empty() { + return Err(DecodeError::InvalidValue); + } + Some(properties) + } + false => None, + }; + + let body = match kind & flag::STATUS != 0 { + true => { + let status = u64::decode(r, version)?; + // Draft-17 on: only a Normal Object may carry Properties. + let legacy = matches!(version, Version::Draft14 | Version::Draft15 | Version::Draft16); + if !legacy && status != 0 && properties.is_some() { + return Err(DecodeError::InvalidValue); + } + if r.has_remaining() { + return Err(DecodeError::TrailingBytes); + } + DatagramBody::Status(status) + } + false => DatagramBody::Payload(r.copy_to_bytes(r.remaining())), + }; + + Ok(Self { + track_alias, + group_id, + object_id, + publisher_priority, + end_of_group: kind & flag::END_OF_GROUP != 0, + properties, + body, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const ALL: [Version; 9] = [ + Version::Draft14, + Version::Draft15, + Version::Draft16, + Version::Draft17, + Version::Draft18, + Version::Draft19, + Version::Draft20, + Version::Draft21, + Version::Draft22, + ]; + + fn single(priority: Option) -> ObjectDatagram { + ObjectDatagram { + track_alias: 3, + group_id: 42, + object_id: None, + publisher_priority: priority, + end_of_group: true, + properties: Some(vec![0x10, 0x05]), + body: DatagramBody::Payload(Bytes::from_static(b"hello")), + } + } + + #[test] + fn roundtrip_every_draft() { + for version in ALL { + let datagram = single(Some(7)); + let mut buf = datagram.encode_bytes(version).unwrap(); + assert_eq!(buf[0], 0x07, "{version}: properties, end of group, object 0"); + let decoded = ObjectDatagram::decode(&mut buf, version).unwrap(); + assert_eq!(decoded, datagram, "{version}"); + assert!(!buf.has_remaining()); + } + } + + #[test] + fn default_priority_needs_draft15() { + assert!(matches!( + single(None).encode_bytes(Version::Draft14), + Err(EncodeError::InvalidState) + )); + assert!(ObjectDatagram::decode(&mut &[0x0C, 0x01, 0x02][..], Version::Draft14).is_err()); + + let mut buf = single(None).encode_bytes(Version::Draft16).unwrap(); + assert_eq!(buf[0], 0x0F); + assert_eq!( + ObjectDatagram::decode(&mut buf, Version::Draft16).unwrap(), + single(None) + ); + } + + #[test] + fn explicit_object_id_and_status() { + let datagram = ObjectDatagram { + track_alias: 1, + group_id: 2, + object_id: Some(0), + publisher_priority: Some(128), + end_of_group: false, + properties: None, + body: DatagramBody::Status(0), + }; + for version in ALL { + let mut buf = datagram.encode_bytes(version).unwrap(); + assert_eq!(buf[0], 0x20, "{version}"); + assert_eq!(ObjectDatagram::decode(&mut buf, version).unwrap(), datagram); + } + } + + #[test] + fn rejects_invalid_types() { + for version in ALL { + // A status that ends the group, the reserved bit, and a bit past the defined ones. + for kind in [0x22u64, 0x10, 0x40] { + let mut bytes = kind.encode_bytes(version).unwrap().to_vec(); + bytes.extend_from_slice(&[0x01, 0x02, 0x03, 0x04]); + assert!( + ObjectDatagram::decode(&mut &bytes[..], version).is_err(), + "{version}: {kind:#x}" + ); + } + } + } + + #[test] + fn rejects_empty_properties() { + // PROPERTIES and ZERO_OBJECT_ID, alias 1, group 2, priority 0, empty block. + let bytes = [0x05, 0x01, 0x02, 0x00, 0x00, b'x']; + assert!(matches!( + ObjectDatagram::decode(&mut &bytes[..], Version::Draft16), + Err(DecodeError::InvalidValue) + )); + } +} diff --git a/rs/moq-net/src/ietf/mod.rs b/rs/moq-net/src/ietf/mod.rs index 8ae7dac1ec..620b5d7fe6 100644 --- a/rs/moq-net/src/ietf/mod.rs +++ b/rs/moq-net/src/ietf/mod.rs @@ -9,6 +9,7 @@ mod parameters; mod adapter; pub mod cluster; mod control; +mod datagram; pub(crate) mod error; mod fetch; mod filter; @@ -34,6 +35,7 @@ mod track; mod version; use control::Control; +pub use datagram::*; pub use fetch::*; pub use filter::*; pub use goaway::*; diff --git a/rs/moq-net/src/ietf/publisher.rs b/rs/moq-net/src/ietf/publisher.rs index 6177770669..44477e366d 100644 --- a/rs/moq-net/src/ietf/publisher.rs +++ b/rs/moq-net/src/ietf/publisher.rs @@ -15,7 +15,7 @@ use web_transport_trait::poll::SendStream as _; use crate::{ AsPath, Error, Timescale, Timestamp, - coding::{Stream, Writer}, + coding::{Encode as _, Stream, Writer}, ietf::{self, Control, EndLocation, FetchHeader, FetchType, Filter, GroupOrder, Location, RequestId}, track::Subscription, util::{MaybeBoxedExt, MaybeSendBox}, @@ -1960,6 +1960,9 @@ struct TrackServe { opened: Arc, /// The track's exclusive end, once its groups ran out because it finished. end: Option, + /// Serve the track's datagrams too, as OBJECT_DATAGRAMs. Off when the transport has no + /// datagrams; there is no stream fallback. + datagrams: bool, } impl TrackServe { @@ -1980,8 +1983,10 @@ impl TrackServe { } } track.end_at(range.end.map_or(Bound::Unbounded, |end| Bound::Included(end.group))); + let datagrams = session.max_datagram_size() > 0; Self { + datagrams, session, track, request_id, @@ -2093,6 +2098,8 @@ impl TrackServe { ); } Poll::Ready(Ok(None)) => { + // Datagrams written before the track finished still go out. + self.poll_datagrams(waiter); self.draining = true; if let Poll::Ready(Ok(end)) = self.track.poll_finished(waiter) { self.end = Some(end); @@ -2103,10 +2110,59 @@ impl TrackServe { Poll::Pending => break, } } + // Groups first, so a burst of datagrams cannot starve them. + self.poll_datagrams(waiter); // Newly created group machines start now rather than on the next wake. let _ = self.children.poll(waiter); Poll::Pending } + + /// Send every buffered datagram as an OBJECT_DATAGRAM, best-effort, like moq-lite. + /// + /// A datagram is Object 0 of a group that has no other, so the Object ends its group. + /// The track's end or failure surfaces through its groups, so this only stops. + fn poll_datagrams(&mut self, waiter: &kio::Waiter) { + if !self.datagrams { + return; + } + while let Poll::Ready(Ok(Some(datagram))) = self.track.poll_recv_datagram(waiter) { + let sequence = datagram.sequence; + let properties = match self.timescale { + Some(timescale) => { + let mut properties = Vec::new(); + if ietf::encode_object_time(&mut properties, datagram.timestamp, timescale, self.version).is_err() { + continue; + } + Some(properties) + } + None => None, + }; + let body = ietf::ObjectDatagram { + track_alias: self.request_id.0, + group_id: sequence, + object_id: None, + publisher_priority: Some(super::priority::to_wire(self.track.info().priority)), + end_of_group: true, + properties, + body: ietf::DatagramBody::Payload(datagram.payload), + }; + let Ok(body) = body.encode_bytes(self.version) else { + continue; + }; + + let max = self.session.max_datagram_size(); + if body.len() > max { + tracing::debug!( + sequence, + size = body.len(), + max, + "dropping datagram larger than the transport limit" + ); + continue; + } + let _ = self.session.send_datagram(&body); + } + } } /// Serves one group on its own unidirectional stream in the moq-transport diff --git a/rs/moq-net/src/ietf/session.rs b/rs/moq-net/src/ietf/session.rs index cb0fd4538a..9887fb99ca 100644 --- a/rs/moq-net/src/ietf/session.rs +++ b/rs/moq-net/src/ietf/session.rs @@ -200,6 +200,7 @@ where subscriber.clone(), version ))); + let mut datagrams = std::pin::pin!(err_only(run_datagrams(adapter.clone(), subscriber.clone()))); // Unsolicited PUBLISH_NAMESPACE unless the peer requires solicitation; // see `Publisher::run_publish_namespaces`. let mut pub_ns_run = std::pin::pin!(err_only(publisher.clone().run_publish_namespaces())); @@ -250,6 +251,9 @@ where if let Poll::Ready(err) = waiter.poll_future(dispatch.as_mut()) { return Poll::Ready(Err(err)); } + if let Poll::Ready(err) = waiter.poll_future(datagrams.as_mut()) { + return Poll::Ready(Err(err)); + } if task_set.poll(waiter).is_ready() { return Poll::Ready(Ok(())); } @@ -340,6 +344,7 @@ where subscriber.clone(), version ))); + let mut datagrams = std::pin::pin!(err_only(run_datagrams(session.clone(), subscriber.clone()))); let mut goaway_recv = std::pin::pin!(err_only(goaway_recv)); let mut setup = std::pin::pin!(setup); // Unsolicited PUBLISH_NAMESPACE unless the peer requires solicitation; @@ -379,6 +384,9 @@ where if let Poll::Ready(err) = waiter.poll_future(dispatch.as_mut()) { return Poll::Ready(Err(err)); } + if let Poll::Ready(err) = waiter.poll_future(datagrams.as_mut()) { + return Poll::Ready(Err(err)); + } if let Poll::Ready(err) = waiter.poll_future(goaway_recv.as_mut()) { return Poll::Ready(Err(err)); } @@ -709,6 +717,23 @@ where } } +/// Receive QUIC datagrams, each an OBJECT_DATAGRAM for one of our subscriptions. +/// +/// A transport without datagrams never delivers one, so this parks. A transport failure +/// or a malformed datagram ends the session. +async fn run_datagrams(mut session: S, subscriber: Subscriber) -> Result<(), Error> +where + S: crate::transport::poll::Boxable, +{ + if session.max_datagram_size() == 0 { + return Ok(()); + } + loop { + let payload = session.recv_datagram().await.map_err(Error::from_transport)?; + subscriber.recv_datagram(payload)?; + } +} + async fn run_uni_group( subscriber: &mut Subscriber, stream: &mut Reader, diff --git a/rs/moq-net/src/ietf/subscriber.rs b/rs/moq-net/src/ietf/subscriber.rs index b2cea32c26..4d434e24c0 100644 --- a/rs/moq-net/src/ietf/subscriber.rs +++ b/rs/moq-net/src/ietf/subscriber.rs @@ -2279,6 +2279,61 @@ where Ok(()) } + + /// Deliver one OBJECT_DATAGRAM as a datagram on its subscription's track: a + /// single-frame group at the Group ID. + /// + /// A malformed datagram is the peer breaking the protocol, so it errors. One the model + /// cannot carry is dropped like any lost datagram: an Object past ID 0 (the group would + /// need a second object), a status other than Normal, or an alias that is not bound + /// yet (the draft lets us drop rather than buffer). + pub fn recv_datagram(&self, payload: bytes::Bytes) -> Result<(), Error> { + let mut buf = payload; + let datagram = ietf::ObjectDatagram::decode(&mut buf, self.version)?; + let (alias, sequence) = (datagram.track_alias, datagram.group_id); + + if datagram.object_id.unwrap_or(0) != 0 { + tracing::debug!(alias, sequence, "dropping a datagram past object 0"); + return Ok(()); + } + let payload = match datagram.body { + ietf::DatagramBody::Payload(payload) => payload, + ietf::DatagramBody::Status(0) => bytes::Bytes::new(), + ietf::DatagramBody::Status(status) => { + tracing::debug!(alias, sequence, status, "dropping a datagram status"); + return Ok(()); + } + }; + + let mut state = self.state.lock(); + let request_id = match state.aliases.read().map.get(&alias) { + Some(Alias::Active(request_id)) => *request_id, + _ => { + tracing::debug!(alias, sequence, "dropping a datagram for an unbound alias"); + return Ok(()); + } + }; + let Some(track) = state.subscribes.get_mut(&request_id) else { + return Ok(()); + }; + + // Like a subgroup object: a track that declared no timescale is stamped on arrival. + let timestamp = match (track.timescale, &datagram.properties) { + (Some(timescale), Some(properties)) => { + ietf::decode_object_time(&mut properties.as_slice(), timescale, self.version)? + } + _ => None, + }; + let timestamp = timestamp.unwrap_or_else(|| crate::Timestamp::from(self.runtime.now())); + + let Some(producer) = track.producer.as_mut() else { + return Ok(()); + }; + if let Err(err) = producer.insert_datagram(sequence, timestamp, payload) { + tracing::debug!(%err, alias, sequence, "dropping datagram"); + } + Ok(()) + } } /// Mark where the track ends, as an END_OF_TRACK object said. @@ -3464,6 +3519,67 @@ mod tests { ); } + /// An OBJECT_DATAGRAM at object 0 is a datagram group at its Group ID; anything the model + /// cannot carry as one is dropped, and a malformed one is the peer's violation. + #[tokio::test] + async fn an_object_datagram_is_a_datagram_group() { + use crate::coding::Encode as _; + use futures::FutureExt as _; + + let subscriber = subscriber_with_tracks(&[(RequestId(11), "cam", "audio")]); + subscriber.register_alias(RequestId(11), 7).unwrap(); + let mut consumer = { + let mut state = subscriber.state.lock(); + let track = state.subscribes.get_mut(&RequestId(11)).unwrap(); + track.timescale = Some(Timescale::default()); + track.producer.as_ref().unwrap().subscribe(None) + }; + + let timestamp = crate::Timestamp::new(96_000, Timescale::default()).unwrap(); + let datagram = |alias: u64, group_id: u64, object_id: Option, body: ietf::DatagramBody| { + let mut properties = Vec::new(); + ietf::encode_object_time(&mut properties, timestamp, Timescale::default(), Version::Draft19).unwrap(); + ietf::ObjectDatagram { + track_alias: alias, + group_id, + object_id, + publisher_priority: None, + // Only a Normal Object may carry Properties, and a status cannot end the group. + end_of_group: matches!(body, ietf::DatagramBody::Payload(_)), + properties: matches!(body, ietf::DatagramBody::Payload(_)).then_some(properties), + body, + } + .encode_bytes(Version::Draft19) + .unwrap() + }; + let payload = |bytes: &'static [u8]| ietf::DatagramBody::Payload(bytes::Bytes::from_static(bytes)); + + // Dropped: a second object in the group, an unbound alias, and a status. + subscriber + .recv_datagram(datagram(7, 4, Some(1), payload(b"no"))) + .unwrap(); + subscriber.recv_datagram(datagram(8, 4, None, payload(b"no"))).unwrap(); + subscriber + .recv_datagram(datagram(7, 4, None, ietf::DatagramBody::Status(END_OF_TRACK))) + .unwrap(); + + subscriber + .recv_datagram(datagram(7, 9, Some(0), payload(b"yes"))) + .unwrap(); + let received = consumer.recv_datagram().now_or_never().unwrap().unwrap().unwrap(); + assert_eq!(received.sequence, 9, "the Group ID is the sequence"); + assert_eq!(received.timestamp, timestamp); + assert_eq!(&received.payload[..], b"yes"); + assert!( + consumer.recv_datagram().now_or_never().is_none(), + "only one got through" + ); + + // A status datagram cannot end the group. + let malformed = bytes::Bytes::from_static(&[0x22, 0x07, 0x04, 0x00]); + assert!(is_protocol_violation(&subscriber.recv_datagram(malformed).unwrap_err())); + } + /// One alias naming two different tracks is the collision section 11.1 makes fatal. #[test] fn an_alias_reused_for_another_track_is_fatal() { diff --git a/rs/moq-net/src/model/datagram.rs b/rs/moq-net/src/model/datagram.rs index f305507207..a8c20b0979 100644 --- a/rs/moq-net/src/model/datagram.rs +++ b/rs/moq-net/src/model/datagram.rs @@ -8,10 +8,10 @@ //! //! Delivery is best-effort per hop: a session drops (with a debug log) any datagram whose encoded //! body exceeds the transport's datagram size, and sessions that can't carry datagrams at all -//! (IETF moq-transport, moq-lite before 05, or stream-only transports like WebSocket) never -//! deliver them. +//! (moq-lite before 05, or stream-only transports like WebSocket) never deliver them. //! -//! Wire counterpart: [`crate::lite::Datagram`]. +//! Wire counterparts: [`crate::lite::Datagram`], and on moq-transport an OBJECT_DATAGRAM at +//! object 0 whose Group ID is the sequence ([`crate::ietf::ObjectDatagram`]). use bytes::Bytes; diff --git a/rs/moq-net/src/model/track.rs b/rs/moq-net/src/model/track.rs index d4f3216a6d..c3309c4fe4 100644 --- a/rs/moq-net/src/model/track.rs +++ b/rs/moq-net/src/model/track.rs @@ -1337,9 +1337,8 @@ impl Producer { /// track's groups but drawing from the same sequence namespace (so interleaving with /// [`Self::append_group`] never reuses a number). There is no group fallback: each /// session drops (with a debug log) any datagram whose encoded body exceeds the - /// transport's datagram size, and sessions that can't carry datagrams at all (IETF - /// moq-transport, moq-lite before 05, or stream-only transports like WebSocket) never - /// deliver them. Keep payloads well under the 1200-byte minimum path MTU. An origin + /// transport's datagram size, and sessions that can't carry datagrams at all (moq-lite + /// before 05, or stream-only transports like WebSocket) never deliver them. Keep payloads well under the 1200-byte minimum path MTU. An origin /// publisher uses this; a relay preserving upstream numbering uses /// [`Self::insert_datagram`]. pub fn append_datagram(&mut self, timestamp: Timestamp, payload: B) -> Result { diff --git a/rs/moq-net/tests/datagram.rs b/rs/moq-net/tests/datagram.rs index 4cc78b1297..d974cc580e 100644 --- a/rs/moq-net/tests/datagram.rs +++ b/rs/moq-net/tests/datagram.rs @@ -1,4 +1,4 @@ -//! MoQ Lite datagram delivery over the in-memory mock transport. +//! Datagram delivery over the in-memory mock transport, on moq-lite and moq-transport. //! //! Covers the whole receive path end to end: publisher encoding, the transport's //! datagram channel, the subscriber's receive loop and `route_datagram`, and @@ -9,7 +9,6 @@ mod support; use std::time::Duration; -use futures::FutureExt as _; use moq_net::{Hop, Timestamp, Version}; use support::harness::{MockConnectOptions, MockPair, connect_mock}; @@ -110,53 +109,74 @@ async fn datagrams_reach_the_subscriber_in_order() { .expect("timed out"); } -/// MoQ Transport has no datagram mapping: groups still flow, inserted datagrams do not. +/// MoQ Transport carries a datagram as an OBJECT_DATAGRAM at object 0, so it arrives as a +/// datagram with its sequence, alongside the groups on streams. #[tokio::test] -async fn ietf_does_not_deliver_datagrams() { - tokio::time::timeout(TEST_TIMEOUT, async { - let publisher = produce_origin(1); - let consumer_origin = produce_origin(2); +async fn ietf_delivers_datagrams() { + for version in [ + "moq-transport-14", + "moq-transport-16", + "moq-transport-17", + "moq-transport-20", + ] { + tokio::time::timeout(TEST_TIMEOUT, ietf_delivers_datagrams_on(version)) + .await + .unwrap_or_else(|_| panic!("{version}: timed out")); + } +} - let broadcast = publisher.create_broadcast("bench").unwrap(); - let mut producer = broadcast.create_track("datagrams", None).unwrap(); - broadcast.announce(Default::default()).unwrap(); +async fn ietf_delivers_datagrams_on(version: &str) { + let publisher = produce_origin(1); + let consumer_origin = produce_origin(2); - let mut options = MockConnectOptions::new("moq-transport-19".parse::().unwrap()); - options.server_publish = Some(publisher.consume()); - options.client_subscribe = Some(consumer_origin.clone()); - let _pair = connect_mock(options).await; + let broadcast = publisher.create_broadcast("bench").unwrap(); + let mut producer = broadcast.create_track("datagrams", None).unwrap(); + broadcast.announce(Default::default()).unwrap(); - let consumer = consumer_origin.consume(); - consumer.routed("bench").await.unwrap(); - let remote = consumer.request_broadcast("bench").await.unwrap(); - let mut subscriber = remote.track("datagrams").unwrap().subscribe(None).await.unwrap(); + let mut options = MockConnectOptions::new(version.parse::().unwrap()); + options.server_publish = Some(publisher.consume()); + options.client_subscribe = Some(consumer_origin.clone()); + let _pair = connect_mock(options).await; - producer - .write_frame(Timestamp::from_millis(0).unwrap(), &b"before"[..]) - .unwrap(); - let before = subscriber.recv_group().await.unwrap().unwrap(); - assert_eq!(before.sequence, 0); + let consumer = consumer_origin.consume(); + consumer.routed("bench").await.unwrap(); + let remote = consumer.request_broadcast("bench").await.unwrap(); + let mut subscriber = remote.track("datagrams").unwrap().subscribe(None).await.unwrap(); - producer - .insert_datagram( - 0, - Timestamp::from_millis(7).unwrap(), - bytes::Bytes::from_static(PAYLOAD), - ) - .unwrap(); - producer - .write_frame(Timestamp::from_millis(1).unwrap(), &b"after"[..]) - .unwrap(); - let after = subscriber.recv_group().await.unwrap().unwrap(); - assert_eq!(after.sequence, 1); + // A group first, so the subscription's alias is bound before any datagram lands. + producer + .write_frame(Timestamp::from_millis(0).unwrap(), &b"before"[..]) + .unwrap(); + let mut before = subscriber.recv_group().await.unwrap().unwrap(); + assert_eq!(before.sequence, 0, "{version}"); + // Drafts without a track timescale stamp objects on arrival, datagrams included. + let stamped = before.next_frame().await.unwrap().unwrap().timestamp.value() == 0; - assert!( - subscriber.recv_datagram().now_or_never().is_none(), - "MoQ Transport must not map datagrams" + producer + .insert_datagram( + 5, + Timestamp::from_millis(7).unwrap(), + bytes::Bytes::from_static(PAYLOAD), + ) + .unwrap(); + let datagram = subscriber.recv_datagram().await.unwrap().unwrap(); + assert_eq!(datagram.sequence, 5, "{version}: the relay must not renumber"); + assert_eq!(&datagram.payload[..], PAYLOAD, "{version}"); + if stamped { + let expected = Timestamp::from_millis(7).unwrap(); + assert_eq!( + datagram.timestamp.convert(expected.scale()).unwrap(), + expected, + "{version}" ); - }) - .await - .expect("timed out"); + } + + // The group sequence continues past the datagram. + producer + .write_frame(Timestamp::from_millis(8).unwrap(), &b"after"[..]) + .unwrap(); + let after = subscriber.recv_group().await.unwrap().unwrap(); + assert_eq!(after.sequence, 6, "{version}"); } /// Explicit insert keeps the origin sequence on the lite wire, including a gap. diff --git a/swift/README.md b/swift/README.md index 766f1ff96e..05c2c6d02e 100644 --- a/swift/README.md +++ b/swift/README.md @@ -72,7 +72,7 @@ Incoming server `Request` values expose the query-free `path` before acceptance. Raw tracks also expose best-effort datagrams: `TrackProducer.appendDatagram(_:timestampUs:)` returns the assigned sequence number, `TrackConsumer.recvDatagram()` receives one datagram, and `TrackConsumer.datagrams` streams them in arrival order. Payloads are capped at 1200 bytes. -Datagrams require a datagram-capable transport and lite-05 or newer moq-lite; IETF moq-transport, +Datagrams require a datagram-capable transport and lite-05 or newer moq-lite, or moq-transport; pre-lite-05, WebSocket, and TCP paths do not deliver them, and there is no stream fallback. JSON tracks carry your own `Codable` types with the framing handled for you. You opt into one of two From a77b18f445f5d10ec09cf5afbd1c5a7a2160c587 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 11:48:55 -0700 Subject: [PATCH 3/5] quest: plan datagram range for both protocols Co-Authored-By: Claude Opus 5.5 --- quest/m1/README.md | 1 + quest/m1/datagram-range.md | 29 +++++++++++++++++++++++++++++ 2 files changed, 30 insertions(+) create mode 100644 quest/m1/datagram-range.md diff --git a/quest/m1/README.md b/quest/m1/README.md index bf58cf9600..8677969dad 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -41,6 +41,7 @@ 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` - [Moxygen compatibility](/quest/m1/moxygen/README.md) - one subgroup per group, whole-group FETCH, and one datagram per group, never a full moxygen pass +- [Datagram range](/quest/m1/datagram-range.md) - a subscriber gets only the datagrams its subscription asked for, on both protocols, not the buffered backlog - [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/datagram-range.md b/quest/m1/datagram-range.md new file mode 100644 index 0000000000..a50a8d4246 --- /dev/null +++ b/quest/m1/datagram-range.md @@ -0,0 +1,29 @@ +# [S] Datagram range + +## Goal + +A subscriber receives only the datagrams its subscription asked for, on +moq-lite and on moq-transport, in Rust and JavaScript. A new subscriber +is not handed datagrams from before its start. + +## Plan + +Datagrams share the group sequence namespace, but their cursor ignores the +subscription's start and end. The model buffers the last 64 per track, and a +new subscriber's cursor starts at the oldest. So a late joiner, or a relay +fanning out a fresh downstream, gets stale datagrams first. Both protocols do +this today, and the moxygen line kept moq-transport matching moq-lite. + +Settle the rule once and apply it to both protocols. It could be the cursor +starting at the live edge, the subscribe range bounding datagrams the way it +bounds groups, or both. Prefer fixing it in the model over filtering in each +session. + +Watch the edge cases: a datagram at the start group when a frame offset +skips object 0, SUBSCRIBE_UPDATE moving the range, and a datagram that +lands before the subscription's alias or id is known. A test must tell a +filtered datagram apart from one dropped for any other reason. + +## Related + +- [Moxygen compatibility](/quest/m1/moxygen/README.md) - brought datagrams to moq-transport with moq-lite's behavior From bb934444d6bb29a7880e982e5b69a7430d7082e7 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 19:29:41 -0700 Subject: [PATCH 4/5] quest: follow up with moq-transport datagrams in @moq/net Rust now carries OBJECT_DATAGRAM; the JS session still does not. Co-Authored-By: Claude Opus 5.5 --- quest/m1/README.md | 1 + quest/m1/e2ee/README.md | 2 +- quest/m1/js-ietf-datagram.md | 21 +++++++++++++++++++++ 3 files changed, 23 insertions(+), 1 deletion(-) create mode 100644 quest/m1/js-ietf-datagram.md diff --git a/quest/m1/README.md b/quest/m1/README.md index 8677969dad..8f3a9fe8da 100644 --- a/quest/m1/README.md +++ b/quest/m1/README.md @@ -43,6 +43,7 @@ transport, benchmark tooling); worktrees isolate commits, not semantics. - [Moxygen compatibility](/quest/m1/moxygen/README.md) - one subgroup per group, whole-group FETCH, and one datagram per group, never a full moxygen pass - [Datagram range](/quest/m1/datagram-range.md) - a subscriber gets only the datagrams its subscription asked for, on both protocols, not the buffered backlog - [JavaScript FETCH](/quest/m1/js-fetch.md) - generic on-demand group serving and IETF FETCH for browser publishers +- [JavaScript moq-transport datagrams](/quest/m1/js-ietf-datagram.md) - `@moq/net` carries datagrams over moq-transport as `OBJECT_DATAGRAM`, like Rust - [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 - [Tooling](/quest/m1/tooling/README.md) - justfiles become a one-line menu over `sh/`, one impact map scopes CI, and every workflow step runs a recipe diff --git a/quest/m1/e2ee/README.md b/quest/m1/e2ee/README.md index 824379b5c5..5293ed17cd 100644 --- a/quest/m1/e2ee/README.md +++ b/quest/m1/e2ee/README.md @@ -40,7 +40,7 @@ The Rust and TypeScript cores expose the same surface, and nothing else: - Deterministic secret-derived physical names hide catalog, codec, role, quality, timeline, and custom-track semantics. Authorized clients derive the encrypted catalog track name, then learn the remaining opaque names from its decrypted contents. Every catalog representation is encrypted; Rust publishers must not emit a plaintext MSF catalog. - A platform that forwards and meters protected bytes must never preview, record, archive, transmux, transcode, transcribe, compose, or inspect them, rejecting those paths before opening a processing session or writing product state. Applications needing those operations terminate E2EE outside the platform. A platform classifies protected broadcasts by its own credential or product state, never by name; the moq.pro (downstream) exclusion classifier and dashboard work stay downstream. -- The first proof covers browser TypeScript and native Rust publication and playback in both directions, with grouped audio and video over both moq-lite and MoQ Transport. Shared vectors cover groups and moq-lite datagrams; MoQ Transport has no datagram delivery. +- The first proof covers browser TypeScript and native Rust publication and playback in both directions, with grouped audio and video over both moq-lite and MoQ Transport. Shared vectors cover groups and moq-lite datagrams; JavaScript has no MoQ Transport datagram delivery yet. ## Quests diff --git a/quest/m1/js-ietf-datagram.md b/quest/m1/js-ietf-datagram.md new file mode 100644 index 0000000000..8f4a0c590e --- /dev/null +++ b/quest/m1/js-ietf-datagram.md @@ -0,0 +1,21 @@ +# [M] JavaScript moq-transport datagrams + +## Goal + +`@moq/net` sends and receives datagrams over moq-transport as +`OBJECT_DATAGRAM`, matching Rust: one Object at ID 0 is a single-frame group +whose Group ID is the sequence. A JavaScript publisher's datagrams reach a Rust +subscriber, and a Rust publisher's reach a JavaScript one. + +## Plan + +Port `rs/moq-net/src/ietf/datagram.rs` and the session's send and receive +loops. Decode every draft's Type flags, drop what the model cannot carry the +same way Rust does, and close the session on a malformed Type. + +The integration test `ietf does not deliver datagrams` flips to delivery on +every supported draft, and `just test interop --all` covers both directions. + +## Related + +- [Datagram range](/quest/m1/datagram-range.md) - the subscribe range for datagrams, settled on both protocols From 8632aebf37b29e0ae4d3b8365b890f31b19f9437 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 20:48:50 -0700 Subject: [PATCH 5/5] docs: moq-transport now carries datagrams, including E2EE ones Map E2EE datagrams onto OBJECT_DATAGRAM in the draft, cover them with a transport test, and drop the stale moq-ffi and libmoq claims. Co-Authored-By: Claude Opus 5.5 --- dart/moq_ffi/lib/src/moq.dart | 2 +- drafts/draft-lcurley-moq-e2ee.md | 3 ++- rs/libmoq/src/api.rs | 2 +- rs/moq-e2ee/tests/transport.rs | 17 +++++++++++++---- rs/moq-ffi/src/consumer.rs | 2 +- 5 files changed, 18 insertions(+), 8 deletions(-) diff --git a/dart/moq_ffi/lib/src/moq.dart b/dart/moq_ffi/lib/src/moq.dart index 1459eadb5e..a044af3e3e 100644 --- a/dart/moq_ffi/lib/src/moq.dart +++ b/dart/moq_ffi/lib/src/moq.dart @@ -12693,7 +12693,7 @@ void _checkApiChecksums() { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqtrackconsumer_recv_datagram() != - 29049) { + 17412) { throw UniffiInternalError.panicked("UniFFI API checksum mismatch"); } if (uniffi_moq_ffi_checksum_method_moqtrackconsumer_recv_group() != 60887) { diff --git a/drafts/draft-lcurley-moq-e2ee.md b/drafts/draft-lcurley-moq-e2ee.md index ec68f21aa9..25ca49c8ea 100644 --- a/drafts/draft-lcurley-moq-e2ee.md +++ b/drafts/draft-lcurley-moq-e2ee.md @@ -227,7 +227,8 @@ An Object ID above `2^32-1` MUST be refused as `identity` before encryption or d ## Datagrams A datagram uses `domain = 0x01`, its 64-bit sequence as `group`, and `frame = 0`. -MoQ Transport has no datagram mapping in this profile; shared vectors cover grouped tracks on both transports and datagrams on moq-lite only. +On MoQ Transport, a datagram is an Object at ID 0 in an OBJECT_DATAGRAM whose Group ID is the sequence. +Shared vectors cover grouped tracks on both transports and datagrams on moq-lite only. ## Nonce The 96-bit AES-GCM nonce is: diff --git a/rs/libmoq/src/api.rs b/rs/libmoq/src/api.rs index d221f6f441..0b56c0a53f 100644 --- a/rs/libmoq/src/api.rs +++ b/rs/libmoq/src/api.rs @@ -3615,7 +3615,7 @@ pub extern "C" fn moq_consume_track_cancel(track: u32) -> i32 { /// touched again, so release `user_data` there. The terminal callback fires even after /// [moq_consume_datagrams_cancel]. Read each datagram with [moq_consume_datagram] and release /// it with [moq_consume_datagram_free]. Datagrams arrive only over datagram-capable -/// transports and lite-05 or newer moq-lite; there is no stream fallback. +/// transports on moq-transport or lite-05 and newer moq-lite; there is no stream fallback. /// /// Returns a non-zero handle to the subscription on success, or a negative code on failure. /// diff --git a/rs/moq-e2ee/tests/transport.rs b/rs/moq-e2ee/tests/transport.rs index bb114105e2..72dde42bc7 100644 --- a/rs/moq-e2ee/tests/transport.rs +++ b/rs/moq-e2ee/tests/transport.rs @@ -1,4 +1,4 @@ -//! Grouped tracks on moq-lite and MoQ Transport; datagrams on moq-lite. +//! Grouped tracks and datagrams on moq-lite and MoQ Transport. mod support; @@ -104,10 +104,9 @@ async fn grouped_over_ietf() { .expect("timed out"); } -#[tokio::test] -async fn datagrams_over_lite() { +async fn datagram_roundtrip(version: &str) { tokio::time::timeout(TEST_TIMEOUT, async { - let mut fixture = connect_protected("moq-lite-05".parse().unwrap(), "audio").await; + let mut fixture = connect_protected(version.parse().unwrap(), "audio").await; fixture .producer .append_datagram(Timestamp::from_millis(9).unwrap(), b"opus") @@ -123,3 +122,13 @@ async fn datagrams_over_lite() { .await .expect("timed out"); } + +#[tokio::test] +async fn datagrams_over_lite() { + datagram_roundtrip("moq-lite-05").await; +} + +#[tokio::test] +async fn datagrams_over_ietf() { + datagram_roundtrip("moq-transport-21").await; +} diff --git a/rs/moq-ffi/src/consumer.rs b/rs/moq-ffi/src/consumer.rs index 84da7e0076..14e441e320 100644 --- a/rs/moq-ffi/src/consumer.rs +++ b/rs/moq-ffi/src/consumer.rs @@ -633,7 +633,7 @@ impl MoqTrackConsumer { /// Receive the next best-effort datagram in arrival order. /// /// Returns `None` when the track ends. Datagram delivery is unavailable over - /// IETF moq-transport, pre-lite-05 moq-lite, and stream-only transports. + /// pre-lite-05 moq-lite and stream-only transports. /// Datagrams are a separate cursor from groups, so this works alongside either /// group order, never commits the track to one, and progresses while a group /// read is pending.