Skip to main content

moqtap_codec/
data_dispatch.rs

1//! Draft-neutral object framing for MoQT data streams.
2//!
3//! [`AnySubgroupObjectReader`],
4//! [`AnySubgroupObjectWriter`],
5//! [`AnyFetchObjectReader`] and
6//! [`AnyFetchObjectWriter`]
7//! present one API over fourteen drafts' object encodings. The values they
8//! produce — `AnySubgroupObject`, `AnySubgroupObjectMeta`, `AnyFetchObject`,
9//! `AnyFetchObjectMeta` — are plain structs of primitives, so a caller can
10//! address objects without naming a `draftNN` type. All of them are also
11//! re-exported from [`crate::dispatch`].
12//!
13//! Objects on drafts 07-13 are standalone: absolute Object IDs, and (on
14//! drafts 11-13) an extension block whose presence is fixed by the stream
15//! type. Drafts 14-20 delta-encode Object IDs against the previous object on
16//! the stream. Both are constructed from the stream's header and read one
17//! object at a time, so the difference stays inside this module.
18//!
19//! Fetch streams split the same way, at a different draft. Through draft-14 a
20//! fetch object spells out its Group ID, Subgroup ID, Object ID and Publisher
21//! Priority on every object, so each one stands alone. Drafts 15-20 put a
22//! Serialization Flags field first and let it leave any of those four off the
23//! wire, meaning "the prior object's" — and from draft-18 the two ID fields
24//! that remain are differences rather than values. So a fetch object on those
25//! drafts is only meaningful in stream order, and
26//! [`AnyFetchObjectReader`]
27//! carries the running state that resolves it. The
28//! values it produces are absolute on every draft.
29//!
30//! # Partial buffers
31//!
32//! Reader state after an error is unspecified. A caller that may be handed an
33//! incomplete object clones the reader, decodes against the clone, and
34//! overwrites the real reader only once the decode succeeds.
35
36use bytes::{Buf, BufMut};
37
38use crate::dispatch::{AnyFetchHeader, AnySubgroupHeader};
39use crate::error::CodecError;
40use crate::varint::VarInt;
41use crate::version::DraftVersion;
42
43// ── Draft-neutral object values ─────────────────────────────
44
45/// One object read from a subgroup data stream, normalised across drafts.
46///
47/// Field semantics are identical on every draft 07-20; the per-draft wire
48/// differences (absolute vs delta object IDs, count- vs length-prefixed
49/// extension blocks, typed vs raw status codes) are resolved by
50/// [`AnySubgroupObjectReader`] before this value is produced.
51#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct AnySubgroupObject {
53    /// Absolute Object ID. Already resolved from delta encoding on drafts
54    /// 14-20; copied verbatim on drafts 07-13.
55    pub object_id: u64,
56    /// The extension-header (draft-17+: "property") block's contents,
57    /// excluding any length or count prefix. Empty when the draft has no
58    /// extension block, when the enclosing header's extensions bit is clear,
59    /// or when the block is present but zero-length.
60    ///
61    /// Opaque and never re-parsed by this crate, and a verbatim copy of the
62    /// wire bytes on every draft except draft-08. Draft-08's block is
63    /// count-prefixed with no byte length, so its contents can only be
64    /// delimited by parsing each extension; the blob is therefore
65    /// re-serialized from that parsed form, which re-encodes every varint
66    /// minimally. A draft-08 extension whose value arrived as a legal
67    /// non-minimal varint is semantically but not byte-identically preserved.
68    pub extension_headers: Vec<u8>,
69    /// Number of extensions in `extension_headers`. `Some` only on draft-08,
70    /// whose extension block is count-prefixed rather than
71    /// byte-length-prefixed, so the count cannot be recovered from the blob
72    /// without re-parsing it. `None` on every other draft.
73    pub extension_count: Option<u64>,
74    /// Object Status as the raw wire code, present only when the payload
75    /// length is zero. `None` means a non-empty payload followed and the
76    /// status is implicitly Normal.
77    ///
78    /// Kept as a raw code rather than a typed enum because the assigned set
79    /// changes across drafts and this value crosses drafts — a relay reads a
80    /// status on one and writes it on another, where the same number may mean
81    /// something else or nothing at all. The code is nonetheless always one
82    /// the *source* draft assigns: every draft's decoder refuses an unassigned
83    /// status, so this field never carries a value its draft's Object Status
84    /// section forbids.
85    pub status: Option<u64>,
86    /// Object payload. Empty when `status` is `Some`.
87    pub payload: Vec<u8>,
88}
89
90/// The framing of one subgroup object, without its payload.
91///
92/// Produced by [`AnySubgroupObjectReader::read_object_meta`] for callers that
93/// forward an object's bytes verbatim and never inspect the payload.
94#[derive(Debug, Clone, Copy, PartialEq, Eq)]
95pub struct AnySubgroupObjectMeta {
96    /// Absolute Object ID, resolved as for [`AnySubgroupObject::object_id`].
97    pub object_id: u64,
98    /// Declared payload length in bytes. Zero when `status` is `Some`.
99    pub payload_length: u64,
100    /// Object Status wire code, as for [`AnySubgroupObject::status`].
101    pub status: Option<u64>,
102    /// Byte length of the extension/property block's contents, excluding its
103    /// prefix. Always the length of the blob
104    /// [`AnySubgroupObject::extension_headers`] would carry, which on draft-08
105    /// is a re-serialized copy rather than the wire bytes.
106    pub extension_headers_len: u64,
107    /// Total bytes this object occupies on the wire, prefix fields included.
108    /// Equals the number of bytes the reader consumed.
109    pub wire_len: u64,
110}
111
112/// What an End of Range indicator asserts about the Locations it covers.
113///
114/// Drafts 16-20 let a fetch stream state that a run of Objects was not
115/// serialized instead of sending them: one frame names the Location that ends
116/// the run, and every Location from the previously serialized Object up to and
117/// including that one is covered. The indicators are the same frame shape with
118/// different claims behind it, and a subscriber may cache the first as a
119/// settled gap while it must not cache the others, so they are carried apart
120/// rather than merged.
121///
122/// Never produced on drafts 07-15, which have no such frame.
123#[derive(Debug, Clone, Copy, PartialEq, Eq)]
124pub enum AnyFetchEndOfRange {
125    /// The covered Objects do not exist.
126    NonExistent,
127    /// The covered Objects' status is unknown to the publisher.
128    Unknown,
129    /// The covered Objects timed out: the relay abandoned them when its
130    /// `FILL_TIMEOUT` budget ran out.
131    ///
132    /// **Draft-20 and later only.** Drafts 16 through 19 have two indicators
133    /// and report the same Objects as [`Self::Unknown`], so this variant is
134    /// never produced on them — which makes its presence a fact about the
135    /// stream's draft as much as about the frame, and is why a receiver cannot
136    /// treat the two as interchangeable.
137    TimedOut,
138}
139
140/// The order a fetch response's groups arrive in.
141///
142/// Drafts 18 and 19 encode an Object's Group ID as a difference from the
143/// previous Object's, and the direction that difference moves is the fetch's
144/// Group Order — which is settled by the control exchange that opened the
145/// fetch and never appears on the data stream. A reader therefore has to be
146/// told, and telling it wrong does not fail to parse: every Object decodes
147/// under a Group ID walking the wrong way.
148///
149/// Ignored on drafts 07-17, whose fetch objects state their Group ID outright.
150#[derive(Debug, Clone, Copy, PartialEq, Eq)]
151pub enum AnyFetchGroupOrder {
152    /// Group IDs increase along the stream; a difference is added.
153    Ascending,
154    /// Group IDs decrease along the stream; a difference is subtracted.
155    Descending,
156}
157
158/// One frame read from a fetch data stream, normalised across drafts.
159///
160/// Usually an object. On drafts 16-20 it may instead be an End of Range
161/// indicator, which carries a Location and no content — [`Self::end_of_range`]
162/// is what tells the two apart, and it is `None` for every object.
163#[derive(Debug, Clone, PartialEq, Eq)]
164pub struct AnyFetchObject {
165    /// Absolute Group ID. Already resolved against the objects before it on
166    /// drafts 15-20, whose fetch objects may omit the field or (on drafts
167    /// 18-19) encode it as a difference; copied verbatim on drafts 07-14.
168    pub group_id: u64,
169    /// Absolute Subgroup ID, resolved as [`Self::group_id`] is. Zero and
170    /// meaningless when [`Self::has_subgroup_id`] is `false`.
171    pub subgroup_id: u64,
172    /// Whether this frame has a Subgroup ID at all.
173    ///
174    /// `true` on every draft 07-15, and for every End of Range indicator's
175    /// predecessor. `false` in two cases drafts 16-20 add: an object whose
176    /// Forwarding Preference is Datagram, which has no Subgroup ID anywhere in
177    /// its framing, and an End of Range indicator, whose Location is a Group
178    /// and Object ID only.
179    ///
180    /// Kept beside `subgroup_id` rather than folded into it because a relay
181    /// that writes this object onto another stream must not invent a Subgroup
182    /// ID of zero for an object that has none: on the receiving draft zero is a
183    /// real subgroup.
184    pub has_subgroup_id: bool,
185    /// Absolute Object ID, resolved as [`Self::group_id`] is.
186    pub object_id: u64,
187    /// Publisher Priority in force for this frame.
188    ///
189    /// Drafts 15-20 let an object omit the field and take the previous
190    /// object's, and an End of Range indicator never carries one. Where
191    /// nothing on the stream has stated a priority, this is 128 — the value
192    /// every draft 15-20 gives a subscription whose Default Publisher Priority
193    /// property is omitted (draft-19 Section 12.4).
194    pub publisher_priority: u8,
195    /// Extension/property block contents, excluding its prefix. Same
196    /// convention as [`AnySubgroupObject::extension_headers`].
197    pub extension_headers: Vec<u8>,
198    /// Number of extensions; `Some` only on draft-08. Same convention as
199    /// [`AnySubgroupObject::extension_count`].
200    pub extension_count: Option<u64>,
201    /// Object Status wire code, present only when the payload is empty.
202    ///
203    /// Always `None` on drafts 16-20: those drafts removed the field from
204    /// fetch objects entirely, stating that Object Status "is only present in
205    /// objects that are delivered via a SUBSCRIPTION, and is absent in Objects
206    /// delivered via a FETCH" (draft-19 Section 11.2.1.1). A zero-length fetch
207    /// object there is an object with no bytes, not a status object.
208    pub status: Option<u64>,
209    /// Which End of Range indicator this frame is, or `None` for an object.
210    ///
211    /// An indicator has a Location and a payload length and nothing else: its
212    /// `payload`, `extension_headers` and `status` are empty, its
213    /// [`Self::has_subgroup_id`] is `false`, and its `publisher_priority` is
214    /// whatever was in force rather than anything it stated.
215    pub end_of_range: Option<AnyFetchEndOfRange>,
216    /// Object payload. Empty when `status` is `Some`.
217    pub payload: Vec<u8>,
218}
219
220/// The framing of one fetch frame, without its payload.
221///
222/// Every field carries the meaning it does on [`AnyFetchObject`].
223#[derive(Debug, Clone, Copy, PartialEq, Eq)]
224pub struct AnyFetchObjectMeta {
225    /// Absolute Group ID.
226    pub group_id: u64,
227    /// Absolute Subgroup ID. Meaningless when `has_subgroup_id` is `false`.
228    pub subgroup_id: u64,
229    /// Whether this frame has a Subgroup ID at all; see
230    /// [`AnyFetchObject::has_subgroup_id`].
231    pub has_subgroup_id: bool,
232    /// Absolute Object ID.
233    pub object_id: u64,
234    /// Publisher Priority in force for this frame; see
235    /// [`AnyFetchObject::publisher_priority`].
236    pub publisher_priority: u8,
237    /// Declared payload length in bytes.
238    pub payload_length: u64,
239    /// Object Status wire code; always `None` on drafts 16-20.
240    pub status: Option<u64>,
241    /// Which End of Range indicator this frame is, or `None` for an object.
242    pub end_of_range: Option<AnyFetchEndOfRange>,
243    /// Byte length of the extension block's contents, excluding its prefix.
244    pub extension_headers_len: u64,
245    /// Total bytes this object occupies on the wire.
246    pub wire_len: u64,
247}
248
249/// The Publisher Priority a fetch frame that states none is read under.
250///
251/// Drafts 16-20 let an object leave the field off the wire and take the
252/// previous object's, and an End of Range indicator carries none at all, so a
253/// stream can reach a frame with no priority ever having been stated. Every one
254/// of those drafts fixes the same fallback for a subscription that never stated
255/// one — draft-19 Section 12.4: "If omitted, the Default Publisher Priority is
256/// 128" — and that is what is reported here.
257///
258/// Draft-15 needs no such fallback: its own reader refuses an object that
259/// inherits a priority with nothing to inherit from, and it has no End of Range
260/// frame, so every draft-15 fetch object has a priority the stream stated.
261///
262/// An End of Range indicator reports the Priority still in force from the last
263/// Object before it, and this constant only when no Object has preceded it.
264/// None of the four drafts decides that. Drafts 17, 18 and 19 say what the
265/// *next* Object inherits — draft-19 Section 11.4.4.2: "Prior Priority: The
266/// Priority from the last actual Object before the End of Range indicator" —
267/// draft-16 does not say even that, and all four agree only that a marker
268/// carries no Priority field. So the answer is chosen here, once, for the four
269/// of them: [`AnyFetchObject::publisher_priority`] is a `u8` with no way to
270/// report "none", and the value in force is one the stream did state, where
271/// this constant would be one it never did.
272#[cfg(any(
273    feature = "draft16",
274    feature = "draft17",
275    feature = "draft18",
276    feature = "draft19",
277    feature = "draft20"
278))]
279const DEFAULT_PUBLISHER_PRIORITY: u8 = 128;
280
281/// Conversions shared by the per-draft glue below. Unused when no draft
282/// feature is enabled.
283#[allow(dead_code)]
284mod conv {
285    use super::{AnySubgroupObject, Buf, CodecError};
286    use crate::varint::VarInt;
287
288    /// Advance `buf` past `len` bytes without copying them.
289    pub fn skip(buf: &mut impl Buf, len: u64) -> Result<(), CodecError> {
290        let len = usize::try_from(len).map_err(|_| CodecError::UnexpectedEnd)?;
291        if buf.remaining() < len {
292            return Err(CodecError::UnexpectedEnd);
293        }
294        buf.advance(len);
295        Ok(())
296    }
297
298    /// Copy `len` bytes out of `buf`.
299    pub fn take(buf: &mut impl Buf, len: u64) -> Result<Vec<u8>, CodecError> {
300        let len = usize::try_from(len).map_err(|_| CodecError::UnexpectedEnd)?;
301        crate::types::read_bytes(buf, len)
302    }
303
304    /// Wrap a value that must fit the varint range.
305    pub fn varint(v: u64) -> Result<VarInt, CodecError> {
306        VarInt::from_u64(v).map_err(|_| CodecError::InvalidField)
307    }
308
309    /// The status code to encode for `object`, or `None` when a payload
310    /// follows instead. An empty payload always carries a status on the
311    /// wire, so a missing one defaults to Normal.
312    ///
313    /// Status `0` with a payload is not a contradiction and is not refused.
314    /// Every draft from 07 to 20 assigns `0x0` to Normal, and every one of them
315    /// encodes a payload-bearing object by leaving the status field off — the
316    /// status is Normal precisely because bytes follow. Saying so explicitly
317    /// asks for the same frame as leaving it out, so both answer `None`, and
318    /// the byte written is identical either way.
319    ///
320    /// Any other status with a payload is refused: those are the statuses whose
321    /// wire form is a status code standing where the payload would be, so there
322    /// is no frame that carries both.
323    pub fn status_to_write(object: &AnySubgroupObject) -> Result<Option<u64>, CodecError> {
324        match (object.status, object.payload.is_empty()) {
325            (Some(0), false) => Ok(None),
326            (Some(_), false) => Err(CodecError::InvalidField),
327            (Some(code), true) => Ok(Some(code)),
328            (None, true) => Ok(Some(0)),
329            (None, false) => Ok(None),
330        }
331    }
332}
333
334// ── Per-draft glue ──────────────────────────────────────────
335
336/// Generates the conversion glue between one draft's standalone
337/// `ObjectHeader` and the draft-neutral object types.
338///
339/// The leading keyword selects the draft's extension-block shape: absent
340/// (draft-07), count-prefixed (draft-08), byte-length-prefixed (drafts
341/// 09/10), or byte-length-prefixed and gated on the stream type (drafts
342/// 11-13).
343macro_rules! legacy_subgroup_glue {
344    (no_extensions $name:ident, $feat:literal, $draft:ident) => {
345        #[cfg(feature = $feat)]
346        mod $name {
347            use super::conv;
348            use super::{AnySubgroupObject, AnySubgroupObjectMeta};
349            use crate::error::CodecError;
350            use crate::$draft::data_stream::ObjectHeader;
351            use crate::$draft::types::ObjectStatus;
352            use bytes::{Buf, BufMut};
353
354            pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
355                let header = ObjectHeader::decode(buf)?;
356                let payload_length = header.payload_length.into_inner();
357                let (status, payload) = if payload_length == 0 {
358                    (Some(header.object_status as u64), Vec::new())
359                } else {
360                    (None, conv::take(buf, payload_length)?)
361                };
362                Ok(AnySubgroupObject {
363                    object_id: header.object_id.into_inner(),
364                    extension_headers: Vec::new(),
365                    extension_count: None,
366                    status,
367                    payload,
368                })
369            }
370
371            pub fn read_object_meta(
372                buf: &mut impl Buf,
373            ) -> Result<AnySubgroupObjectMeta, CodecError> {
374                let start = buf.remaining();
375                let header = ObjectHeader::decode(buf)?;
376                let payload_length = header.payload_length.into_inner();
377                let status = if payload_length == 0 {
378                    Some(header.object_status as u64)
379                } else {
380                    conv::skip(buf, payload_length)?;
381                    None
382                };
383                Ok(AnySubgroupObjectMeta {
384                    object_id: header.object_id.into_inner(),
385                    payload_length,
386                    status,
387                    extension_headers_len: 0,
388                    wire_len: (start - buf.remaining()) as u64,
389                })
390            }
391
392            pub fn write_object(
393                object: &AnySubgroupObject,
394                buf: &mut impl BufMut,
395            ) -> Result<(), CodecError> {
396                if !object.extension_headers.is_empty() {
397                    return Err(CodecError::InvalidField);
398                }
399                let object_status = match conv::status_to_write(object)? {
400                    Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
401                    None => ObjectStatus::Normal,
402                };
403                ObjectHeader {
404                    object_id: conv::varint(object.object_id)?,
405                    payload_length: conv::varint(object.payload.len() as u64)?,
406                    object_status,
407                }
408                .encode(buf);
409                buf.put_slice(&object.payload);
410                Ok(())
411            }
412        }
413    };
414
415    (count_extensions $name:ident, $feat:literal, $draft:ident) => {
416        #[cfg(feature = $feat)]
417        mod $name {
418            use super::conv;
419            use super::{AnySubgroupObject, AnySubgroupObjectMeta};
420            use crate::error::CodecError;
421            use crate::$draft::data_stream::ObjectHeader;
422            use crate::$draft::types::ObjectStatus;
423            use bytes::{Buf, BufMut};
424
425            pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
426                let header = ObjectHeader::decode(buf)?;
427                let payload_length = header.payload_length.into_inner();
428                let (status, payload) = if payload_length == 0 {
429                    (Some(header.object_status as u64), Vec::new())
430                } else {
431                    (None, conv::take(buf, payload_length)?)
432                };
433                Ok(AnySubgroupObject {
434                    object_id: header.object_id.into_inner(),
435                    extension_headers: header.extensions,
436                    extension_count: Some(header.extension_count.into_inner()),
437                    status,
438                    payload,
439                })
440            }
441
442            pub fn read_object_meta(
443                buf: &mut impl Buf,
444            ) -> Result<AnySubgroupObjectMeta, CodecError> {
445                let start = buf.remaining();
446                let header = ObjectHeader::decode(buf)?;
447                let payload_length = header.payload_length.into_inner();
448                let status = if payload_length == 0 {
449                    Some(header.object_status as u64)
450                } else {
451                    conv::skip(buf, payload_length)?;
452                    None
453                };
454                Ok(AnySubgroupObjectMeta {
455                    object_id: header.object_id.into_inner(),
456                    payload_length,
457                    status,
458                    extension_headers_len: header.extensions.len() as u64,
459                    wire_len: (start - buf.remaining()) as u64,
460                })
461            }
462
463            pub fn write_object(
464                object: &AnySubgroupObject,
465                buf: &mut impl BufMut,
466            ) -> Result<(), CodecError> {
467                let extension_count = match object.extension_count {
468                    Some(count) => count,
469                    None if object.extension_headers.is_empty() => 0,
470                    None => return Err(CodecError::InvalidField),
471                };
472                let object_status = match conv::status_to_write(object)? {
473                    Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
474                    None => ObjectStatus::Normal,
475                };
476                ObjectHeader {
477                    object_id: conv::varint(object.object_id)?,
478                    extension_count: conv::varint(extension_count)?,
479                    extensions: object.extension_headers.clone(),
480                    payload_length: conv::varint(object.payload.len() as u64)?,
481                    object_status,
482                }
483                .encode(buf);
484                buf.put_slice(&object.payload);
485                Ok(())
486            }
487        }
488    };
489
490    (length_extensions $name:ident, $feat:literal, $draft:ident) => {
491        #[cfg(feature = $feat)]
492        mod $name {
493            use super::conv;
494            use super::{AnySubgroupObject, AnySubgroupObjectMeta};
495            use crate::error::CodecError;
496            use crate::$draft::data_stream::ObjectHeader;
497            use crate::$draft::types::ObjectStatus;
498            use bytes::{Buf, BufMut};
499
500            pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
501                let header = ObjectHeader::decode(buf)?;
502                let payload_length = header.payload_length.into_inner();
503                let (status, payload) = if payload_length == 0 {
504                    (Some(header.object_status as u64), Vec::new())
505                } else {
506                    (None, conv::take(buf, payload_length)?)
507                };
508                Ok(AnySubgroupObject {
509                    object_id: header.object_id.into_inner(),
510                    extension_headers: header.extensions,
511                    extension_count: None,
512                    status,
513                    payload,
514                })
515            }
516
517            pub fn read_object_meta(
518                buf: &mut impl Buf,
519            ) -> Result<AnySubgroupObjectMeta, CodecError> {
520                let start = buf.remaining();
521                let header = ObjectHeader::decode(buf)?;
522                let payload_length = header.payload_length.into_inner();
523                let status = if payload_length == 0 {
524                    Some(header.object_status as u64)
525                } else {
526                    conv::skip(buf, payload_length)?;
527                    None
528                };
529                Ok(AnySubgroupObjectMeta {
530                    object_id: header.object_id.into_inner(),
531                    payload_length,
532                    status,
533                    extension_headers_len: header.extension_headers_length.into_inner(),
534                    wire_len: (start - buf.remaining()) as u64,
535                })
536            }
537
538            pub fn write_object(
539                object: &AnySubgroupObject,
540                buf: &mut impl BufMut,
541            ) -> Result<(), CodecError> {
542                let object_status = match conv::status_to_write(object)? {
543                    Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
544                    None => ObjectStatus::Normal,
545                };
546                ObjectHeader {
547                    object_id: conv::varint(object.object_id)?,
548                    extension_headers_length: conv::varint(object.extension_headers.len() as u64)?,
549                    extensions: object.extension_headers.clone(),
550                    payload_length: conv::varint(object.payload.len() as u64)?,
551                    object_status,
552                }
553                .encode(buf);
554                buf.put_slice(&object.payload);
555                Ok(())
556            }
557        }
558    };
559
560    (gated_extensions $name:ident, $feat:literal, $draft:ident) => {
561        #[cfg(feature = $feat)]
562        mod $name {
563            use super::conv;
564            use super::{AnySubgroupObject, AnySubgroupObjectMeta};
565            use crate::error::CodecError;
566            use crate::$draft::data_stream::ObjectHeader;
567            use crate::$draft::types::ObjectStatus;
568            use bytes::{Buf, BufMut};
569
570            pub fn read_object(
571                extensions: bool,
572                buf: &mut impl Buf,
573            ) -> Result<AnySubgroupObject, CodecError> {
574                let header = ObjectHeader::decode_with_extensions(extensions, buf)?;
575                let payload_length = header.payload_length.into_inner();
576                let (status, payload) = if payload_length == 0 {
577                    (Some(header.object_status as u64), Vec::new())
578                } else {
579                    (None, conv::take(buf, payload_length)?)
580                };
581                Ok(AnySubgroupObject {
582                    object_id: header.object_id.into_inner(),
583                    extension_headers: header.extensions,
584                    extension_count: None,
585                    status,
586                    payload,
587                })
588            }
589
590            pub fn read_object_meta(
591                extensions: bool,
592                buf: &mut impl Buf,
593            ) -> Result<AnySubgroupObjectMeta, CodecError> {
594                let start = buf.remaining();
595                let header = ObjectHeader::decode_with_extensions(extensions, buf)?;
596                let payload_length = header.payload_length.into_inner();
597                let status = if payload_length == 0 {
598                    Some(header.object_status as u64)
599                } else {
600                    conv::skip(buf, payload_length)?;
601                    None
602                };
603                Ok(AnySubgroupObjectMeta {
604                    object_id: header.object_id.into_inner(),
605                    payload_length,
606                    status,
607                    extension_headers_len: header.extension_headers_length.into_inner(),
608                    wire_len: (start - buf.remaining()) as u64,
609                })
610            }
611
612            pub fn write_object(
613                extensions: bool,
614                object: &AnySubgroupObject,
615                buf: &mut impl BufMut,
616            ) -> Result<(), CodecError> {
617                if !extensions && !object.extension_headers.is_empty() {
618                    return Err(CodecError::InvalidField);
619                }
620                let object_status = match conv::status_to_write(object)? {
621                    Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
622                    None => ObjectStatus::Normal,
623                };
624                ObjectHeader {
625                    object_id: conv::varint(object.object_id)?,
626                    extension_headers_length: conv::varint(object.extension_headers.len() as u64)?,
627                    extensions: object.extension_headers.clone(),
628                    payload_length: conv::varint(object.payload.len() as u64)?,
629                    object_status,
630                }
631                .encode_with_extensions(extensions, buf);
632                buf.put_slice(&object.payload);
633                Ok(())
634            }
635        }
636    };
637}
638
639legacy_subgroup_glue!(no_extensions sg07, "draft07", draft07);
640legacy_subgroup_glue!(count_extensions sg08, "draft08", draft08);
641legacy_subgroup_glue!(length_extensions sg09, "draft09", draft09);
642legacy_subgroup_glue!(length_extensions sg10, "draft10", draft10);
643legacy_subgroup_glue!(gated_extensions sg11, "draft11", draft11);
644legacy_subgroup_glue!(gated_extensions sg12, "draft12", draft12);
645legacy_subgroup_glue!(gated_extensions sg13, "draft13", draft13);
646
647/// Generates the conversion glue between one draft's stateful
648/// `SubgroupObjectReader` and the draft-neutral object types.
649///
650/// Every draft from 14 on carries a typed `ObjectStatus`, so the leading
651/// keyword selects how the payload length reaches the wire instead: draft-14
652/// derives it from the payload, while drafts 15-20 carry an explicit
653/// payload-length field, which this glue always sets from the payload.
654///
655/// Both arms funnel the draft-neutral `AnySubgroupObject`, whose status is a
656/// bare `u64`, through the target draft's `ObjectStatus::from_u64`. That is
657/// the one place a status code the target draft does not assign can still be
658/// offered to an encoder at run time — relaying an object between drafts, for
659/// instance — and it is refused there with `CodecError::InvalidField`.
660macro_rules! modern_subgroup_glue {
661    (derived_length $name:ident, $feat:literal, $draft:ident) => {
662        #[cfg(feature = $feat)]
663        mod $name {
664            use super::conv;
665            use super::{AnySubgroupObject, AnySubgroupObjectMeta};
666            use crate::error::CodecError;
667            use crate::$draft::data_stream::{SubgroupObject, SubgroupObjectReader};
668            use crate::$draft::types::ObjectStatus;
669            use bytes::{Buf, BufMut};
670
671            pub fn read_object(
672                reader: &mut SubgroupObjectReader,
673                buf: &mut impl Buf,
674            ) -> Result<AnySubgroupObject, CodecError> {
675                let object = reader.read_object(buf)?;
676                Ok(AnySubgroupObject {
677                    object_id: object.object_id.into_inner(),
678                    extension_headers: object.extension_headers,
679                    extension_count: None,
680                    status: object.status.map(ObjectStatus::as_u64),
681                    payload: object.payload,
682                })
683            }
684
685            pub fn read_object_meta(
686                reader: &mut SubgroupObjectReader,
687                buf: &mut impl Buf,
688            ) -> Result<AnySubgroupObjectMeta, CodecError> {
689                let meta = reader.read_object_meta(buf)?;
690                Ok(AnySubgroupObjectMeta {
691                    object_id: meta.object_id,
692                    payload_length: meta.payload_length,
693                    status: meta.status,
694                    extension_headers_len: meta.extension_headers_len,
695                    wire_len: meta.wire_len,
696                })
697            }
698
699            pub fn write_object(
700                writer: &mut SubgroupObjectReader,
701                object: &AnySubgroupObject,
702                buf: &mut impl BufMut,
703            ) -> Result<(), CodecError> {
704                let status = match conv::status_to_write(object)? {
705                    Some(code) => {
706                        Some(ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?)
707                    }
708                    None => None,
709                };
710                writer.write_object(
711                    &SubgroupObject {
712                        object_id: conv::varint(object.object_id)?,
713                        extension_headers: object.extension_headers.clone(),
714                        status,
715                        payload: object.payload.clone(),
716                    },
717                    buf,
718                )
719            }
720        }
721    };
722
723    (explicit_length $name:ident, $feat:literal, $draft:ident) => {
724        #[cfg(feature = $feat)]
725        mod $name {
726            use super::conv;
727            use super::{AnySubgroupObject, AnySubgroupObjectMeta};
728            use crate::error::CodecError;
729            use crate::$draft::data_stream::{SubgroupObject, SubgroupObjectReader};
730            use crate::$draft::types::ObjectStatus;
731            use bytes::{Buf, BufMut};
732
733            pub fn read_object(
734                reader: &mut SubgroupObjectReader,
735                buf: &mut impl Buf,
736            ) -> Result<AnySubgroupObject, CodecError> {
737                let object = reader.read_object(buf)?;
738                Ok(AnySubgroupObject {
739                    object_id: object.object_id.into_inner(),
740                    extension_headers: object.extension_headers,
741                    extension_count: None,
742                    status: object.object_status.map(ObjectStatus::as_u64),
743                    payload: object.payload,
744                })
745            }
746
747            pub fn read_object_meta(
748                reader: &mut SubgroupObjectReader,
749                buf: &mut impl Buf,
750            ) -> Result<AnySubgroupObjectMeta, CodecError> {
751                let meta = reader.read_object_meta(buf)?;
752                Ok(AnySubgroupObjectMeta {
753                    object_id: meta.object_id,
754                    payload_length: meta.payload_length,
755                    status: meta.status,
756                    extension_headers_len: meta.extension_headers_len,
757                    wire_len: meta.wire_len,
758                })
759            }
760
761            pub fn write_object(
762                writer: &mut SubgroupObjectReader,
763                object: &AnySubgroupObject,
764                buf: &mut impl BufMut,
765            ) -> Result<(), CodecError> {
766                let object_status = match conv::status_to_write(object)? {
767                    Some(code) => {
768                        Some(ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?)
769                    }
770                    None => None,
771                };
772                writer.write_object(
773                    &SubgroupObject {
774                        object_id: conv::varint(object.object_id)?,
775                        extension_headers: object.extension_headers.clone(),
776                        payload_length: conv::varint(object.payload.len() as u64)?,
777                        object_status,
778                        payload: object.payload.clone(),
779                    },
780                    buf,
781                )
782            }
783        }
784    };
785}
786
787modern_subgroup_glue!(derived_length sg14, "draft14", draft14);
788modern_subgroup_glue!(explicit_length sg15, "draft15", draft15);
789modern_subgroup_glue!(explicit_length sg16, "draft16", draft16);
790modern_subgroup_glue!(explicit_length sg17, "draft17", draft17);
791modern_subgroup_glue!(explicit_length sg18, "draft18", draft18);
792modern_subgroup_glue!(explicit_length sg19, "draft19", draft19);
793modern_subgroup_glue!(explicit_length sg20, "draft20", draft20);
794
795/// Generates the conversion glue for one draft's fetch objects.
796///
797/// The leading keyword selects the extension-block shape, as for
798/// `legacy_subgroup_glue!`. Unlike subgroup objects, the fetch extension
799/// block is unconditional on drafts 09-13 — it is never gated on the stream
800/// type.
801macro_rules! fetch_glue {
802    (no_extensions $name:ident, $feat:literal, $draft:ident) => {
803        #[cfg(feature = $feat)]
804        mod $name {
805            use super::conv;
806            use super::{AnyFetchObject, AnyFetchObjectMeta};
807            use crate::error::CodecError;
808            use crate::$draft::data_stream::FetchObjectHeader;
809            use bytes::Buf;
810
811            pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
812                let header = FetchObjectHeader::decode(buf)?;
813                let payload_length = header.payload_length.into_inner();
814                let (status, payload) = if payload_length == 0 {
815                    (Some(header.object_status as u64), Vec::new())
816                } else {
817                    (None, conv::take(buf, payload_length)?)
818                };
819                Ok(AnyFetchObject {
820                    group_id: header.group_id.into_inner(),
821                    subgroup_id: header.subgroup_id.into_inner(),
822                    has_subgroup_id: true,
823                    object_id: header.object_id.into_inner(),
824                    publisher_priority: header.publisher_priority,
825                    extension_headers: Vec::new(),
826                    extension_count: None,
827                    status,
828                    end_of_range: None,
829                    payload,
830                })
831            }
832
833            pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
834                let start = buf.remaining();
835                let header = FetchObjectHeader::decode(buf)?;
836                let payload_length = header.payload_length.into_inner();
837                let status = if payload_length == 0 {
838                    Some(header.object_status as u64)
839                } else {
840                    conv::skip(buf, payload_length)?;
841                    None
842                };
843                Ok(AnyFetchObjectMeta {
844                    group_id: header.group_id.into_inner(),
845                    subgroup_id: header.subgroup_id.into_inner(),
846                    has_subgroup_id: true,
847                    object_id: header.object_id.into_inner(),
848                    publisher_priority: header.publisher_priority,
849                    payload_length,
850                    status,
851                    end_of_range: None,
852                    extension_headers_len: 0,
853                    wire_len: (start - buf.remaining()) as u64,
854                })
855            }
856        }
857    };
858
859    (count_extensions $name:ident, $feat:literal, $draft:ident) => {
860        #[cfg(feature = $feat)]
861        mod $name {
862            use super::conv;
863            use super::{AnyFetchObject, AnyFetchObjectMeta};
864            use crate::error::CodecError;
865            use crate::$draft::data_stream::FetchObjectHeader;
866            use bytes::Buf;
867
868            pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
869                let header = FetchObjectHeader::decode(buf)?;
870                let payload_length = header.payload_length.into_inner();
871                let (status, payload) = if payload_length == 0 {
872                    (Some(header.object_status as u64), Vec::new())
873                } else {
874                    (None, conv::take(buf, payload_length)?)
875                };
876                Ok(AnyFetchObject {
877                    group_id: header.group_id.into_inner(),
878                    subgroup_id: header.subgroup_id.into_inner(),
879                    has_subgroup_id: true,
880                    object_id: header.object_id.into_inner(),
881                    publisher_priority: header.publisher_priority,
882                    extension_headers: header.extensions,
883                    extension_count: Some(header.extension_count.into_inner()),
884                    status,
885                    end_of_range: None,
886                    payload,
887                })
888            }
889
890            pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
891                let start = buf.remaining();
892                let header = FetchObjectHeader::decode(buf)?;
893                let payload_length = header.payload_length.into_inner();
894                let status = if payload_length == 0 {
895                    Some(header.object_status as u64)
896                } else {
897                    conv::skip(buf, payload_length)?;
898                    None
899                };
900                Ok(AnyFetchObjectMeta {
901                    group_id: header.group_id.into_inner(),
902                    subgroup_id: header.subgroup_id.into_inner(),
903                    has_subgroup_id: true,
904                    object_id: header.object_id.into_inner(),
905                    publisher_priority: header.publisher_priority,
906                    payload_length,
907                    status,
908                    end_of_range: None,
909                    extension_headers_len: header.extensions.len() as u64,
910                    wire_len: (start - buf.remaining()) as u64,
911                })
912            }
913        }
914    };
915
916    (length_extensions $name:ident, $feat:literal, $draft:ident) => {
917        #[cfg(feature = $feat)]
918        mod $name {
919            use super::conv;
920            use super::{AnyFetchObject, AnyFetchObjectMeta};
921            use crate::error::CodecError;
922            use crate::$draft::data_stream::FetchObjectHeader;
923            use bytes::Buf;
924
925            pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
926                let header = FetchObjectHeader::decode(buf)?;
927                let payload_length = header.payload_length.into_inner();
928                let (status, payload) = if payload_length == 0 {
929                    (Some(header.object_status as u64), Vec::new())
930                } else {
931                    (None, conv::take(buf, payload_length)?)
932                };
933                Ok(AnyFetchObject {
934                    group_id: header.group_id.into_inner(),
935                    subgroup_id: header.subgroup_id.into_inner(),
936                    has_subgroup_id: true,
937                    object_id: header.object_id.into_inner(),
938                    publisher_priority: header.publisher_priority,
939                    extension_headers: header.extensions,
940                    extension_count: None,
941                    status,
942                    end_of_range: None,
943                    payload,
944                })
945            }
946
947            pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
948                let start = buf.remaining();
949                let header = FetchObjectHeader::decode(buf)?;
950                let payload_length = header.payload_length.into_inner();
951                let status = if payload_length == 0 {
952                    Some(header.object_status as u64)
953                } else {
954                    conv::skip(buf, payload_length)?;
955                    None
956                };
957                Ok(AnyFetchObjectMeta {
958                    group_id: header.group_id.into_inner(),
959                    subgroup_id: header.subgroup_id.into_inner(),
960                    has_subgroup_id: true,
961                    object_id: header.object_id.into_inner(),
962                    publisher_priority: header.publisher_priority,
963                    payload_length,
964                    status,
965                    end_of_range: None,
966                    extension_headers_len: header.extension_headers_length.into_inner(),
967                    wire_len: (start - buf.remaining()) as u64,
968                })
969            }
970        }
971    };
972}
973
974fetch_glue!(no_extensions fo07, "draft07", draft07);
975fetch_glue!(count_extensions fo08, "draft08", draft08);
976fetch_glue!(length_extensions fo09, "draft09", draft09);
977fetch_glue!(length_extensions fo10, "draft10", draft10);
978fetch_glue!(length_extensions fo11, "draft11", draft11);
979fetch_glue!(length_extensions fo12, "draft12", draft12);
980fetch_glue!(length_extensions fo13, "draft13", draft13);
981
982#[cfg(feature = "draft14")]
983mod fo14 {
984    use super::{AnyFetchObject, AnyFetchObjectMeta};
985    use crate::draft14::data_stream::FetchObject;
986    use crate::draft14::types::ObjectStatus;
987    use crate::error::CodecError;
988    use bytes::Buf;
989
990    pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
991        let object = FetchObject::decode(buf)?;
992        Ok(AnyFetchObject {
993            group_id: object.group_id.into_inner(),
994            subgroup_id: object.subgroup_id.into_inner(),
995            has_subgroup_id: true,
996            object_id: object.object_id.into_inner(),
997            publisher_priority: object.publisher_priority,
998            extension_headers: object.extension_headers,
999            extension_count: None,
1000            status: object.status.map(ObjectStatus::as_u64),
1001            end_of_range: None,
1002            payload: object.payload,
1003        })
1004    }
1005
1006    pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
1007        let meta = FetchObject::decode_meta(buf)?;
1008        Ok(AnyFetchObjectMeta {
1009            group_id: meta.group_id,
1010            subgroup_id: meta.subgroup_id,
1011            has_subgroup_id: true,
1012            object_id: meta.object_id,
1013            publisher_priority: meta.publisher_priority,
1014            payload_length: meta.payload_length,
1015            status: meta.status,
1016            end_of_range: None,
1017            extension_headers_len: meta.extension_headers_len,
1018            wire_len: meta.wire_len,
1019        })
1020    }
1021}
1022
1023#[cfg(feature = "draft15")]
1024mod fo15 {
1025    use super::{conv, AnyFetchObject, AnyFetchObjectMeta};
1026    use crate::draft15::data_stream::FetchObjectReader;
1027    use crate::draft15::types::ObjectStatus;
1028    use crate::error::CodecError;
1029    use bytes::Buf;
1030
1031    pub fn read_object(
1032        reader: &mut FetchObjectReader,
1033        buf: &mut impl Buf,
1034    ) -> Result<AnyFetchObject, CodecError> {
1035        let header = reader.read_object_header(buf)?;
1036        // The draft-15 reader stops at the payload length so a caller can
1037        // forward the bytes without copying them; the draft-neutral object
1038        // holds the payload, so the copy happens here instead.
1039        let payload = conv::take(buf, header.payload_length.into_inner())?;
1040        Ok(AnyFetchObject {
1041            group_id: header.group_id.into_inner(),
1042            subgroup_id: header.subgroup_id.into_inner(),
1043            has_subgroup_id: true,
1044            object_id: header.object_id.into_inner(),
1045            publisher_priority: header.publisher_priority,
1046            extension_headers: header.extension_headers,
1047            extension_count: None,
1048            // Already `Some` only for a zero-length object, which is the
1049            // draft-neutral convention too.
1050            status: header.object_status.map(ObjectStatus::as_u64),
1051            end_of_range: None,
1052            payload,
1053        })
1054    }
1055
1056    pub fn read_object_frame(
1057        reader: &mut FetchObjectReader,
1058        buf: &mut impl Buf,
1059    ) -> Result<super::AnyFetchFrame, CodecError> {
1060        let start = buf.remaining();
1061        let header = reader.read_object_header(buf)?;
1062        let payload_length = header.payload_length.into_inner();
1063        conv::skip(buf, payload_length)?;
1064        let meta = AnyFetchObjectMeta {
1065            group_id: header.group_id.into_inner(),
1066            subgroup_id: header.subgroup_id.into_inner(),
1067            has_subgroup_id: true,
1068            object_id: header.object_id.into_inner(),
1069            publisher_priority: header.publisher_priority,
1070            payload_length,
1071            status: header.object_status.map(ObjectStatus::as_u64),
1072            end_of_range: None,
1073            extension_headers_len: header.extension_headers.len() as u64,
1074            wire_len: (start - buf.remaining()) as u64,
1075        };
1076        Ok(super::AnyFetchFrame {
1077            meta,
1078            draft: crate::version::DraftVersion::Draft15,
1079            shape: super::FetchFrameShape::Draft15(header),
1080        })
1081    }
1082
1083    pub fn read_object_meta(
1084        reader: &mut FetchObjectReader,
1085        buf: &mut impl Buf,
1086    ) -> Result<AnyFetchObjectMeta, CodecError> {
1087        read_object_frame(reader, buf).map(|frame| frame.meta)
1088    }
1089}
1090
1091#[cfg(feature = "draft16")]
1092mod fo16 {
1093    use super::DEFAULT_PUBLISHER_PRIORITY;
1094    use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1095    use crate::draft16::data_stream::{
1096        FetchEndOfRange, FetchObjectHeader, FetchObjectLocation, FetchObjectReader,
1097    };
1098    use crate::error::CodecError;
1099    use bytes::Buf;
1100
1101    /// Draft-16 keeps the frame and its resolved Location apart, and both are
1102    /// carried out of here: the Location fills the draft-neutral value, and the
1103    /// pair of them is what re-encoding this frame later takes.
1104    fn parts(
1105        reader: &mut FetchObjectReader,
1106        buf: &mut impl Buf,
1107    ) -> Result<(FetchObjectHeader, FetchObjectLocation), CodecError> {
1108        let header = FetchObjectHeader::decode(buf)?;
1109        let location = reader.resolve(&header)?;
1110        Ok((header, location))
1111    }
1112
1113    fn resolved(location: &FetchObjectLocation) -> super::Resolved {
1114        let end_of_range = location.end_of_range;
1115        super::Resolved {
1116            group_id: location.group_id,
1117            // Draft-16 Section 10.4.4.2 gives an End of Range indicator a Group
1118            // ID and an Object ID and says "Subgroup ID, Priority and Extensions
1119            // are not present". Its per-draft resolver still reads the two low
1120            // flag bits of a marker as Subgroup ID mode zero and answers zero,
1121            // which is a real Subgroup ID; the draft-neutral value says the
1122            // marker has none, as drafts 17-20 do.
1123            subgroup_id: location.subgroup_id.filter(|_| end_of_range.is_none()),
1124            object_id: location.object_id,
1125            publisher_priority: location.publisher_priority.unwrap_or(DEFAULT_PUBLISHER_PRIORITY),
1126            end_of_range: end_of_range.map(|r| match r {
1127                FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1128                FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1129            }),
1130        }
1131    }
1132
1133    pub fn read_object(
1134        reader: &mut FetchObjectReader,
1135        buf: &mut impl Buf,
1136    ) -> Result<AnyFetchObject, CodecError> {
1137        let (header, location) = parts(reader, buf)?;
1138        let payload = conv::take(buf, header.payload_length.into_inner())?;
1139        Ok(resolved(&location).into_object(header.extensions.unwrap_or_default(), payload))
1140    }
1141
1142    pub fn read_object_frame(
1143        reader: &mut FetchObjectReader,
1144        buf: &mut impl Buf,
1145    ) -> Result<super::AnyFetchFrame, CodecError> {
1146        let start = buf.remaining();
1147        let (header, location) = parts(reader, buf)?;
1148        let payload_length = header.payload_length.into_inner();
1149        conv::skip(buf, payload_length)?;
1150        let meta = resolved(&location).into_meta(
1151            header.extensions.as_ref().map_or(0, |e| e.len() as u64),
1152            payload_length,
1153            (start - buf.remaining()) as u64,
1154        );
1155        Ok(super::AnyFetchFrame {
1156            meta,
1157            draft: crate::version::DraftVersion::Draft16,
1158            shape: super::FetchFrameShape::Draft16(header, location),
1159        })
1160    }
1161
1162    pub fn read_object_meta(
1163        reader: &mut FetchObjectReader,
1164        buf: &mut impl Buf,
1165    ) -> Result<AnyFetchObjectMeta, CodecError> {
1166        read_object_frame(reader, buf).map(|frame| frame.meta)
1167    }
1168}
1169
1170#[cfg(feature = "draft17")]
1171mod fo17 {
1172    use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1173    use crate::draft17::data_stream::{EndOfRange, FetchObject, FetchObjectReader};
1174    use crate::error::CodecError;
1175    use bytes::Buf;
1176
1177    fn resolved(object: &FetchObject) -> super::Resolved {
1178        super::Resolved {
1179            group_id: object.group_id,
1180            subgroup_id: object.subgroup_id,
1181            object_id: object.object_id,
1182            publisher_priority: object
1183                .publisher_priority
1184                .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1185            end_of_range: object.header.end_of_range().map(|r| match r {
1186                EndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1187                EndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1188            }),
1189        }
1190    }
1191
1192    pub fn read_object(
1193        reader: &mut FetchObjectReader,
1194        buf: &mut impl Buf,
1195    ) -> Result<AnyFetchObject, CodecError> {
1196        let object = reader.read_object_header(buf)?;
1197        let resolved = resolved(&object);
1198        let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1199        Ok(resolved.into_object(object.header.properties, payload))
1200    }
1201
1202    pub fn read_object_frame(
1203        reader: &mut FetchObjectReader,
1204        buf: &mut impl Buf,
1205    ) -> Result<super::AnyFetchFrame, CodecError> {
1206        let start = buf.remaining();
1207        let object = reader.read_object_header(buf)?;
1208        let resolved = resolved(&object);
1209        let payload_length = object.header.payload_length.into_inner();
1210        conv::skip(buf, payload_length)?;
1211        let meta = resolved.into_meta(
1212            object.header.properties.len() as u64,
1213            payload_length,
1214            (start - buf.remaining()) as u64,
1215        );
1216        Ok(super::AnyFetchFrame {
1217            meta,
1218            draft: crate::version::DraftVersion::Draft17,
1219            shape: super::FetchFrameShape::Draft17(object),
1220        })
1221    }
1222
1223    pub fn read_object_meta(
1224        reader: &mut FetchObjectReader,
1225        buf: &mut impl Buf,
1226    ) -> Result<AnyFetchObjectMeta, CodecError> {
1227        read_object_frame(reader, buf).map(|frame| frame.meta)
1228    }
1229}
1230
1231#[cfg(feature = "draft18")]
1232mod fo18 {
1233    use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1234    use crate::draft18::data_stream::{EndOfRange, FetchObject, FetchObjectReader};
1235    use crate::error::CodecError;
1236    use bytes::Buf;
1237
1238    fn resolved(object: &FetchObject) -> super::Resolved {
1239        super::Resolved {
1240            group_id: object.group_id,
1241            subgroup_id: object.subgroup_id,
1242            object_id: object.object_id,
1243            publisher_priority: object
1244                .publisher_priority
1245                .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1246            end_of_range: object.header.end_of_range().map(|r| match r {
1247                EndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1248                EndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1249            }),
1250        }
1251    }
1252
1253    pub fn read_object(
1254        reader: &mut FetchObjectReader,
1255        buf: &mut impl Buf,
1256    ) -> Result<AnyFetchObject, CodecError> {
1257        let object = reader.read_object_header(buf)?;
1258        let resolved = resolved(&object);
1259        let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1260        Ok(resolved.into_object(object.header.properties, payload))
1261    }
1262
1263    pub fn read_object_frame(
1264        reader: &mut FetchObjectReader,
1265        buf: &mut impl Buf,
1266    ) -> Result<super::AnyFetchFrame, CodecError> {
1267        let start = buf.remaining();
1268        let object = reader.read_object_header(buf)?;
1269        let resolved = resolved(&object);
1270        let payload_length = object.header.payload_length.into_inner();
1271        conv::skip(buf, payload_length)?;
1272        let meta = resolved.into_meta(
1273            object.header.properties.len() as u64,
1274            payload_length,
1275            (start - buf.remaining()) as u64,
1276        );
1277        Ok(super::AnyFetchFrame {
1278            meta,
1279            draft: crate::version::DraftVersion::Draft18,
1280            shape: super::FetchFrameShape::Draft18(object),
1281        })
1282    }
1283
1284    pub fn read_object_meta(
1285        reader: &mut FetchObjectReader,
1286        buf: &mut impl Buf,
1287    ) -> Result<AnyFetchObjectMeta, CodecError> {
1288        read_object_frame(reader, buf).map(|frame| frame.meta)
1289    }
1290}
1291
1292#[cfg(feature = "draft19")]
1293mod fo19 {
1294    use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1295    use crate::draft19::data_stream::{FetchEndOfRange, FetchObject, FetchObjectReader};
1296    use crate::error::CodecError;
1297    use bytes::Buf;
1298
1299    fn resolved(object: &FetchObject) -> super::Resolved {
1300        super::Resolved {
1301            group_id: object.group_id,
1302            subgroup_id: object.subgroup_id,
1303            object_id: object.object_id,
1304            publisher_priority: object
1305                .publisher_priority
1306                .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1307            end_of_range: object.header.end_of_range().map(|r| match r {
1308                FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1309                FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1310            }),
1311        }
1312    }
1313
1314    pub fn read_object(
1315        reader: &mut FetchObjectReader,
1316        buf: &mut impl Buf,
1317    ) -> Result<AnyFetchObject, CodecError> {
1318        let object = reader.read_object_header(buf)?;
1319        let resolved = resolved(&object);
1320        let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1321        Ok(resolved.into_object(object.header.properties.unwrap_or_default(), payload))
1322    }
1323
1324    pub fn read_object_frame(
1325        reader: &mut FetchObjectReader,
1326        buf: &mut impl Buf,
1327    ) -> Result<super::AnyFetchFrame, CodecError> {
1328        let start = buf.remaining();
1329        let object = reader.read_object_header(buf)?;
1330        let resolved = resolved(&object);
1331        let payload_length = object.header.payload_length.into_inner();
1332        conv::skip(buf, payload_length)?;
1333        let meta = resolved.into_meta(
1334            object.header.properties.as_ref().map_or(0, |p| p.len() as u64),
1335            payload_length,
1336            (start - buf.remaining()) as u64,
1337        );
1338        Ok(super::AnyFetchFrame {
1339            meta,
1340            draft: crate::version::DraftVersion::Draft19,
1341            shape: super::FetchFrameShape::Draft19(object),
1342        })
1343    }
1344
1345    pub fn read_object_meta(
1346        reader: &mut FetchObjectReader,
1347        buf: &mut impl Buf,
1348    ) -> Result<AnyFetchObjectMeta, CodecError> {
1349        read_object_frame(reader, buf).map(|frame| frame.meta)
1350    }
1351}
1352
1353#[cfg(feature = "draft20")]
1354mod fo20 {
1355    use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1356    use crate::draft20::data_stream::{FetchEndOfRange, FetchObject, FetchObjectReader};
1357    use crate::error::CodecError;
1358    use bytes::Buf;
1359
1360    fn resolved(object: &FetchObject) -> super::Resolved {
1361        super::Resolved {
1362            group_id: object.group_id,
1363            subgroup_id: object.subgroup_id,
1364            object_id: object.object_id,
1365            publisher_priority: object
1366                .publisher_priority
1367                .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1368            end_of_range: object.header.end_of_range().map(|r| match r {
1369                FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1370                FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1371                // Draft-20's third marker, Table 7's 0x20C. Only this draft's
1372                // arm can produce it.
1373                FetchEndOfRange::TimedOut => AnyFetchEndOfRange::TimedOut,
1374            }),
1375        }
1376    }
1377
1378    pub fn read_object(
1379        reader: &mut FetchObjectReader,
1380        buf: &mut impl Buf,
1381    ) -> Result<AnyFetchObject, CodecError> {
1382        let object = reader.read_object_header(buf)?;
1383        let resolved = resolved(&object);
1384        let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1385        Ok(resolved.into_object(object.header.properties.unwrap_or_default(), payload))
1386    }
1387
1388    pub fn read_object_frame(
1389        reader: &mut FetchObjectReader,
1390        buf: &mut impl Buf,
1391    ) -> Result<super::AnyFetchFrame, CodecError> {
1392        let start = buf.remaining();
1393        let object = reader.read_object_header(buf)?;
1394        let resolved = resolved(&object);
1395        let payload_length = object.header.payload_length.into_inner();
1396        conv::skip(buf, payload_length)?;
1397        let meta = resolved.into_meta(
1398            object.header.properties.as_ref().map_or(0, |p| p.len() as u64),
1399            payload_length,
1400            (start - buf.remaining()) as u64,
1401        );
1402        Ok(super::AnyFetchFrame {
1403            meta,
1404            draft: crate::version::DraftVersion::Draft20,
1405            shape: super::FetchFrameShape::Draft20(object),
1406        })
1407    }
1408
1409    pub fn read_object_meta(
1410        reader: &mut FetchObjectReader,
1411        buf: &mut impl Buf,
1412    ) -> Result<AnyFetchObjectMeta, CodecError> {
1413        read_object_frame(reader, buf).map(|frame| frame.meta)
1414    }
1415}
1416
1417/// A drafts-16-to-19 fetch frame's identity once the fields its Serialization
1418/// Flags left off the wire have been filled in.
1419///
1420/// The four drafts read those flags differently enough to need a resolver
1421/// each, but they all end up saying the same five things, and turning that into
1422/// a draft-neutral value is the same work every time. `subgroup_id` is `None`
1423/// for the frames that have none at all rather than zero, which is a real
1424/// Subgroup ID; see [`AnyFetchObject::has_subgroup_id`].
1425#[cfg(any(
1426    feature = "draft16",
1427    feature = "draft17",
1428    feature = "draft18",
1429    feature = "draft19",
1430    feature = "draft20"
1431))]
1432struct Resolved {
1433    group_id: u64,
1434    subgroup_id: Option<u64>,
1435    object_id: u64,
1436    publisher_priority: u8,
1437    end_of_range: Option<AnyFetchEndOfRange>,
1438}
1439
1440#[cfg(any(
1441    feature = "draft16",
1442    feature = "draft17",
1443    feature = "draft18",
1444    feature = "draft19",
1445    feature = "draft20"
1446))]
1447impl Resolved {
1448    fn into_object(self, extension_headers: Vec<u8>, payload: Vec<u8>) -> AnyFetchObject {
1449        AnyFetchObject {
1450            group_id: self.group_id,
1451            subgroup_id: self.subgroup_id.unwrap_or(0),
1452            has_subgroup_id: self.subgroup_id.is_some(),
1453            object_id: self.object_id,
1454            publisher_priority: self.publisher_priority,
1455            extension_headers,
1456            extension_count: None,
1457            // Drafts 16-20 removed the Object Status field from fetch objects;
1458            // a zero-length payload here carries no code to report.
1459            status: None,
1460            end_of_range: self.end_of_range,
1461            payload,
1462        }
1463    }
1464
1465    fn into_meta(
1466        self,
1467        extension_headers_len: u64,
1468        payload_length: u64,
1469        wire_len: u64,
1470    ) -> AnyFetchObjectMeta {
1471        AnyFetchObjectMeta {
1472            group_id: self.group_id,
1473            subgroup_id: self.subgroup_id.unwrap_or(0),
1474            has_subgroup_id: self.subgroup_id.is_some(),
1475            object_id: self.object_id,
1476            publisher_priority: self.publisher_priority,
1477            payload_length,
1478            status: None,
1479            end_of_range: self.end_of_range,
1480            extension_headers_len,
1481            wire_len,
1482        }
1483    }
1484}
1485
1486// ── Subgroup object reader ──────────────────────────────────
1487
1488/// Per-draft reader state. Drafts 07-10 need none, drafts 11-13 need the
1489/// stream type's extensions-present flag, drafts 14-20 own a stateful
1490/// per-draft reader that tracks the Object ID delta.
1491#[derive(Debug, Clone)]
1492enum SubgroupReaderState {
1493    #[cfg(feature = "draft07")]
1494    Draft07,
1495    #[cfg(feature = "draft08")]
1496    Draft08,
1497    #[cfg(feature = "draft09")]
1498    Draft09,
1499    #[cfg(feature = "draft10")]
1500    Draft10,
1501    #[cfg(feature = "draft11")]
1502    Draft11 { extensions: bool },
1503    #[cfg(feature = "draft12")]
1504    Draft12 { extensions: bool },
1505    #[cfg(feature = "draft13")]
1506    Draft13 { extensions: bool },
1507    #[cfg(feature = "draft14")]
1508    Draft14(crate::draft14::data_stream::SubgroupObjectReader),
1509    #[cfg(feature = "draft15")]
1510    Draft15(crate::draft15::data_stream::SubgroupObjectReader),
1511    #[cfg(feature = "draft16")]
1512    Draft16(crate::draft16::data_stream::SubgroupObjectReader),
1513    #[cfg(feature = "draft17")]
1514    Draft17(crate::draft17::data_stream::SubgroupObjectReader),
1515    #[cfg(feature = "draft18")]
1516    Draft18(crate::draft18::data_stream::SubgroupObjectReader),
1517    #[cfg(feature = "draft19")]
1518    Draft19(crate::draft19::data_stream::SubgroupObjectReader),
1519    #[cfg(feature = "draft20")]
1520    Draft20(crate::draft20::data_stream::SubgroupObjectReader),
1521}
1522
1523/// Stateful reader for the objects on a subgroup data stream, for any
1524/// enabled draft.
1525///
1526/// Drafts 07-13 encode absolute object IDs and need no state, drafts 14-20
1527/// delta-encode them against the previous object. This reader presents both
1528/// as the same API: construct it from the stream's header, then call
1529/// [`read_object`](Self::read_object) once per object.
1530///
1531/// The reader is [`Clone`] specifically so callers can probe a partial buffer
1532/// against a copy and commit only on success; see the module docs.
1533#[derive(Debug, Clone)]
1534pub struct AnySubgroupObjectReader {
1535    state: SubgroupReaderState,
1536}
1537
1538impl AnySubgroupObjectReader {
1539    /// Create a reader seeded from the stream's subgroup header.
1540    ///
1541    /// Returns [`CodecError::UnsupportedDraft`] when the header's draft is
1542    /// not compiled in, and [`CodecError::InvalidField`] when the header's
1543    /// stream type is not a subgroup type.
1544    #[allow(unused_variables, unreachable_code)]
1545    pub fn new(header: &AnySubgroupHeader) -> Result<Self, CodecError> {
1546        let state = match header {
1547            #[cfg(feature = "draft07")]
1548            AnySubgroupHeader::Draft07(_) => SubgroupReaderState::Draft07,
1549            #[cfg(feature = "draft08")]
1550            AnySubgroupHeader::Draft08(_) => SubgroupReaderState::Draft08,
1551            #[cfg(feature = "draft09")]
1552            AnySubgroupHeader::Draft09(_) => SubgroupReaderState::Draft09,
1553            #[cfg(feature = "draft10")]
1554            AnySubgroupHeader::Draft10(_) => SubgroupReaderState::Draft10,
1555            #[cfg(feature = "draft11")]
1556            AnySubgroupHeader::Draft11(h) => {
1557                SubgroupReaderState::Draft11 { extensions: subgroup_extensions_11(h)? }
1558            }
1559            #[cfg(feature = "draft12")]
1560            AnySubgroupHeader::Draft12(h) => {
1561                SubgroupReaderState::Draft12 { extensions: subgroup_extensions_12(h)? }
1562            }
1563            #[cfg(feature = "draft13")]
1564            AnySubgroupHeader::Draft13(h) => {
1565                SubgroupReaderState::Draft13 { extensions: subgroup_extensions_13(h)? }
1566            }
1567            #[cfg(feature = "draft14")]
1568            AnySubgroupHeader::Draft14(h) => SubgroupReaderState::Draft14(
1569                crate::draft14::data_stream::SubgroupObjectReader::new(h),
1570            ),
1571            #[cfg(feature = "draft15")]
1572            AnySubgroupHeader::Draft15(h) => SubgroupReaderState::Draft15(
1573                crate::draft15::data_stream::SubgroupObjectReader::new(h),
1574            ),
1575            #[cfg(feature = "draft16")]
1576            AnySubgroupHeader::Draft16(h) => SubgroupReaderState::Draft16(
1577                crate::draft16::data_stream::SubgroupObjectReader::new(h),
1578            ),
1579            #[cfg(feature = "draft17")]
1580            AnySubgroupHeader::Draft17(h) => SubgroupReaderState::Draft17(
1581                crate::draft17::data_stream::SubgroupObjectReader::new(h),
1582            ),
1583            #[cfg(feature = "draft18")]
1584            AnySubgroupHeader::Draft18(h) => SubgroupReaderState::Draft18(
1585                crate::draft18::data_stream::SubgroupObjectReader::new(h),
1586            ),
1587            #[cfg(feature = "draft19")]
1588            AnySubgroupHeader::Draft19(h) => SubgroupReaderState::Draft19(
1589                crate::draft19::data_stream::SubgroupObjectReader::new(h),
1590            ),
1591            #[cfg(feature = "draft20")]
1592            AnySubgroupHeader::Draft20(h) => SubgroupReaderState::Draft20(
1593                crate::draft20::data_stream::SubgroupObjectReader::new(h),
1594            ),
1595            #[allow(unreachable_patterns)]
1596            _ => {
1597                return Err(CodecError::UnsupportedDraft(format!(
1598                    "draft {:?} not enabled via feature flag",
1599                    header.draft()
1600                )));
1601            }
1602        };
1603        Ok(Self { state })
1604    }
1605
1606    /// The draft this reader decodes.
1607    #[allow(unreachable_code)]
1608    pub fn draft(&self) -> DraftVersion {
1609        match &self.state {
1610            #[cfg(feature = "draft07")]
1611            SubgroupReaderState::Draft07 => DraftVersion::Draft07,
1612            #[cfg(feature = "draft08")]
1613            SubgroupReaderState::Draft08 => DraftVersion::Draft08,
1614            #[cfg(feature = "draft09")]
1615            SubgroupReaderState::Draft09 => DraftVersion::Draft09,
1616            #[cfg(feature = "draft10")]
1617            SubgroupReaderState::Draft10 => DraftVersion::Draft10,
1618            #[cfg(feature = "draft11")]
1619            SubgroupReaderState::Draft11 { .. } => DraftVersion::Draft11,
1620            #[cfg(feature = "draft12")]
1621            SubgroupReaderState::Draft12 { .. } => DraftVersion::Draft12,
1622            #[cfg(feature = "draft13")]
1623            SubgroupReaderState::Draft13 { .. } => DraftVersion::Draft13,
1624            #[cfg(feature = "draft14")]
1625            SubgroupReaderState::Draft14(_) => DraftVersion::Draft14,
1626            #[cfg(feature = "draft15")]
1627            SubgroupReaderState::Draft15(_) => DraftVersion::Draft15,
1628            #[cfg(feature = "draft16")]
1629            SubgroupReaderState::Draft16(_) => DraftVersion::Draft16,
1630            #[cfg(feature = "draft17")]
1631            SubgroupReaderState::Draft17(_) => DraftVersion::Draft17,
1632            #[cfg(feature = "draft18")]
1633            SubgroupReaderState::Draft18(_) => DraftVersion::Draft18,
1634            #[cfg(feature = "draft19")]
1635            SubgroupReaderState::Draft19(_) => DraftVersion::Draft19,
1636            #[cfg(feature = "draft20")]
1637            SubgroupReaderState::Draft20(_) => DraftVersion::Draft20,
1638            #[allow(unreachable_patterns)]
1639            _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1640        }
1641    }
1642
1643    /// Decode the next object, including its payload.
1644    ///
1645    /// Returns [`CodecError::UnexpectedEnd`] when `buf` holds only part of an
1646    /// object; the reader's state is unspecified after such an error, so
1647    /// callers that may be fed partial buffers must probe against a clone.
1648    #[allow(unused_variables, unreachable_code)]
1649    pub fn read_object(&mut self, buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
1650        match &mut self.state {
1651            #[cfg(feature = "draft07")]
1652            SubgroupReaderState::Draft07 => sg07::read_object(buf),
1653            #[cfg(feature = "draft08")]
1654            SubgroupReaderState::Draft08 => sg08::read_object(buf),
1655            #[cfg(feature = "draft09")]
1656            SubgroupReaderState::Draft09 => sg09::read_object(buf),
1657            #[cfg(feature = "draft10")]
1658            SubgroupReaderState::Draft10 => sg10::read_object(buf),
1659            #[cfg(feature = "draft11")]
1660            SubgroupReaderState::Draft11 { extensions } => sg11::read_object(*extensions, buf),
1661            #[cfg(feature = "draft12")]
1662            SubgroupReaderState::Draft12 { extensions } => sg12::read_object(*extensions, buf),
1663            #[cfg(feature = "draft13")]
1664            SubgroupReaderState::Draft13 { extensions } => sg13::read_object(*extensions, buf),
1665            #[cfg(feature = "draft14")]
1666            SubgroupReaderState::Draft14(inner) => sg14::read_object(inner, buf),
1667            #[cfg(feature = "draft15")]
1668            SubgroupReaderState::Draft15(inner) => sg15::read_object(inner, buf),
1669            #[cfg(feature = "draft16")]
1670            SubgroupReaderState::Draft16(inner) => sg16::read_object(inner, buf),
1671            #[cfg(feature = "draft17")]
1672            SubgroupReaderState::Draft17(inner) => sg17::read_object(inner, buf),
1673            #[cfg(feature = "draft18")]
1674            SubgroupReaderState::Draft18(inner) => sg18::read_object(inner, buf),
1675            #[cfg(feature = "draft19")]
1676            SubgroupReaderState::Draft19(inner) => sg19::read_object(inner, buf),
1677            #[cfg(feature = "draft20")]
1678            SubgroupReaderState::Draft20(inner) => sg20::read_object(inner, buf),
1679            #[allow(unreachable_patterns)]
1680            _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1681        }
1682    }
1683
1684    /// Decode the next object's framing without copying its payload.
1685    ///
1686    /// Advances `buf` past the whole object exactly as
1687    /// [`read_object`](Self::read_object) does, but returns only scalars.
1688    /// This is the path a relay uses when it forwards the object's bytes
1689    /// verbatim and never inspects the payload.
1690    #[allow(unused_variables, unreachable_code)]
1691    pub fn read_object_meta(
1692        &mut self,
1693        buf: &mut impl Buf,
1694    ) -> Result<AnySubgroupObjectMeta, CodecError> {
1695        match &mut self.state {
1696            #[cfg(feature = "draft07")]
1697            SubgroupReaderState::Draft07 => sg07::read_object_meta(buf),
1698            #[cfg(feature = "draft08")]
1699            SubgroupReaderState::Draft08 => sg08::read_object_meta(buf),
1700            #[cfg(feature = "draft09")]
1701            SubgroupReaderState::Draft09 => sg09::read_object_meta(buf),
1702            #[cfg(feature = "draft10")]
1703            SubgroupReaderState::Draft10 => sg10::read_object_meta(buf),
1704            #[cfg(feature = "draft11")]
1705            SubgroupReaderState::Draft11 { extensions } => sg11::read_object_meta(*extensions, buf),
1706            #[cfg(feature = "draft12")]
1707            SubgroupReaderState::Draft12 { extensions } => sg12::read_object_meta(*extensions, buf),
1708            #[cfg(feature = "draft13")]
1709            SubgroupReaderState::Draft13 { extensions } => sg13::read_object_meta(*extensions, buf),
1710            #[cfg(feature = "draft14")]
1711            SubgroupReaderState::Draft14(inner) => sg14::read_object_meta(inner, buf),
1712            #[cfg(feature = "draft15")]
1713            SubgroupReaderState::Draft15(inner) => sg15::read_object_meta(inner, buf),
1714            #[cfg(feature = "draft16")]
1715            SubgroupReaderState::Draft16(inner) => sg16::read_object_meta(inner, buf),
1716            #[cfg(feature = "draft17")]
1717            SubgroupReaderState::Draft17(inner) => sg17::read_object_meta(inner, buf),
1718            #[cfg(feature = "draft18")]
1719            SubgroupReaderState::Draft18(inner) => sg18::read_object_meta(inner, buf),
1720            #[cfg(feature = "draft19")]
1721            SubgroupReaderState::Draft19(inner) => sg19::read_object_meta(inner, buf),
1722            #[cfg(feature = "draft20")]
1723            SubgroupReaderState::Draft20(inner) => sg20::read_object_meta(inner, buf),
1724            #[allow(unreachable_patterns)]
1725            _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1726        }
1727    }
1728}
1729
1730// ── Subgroup object writer ──────────────────────────────────
1731
1732/// Per-draft writer state. Mirrors [`SubgroupReaderState`]; drafts 14-20
1733/// reuse each draft's `SubgroupObjectReader`, which owns both directions of
1734/// the delta state. Drafts 07-13 encode absolute IDs, so nothing on the wire
1735/// forces them to increase and they carry a `prev_object_id` of their own —
1736/// see `advance_absolute_id`.
1737#[derive(Debug, Clone)]
1738enum SubgroupWriterState {
1739    #[cfg(feature = "draft07")]
1740    Draft07 { prev_object_id: Option<u64> },
1741    #[cfg(feature = "draft08")]
1742    Draft08 { prev_object_id: Option<u64> },
1743    #[cfg(feature = "draft09")]
1744    Draft09 { prev_object_id: Option<u64> },
1745    #[cfg(feature = "draft10")]
1746    Draft10 { prev_object_id: Option<u64> },
1747    #[cfg(feature = "draft11")]
1748    Draft11 { extensions: bool, prev_object_id: Option<u64> },
1749    #[cfg(feature = "draft12")]
1750    Draft12 { extensions: bool, prev_object_id: Option<u64> },
1751    #[cfg(feature = "draft13")]
1752    Draft13 { extensions: bool, prev_object_id: Option<u64> },
1753    #[cfg(feature = "draft14")]
1754    Draft14 { inner: crate::draft14::data_stream::SubgroupObjectReader, extensions: bool },
1755    #[cfg(feature = "draft15")]
1756    Draft15 { inner: crate::draft15::data_stream::SubgroupObjectReader, extensions: bool },
1757    #[cfg(feature = "draft16")]
1758    Draft16 { inner: crate::draft16::data_stream::SubgroupObjectReader, extensions: bool },
1759    #[cfg(feature = "draft17")]
1760    Draft17 { inner: crate::draft17::data_stream::SubgroupObjectReader, extensions: bool },
1761    #[cfg(feature = "draft18")]
1762    Draft18 { inner: crate::draft18::data_stream::SubgroupObjectReader, extensions: bool },
1763    #[cfg(feature = "draft19")]
1764    Draft19 { inner: crate::draft19::data_stream::SubgroupObjectReader, extensions: bool },
1765    #[cfg(feature = "draft20")]
1766    Draft20 { inner: crate::draft20::data_stream::SubgroupObjectReader, extensions: bool },
1767}
1768
1769/// Serializer for the objects on a subgroup data stream, for any enabled
1770/// draft.
1771///
1772/// Mirrors [`AnySubgroupObjectReader`]. On drafts 14-20 it tracks the
1773/// previous Object ID so successive writes produce correct deltas; on drafts
1774/// 07-13 object IDs are absolute and the same state only enforces that they
1775/// increase.
1776///
1777/// # Eliding objects
1778///
1779/// To remove an object from a stream, read it and then simply do not write
1780/// it. The writer's delta state advances only when
1781/// [`write_object`](Self::write_object) succeeds, so the next object written
1782/// re-derives its delta against the last *retained* object automatically.
1783/// See [`Self::write_object`] for the exact invariant.
1784#[derive(Debug, Clone)]
1785pub struct AnySubgroupObjectWriter {
1786    state: SubgroupWriterState,
1787}
1788
1789impl AnySubgroupObjectWriter {
1790    /// Create a writer for a stream with the given header.
1791    ///
1792    /// The header fixes the extension-presence and (on drafts 11-13) stream
1793    /// type used for every object written, exactly as it does for
1794    /// [`AnySubgroupObjectReader::new`].
1795    #[allow(unused_variables, unreachable_code)]
1796    pub fn new(header: &AnySubgroupHeader) -> Result<Self, CodecError> {
1797        let state = match header {
1798            #[cfg(feature = "draft07")]
1799            AnySubgroupHeader::Draft07(_) => SubgroupWriterState::Draft07 { prev_object_id: None },
1800            #[cfg(feature = "draft08")]
1801            AnySubgroupHeader::Draft08(_) => SubgroupWriterState::Draft08 { prev_object_id: None },
1802            #[cfg(feature = "draft09")]
1803            AnySubgroupHeader::Draft09(_) => SubgroupWriterState::Draft09 { prev_object_id: None },
1804            #[cfg(feature = "draft10")]
1805            AnySubgroupHeader::Draft10(_) => SubgroupWriterState::Draft10 { prev_object_id: None },
1806            #[cfg(feature = "draft11")]
1807            AnySubgroupHeader::Draft11(h) => SubgroupWriterState::Draft11 {
1808                extensions: subgroup_extensions_11(h)?,
1809                prev_object_id: None,
1810            },
1811            #[cfg(feature = "draft12")]
1812            AnySubgroupHeader::Draft12(h) => SubgroupWriterState::Draft12 {
1813                extensions: subgroup_extensions_12(h)?,
1814                prev_object_id: None,
1815            },
1816            #[cfg(feature = "draft13")]
1817            AnySubgroupHeader::Draft13(h) => SubgroupWriterState::Draft13 {
1818                extensions: subgroup_extensions_13(h)?,
1819                prev_object_id: None,
1820            },
1821            #[cfg(feature = "draft14")]
1822            AnySubgroupHeader::Draft14(h) => SubgroupWriterState::Draft14 {
1823                inner: crate::draft14::data_stream::SubgroupObjectReader::new(h),
1824                extensions: h.stream_type.extensions_present(),
1825            },
1826            #[cfg(feature = "draft15")]
1827            AnySubgroupHeader::Draft15(h) => SubgroupWriterState::Draft15 {
1828                inner: crate::draft15::data_stream::SubgroupObjectReader::new(h),
1829                extensions: h.has_extensions(),
1830            },
1831            #[cfg(feature = "draft16")]
1832            AnySubgroupHeader::Draft16(h) => SubgroupWriterState::Draft16 {
1833                inner: crate::draft16::data_stream::SubgroupObjectReader::new(h),
1834                extensions: h.has_extensions(),
1835            },
1836            #[cfg(feature = "draft17")]
1837            AnySubgroupHeader::Draft17(h) => SubgroupWriterState::Draft17 {
1838                inner: crate::draft17::data_stream::SubgroupObjectReader::new(h),
1839                extensions: h.has_properties(),
1840            },
1841            #[cfg(feature = "draft18")]
1842            AnySubgroupHeader::Draft18(h) => SubgroupWriterState::Draft18 {
1843                inner: crate::draft18::data_stream::SubgroupObjectReader::new(h),
1844                extensions: h.has_properties(),
1845            },
1846            #[cfg(feature = "draft19")]
1847            AnySubgroupHeader::Draft19(h) => SubgroupWriterState::Draft19 {
1848                inner: crate::draft19::data_stream::SubgroupObjectReader::new(h),
1849                extensions: h.has_properties(),
1850            },
1851            #[cfg(feature = "draft20")]
1852            AnySubgroupHeader::Draft20(h) => SubgroupWriterState::Draft20 {
1853                inner: crate::draft20::data_stream::SubgroupObjectReader::new(h),
1854                extensions: h.has_properties(),
1855            },
1856            #[allow(unreachable_patterns)]
1857            _ => {
1858                return Err(CodecError::UnsupportedDraft(format!(
1859                    "draft {:?} not enabled via feature flag",
1860                    header.draft()
1861                )));
1862            }
1863        };
1864        Ok(Self { state })
1865    }
1866
1867    /// The draft this writer encodes.
1868    #[allow(unreachable_code)]
1869    pub fn draft(&self) -> DraftVersion {
1870        match &self.state {
1871            #[cfg(feature = "draft07")]
1872            SubgroupWriterState::Draft07 { .. } => DraftVersion::Draft07,
1873            #[cfg(feature = "draft08")]
1874            SubgroupWriterState::Draft08 { .. } => DraftVersion::Draft08,
1875            #[cfg(feature = "draft09")]
1876            SubgroupWriterState::Draft09 { .. } => DraftVersion::Draft09,
1877            #[cfg(feature = "draft10")]
1878            SubgroupWriterState::Draft10 { .. } => DraftVersion::Draft10,
1879            #[cfg(feature = "draft11")]
1880            SubgroupWriterState::Draft11 { .. } => DraftVersion::Draft11,
1881            #[cfg(feature = "draft12")]
1882            SubgroupWriterState::Draft12 { .. } => DraftVersion::Draft12,
1883            #[cfg(feature = "draft13")]
1884            SubgroupWriterState::Draft13 { .. } => DraftVersion::Draft13,
1885            #[cfg(feature = "draft14")]
1886            SubgroupWriterState::Draft14 { .. } => DraftVersion::Draft14,
1887            #[cfg(feature = "draft15")]
1888            SubgroupWriterState::Draft15 { .. } => DraftVersion::Draft15,
1889            #[cfg(feature = "draft16")]
1890            SubgroupWriterState::Draft16 { .. } => DraftVersion::Draft16,
1891            #[cfg(feature = "draft17")]
1892            SubgroupWriterState::Draft17 { .. } => DraftVersion::Draft17,
1893            #[cfg(feature = "draft18")]
1894            SubgroupWriterState::Draft18 { .. } => DraftVersion::Draft18,
1895            #[cfg(feature = "draft19")]
1896            SubgroupWriterState::Draft19 { .. } => DraftVersion::Draft19,
1897            #[cfg(feature = "draft20")]
1898            SubgroupWriterState::Draft20 { .. } => DraftVersion::Draft20,
1899            #[allow(unreachable_patterns)]
1900            _ => unreachable!("AnySubgroupObjectWriter has no enabled variants"),
1901        }
1902    }
1903
1904    /// Encode one object, advancing the delta state.
1905    ///
1906    /// # Invariant
1907    ///
1908    /// Let a stream's objects decode to absolute IDs `a_0, a_1, .., a_n`.
1909    /// Feeding any strictly-increasing subsequence of those objects through
1910    /// one writer, in order, produces a byte stream that decodes back to
1911    /// exactly that subsequence of absolute IDs, on every draft 07-20.
1912    ///
1913    /// Concretely: dropping `a_2` from `0,1,2,3,4` yields a stream decoding
1914    /// to `0,1,3,4` — not `0,1,2,3`.
1915    ///
1916    /// # Errors
1917    ///
1918    /// [`CodecError::InvalidField`] when `object.object_id` is not strictly
1919    /// greater than the previously written object's ID (two objects on a
1920    /// subgroup stream can never share an ID, so no valid delta exists), when
1921    /// a computed delta or length exceeds the varint range, when the object
1922    /// carries extension bytes that a stream without an extension block
1923    /// cannot represent, or when a non-empty payload is paired with a status.
1924    ///
1925    /// Also [`CodecError::InvalidField`] when `object.status` holds a code the
1926    /// draft being written does not assign. [`AnySubgroupObject::status`] is a
1927    /// raw wire code because it crosses drafts, and the assigned set moves
1928    /// between them, so a status read off one draft's stream is not
1929    /// necessarily writable onto another's: forwarding a draft-15 Object Does
1930    /// Not Exist (0x1) onto a draft-16 or later stream is refused here rather
1931    /// than emitted as a byte the peer must close the session over.
1932    #[allow(unused_variables, unreachable_code)]
1933    pub fn write_object(
1934        &mut self,
1935        object: &AnySubgroupObject,
1936        buf: &mut impl BufMut,
1937    ) -> Result<(), CodecError> {
1938        match &mut self.state {
1939            #[cfg(feature = "draft07")]
1940            SubgroupWriterState::Draft07 { prev_object_id } => {
1941                advance_absolute_id(prev_object_id, object, |o| sg07::write_object(o, buf))
1942            }
1943            #[cfg(feature = "draft08")]
1944            SubgroupWriterState::Draft08 { prev_object_id } => {
1945                advance_absolute_id(prev_object_id, object, |o| sg08::write_object(o, buf))
1946            }
1947            #[cfg(feature = "draft09")]
1948            SubgroupWriterState::Draft09 { prev_object_id } => {
1949                advance_absolute_id(prev_object_id, object, |o| sg09::write_object(o, buf))
1950            }
1951            #[cfg(feature = "draft10")]
1952            SubgroupWriterState::Draft10 { prev_object_id } => {
1953                advance_absolute_id(prev_object_id, object, |o| sg10::write_object(o, buf))
1954            }
1955            #[cfg(feature = "draft11")]
1956            SubgroupWriterState::Draft11 { extensions, prev_object_id } => {
1957                let extensions = *extensions;
1958                advance_absolute_id(prev_object_id, object, |o| {
1959                    sg11::write_object(extensions, o, buf)
1960                })
1961            }
1962            #[cfg(feature = "draft12")]
1963            SubgroupWriterState::Draft12 { extensions, prev_object_id } => {
1964                let extensions = *extensions;
1965                advance_absolute_id(prev_object_id, object, |o| {
1966                    sg12::write_object(extensions, o, buf)
1967                })
1968            }
1969            #[cfg(feature = "draft13")]
1970            SubgroupWriterState::Draft13 { extensions, prev_object_id } => {
1971                let extensions = *extensions;
1972                advance_absolute_id(prev_object_id, object, |o| {
1973                    sg13::write_object(extensions, o, buf)
1974                })
1975            }
1976            #[cfg(feature = "draft14")]
1977            SubgroupWriterState::Draft14 { inner, extensions } => {
1978                reject_unrepresentable_extensions(*extensions, object)?;
1979                sg14::write_object(inner, object, buf)
1980            }
1981            #[cfg(feature = "draft15")]
1982            SubgroupWriterState::Draft15 { inner, extensions } => {
1983                reject_unrepresentable_extensions(*extensions, object)?;
1984                sg15::write_object(inner, object, buf)
1985            }
1986            #[cfg(feature = "draft16")]
1987            SubgroupWriterState::Draft16 { inner, extensions } => {
1988                reject_unrepresentable_extensions(*extensions, object)?;
1989                sg16::write_object(inner, object, buf)
1990            }
1991            #[cfg(feature = "draft17")]
1992            SubgroupWriterState::Draft17 { inner, extensions } => {
1993                reject_unrepresentable_extensions(*extensions, object)?;
1994                sg17::write_object(inner, object, buf)
1995            }
1996            #[cfg(feature = "draft18")]
1997            SubgroupWriterState::Draft18 { inner, extensions } => {
1998                reject_unrepresentable_extensions(*extensions, object)?;
1999                sg18::write_object(inner, object, buf)
2000            }
2001            #[cfg(feature = "draft19")]
2002            SubgroupWriterState::Draft19 { inner, extensions } => {
2003                reject_unrepresentable_extensions(*extensions, object)?;
2004                sg19::write_object(inner, object, buf)
2005            }
2006            #[cfg(feature = "draft20")]
2007            SubgroupWriterState::Draft20 { inner, extensions } => {
2008                reject_unrepresentable_extensions(*extensions, object)?;
2009                sg20::write_object(inner, object, buf)
2010            }
2011            #[allow(unreachable_patterns)]
2012            _ => unreachable!("AnySubgroupObjectWriter has no enabled variants"),
2013        }
2014    }
2015}
2016
2017// ── Re-emitting an object whose bytes are already known ─────
2018
2019/// What [`reemit_subgroup_object`] had to do.
2020#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2021pub enum Reemit {
2022    /// The bytes were copied unchanged.
2023    Verbatim,
2024    /// Only the leading Object ID field changed.
2025    Reencoded {
2026        /// Bytes the ID field occupied in `raw`.
2027        id_bytes_before: usize,
2028        /// Bytes it occupies in the output.
2029        id_bytes_after: usize,
2030    },
2031}
2032
2033/// Re-emit a subgroup object whose wire bytes are already known, adjusting
2034/// only this draft's encoding of its identity.
2035///
2036/// This is the whole of what removing an object from a subgroup stream
2037/// costs. Drafts 07-13 encode absolute Object IDs, so every survivor's
2038/// bytes are already correct and this copies `raw` unchanged after checking
2039/// that IDs still increase. Drafts 14-20 encode `id - prev - 1`, so the
2040/// leading varint is recomputed against `prev_forwarded`; when its minimal
2041/// encoding is byte-identical to the one in `raw` the bytes are still
2042/// copied unchanged. Everything after the ID field — extension block,
2043/// length, status, payload — is always copied verbatim.
2044///
2045/// After one object is re-emitted following an elided run, the writer's
2046/// cursor re-converges with the reader's, so every later object's original
2047/// bytes remain correct. An elide therefore costs at most one fix-up, not a
2048/// re-encode of the stream's tail.
2049///
2050/// `prev_forwarded` is the absolute Object ID of the last object actually
2051/// forwarded on this stream, or `None` when none has been.
2052///
2053/// # `raw` need not be a complete object
2054///
2055/// **Any prefix is legal provided the whole leading Object ID field is
2056/// present.** Everything after that field is copied byte-for-byte, however
2057/// many bytes there are, and **no length validation is performed** — this
2058/// function never reads the extension block, never reads the payload
2059/// length field and never compares it to `raw.len()`. It cannot: on drafts
2060/// 07-13 it does not decode past the ID at all, and on 14-20 it decodes
2061/// exactly one varint.
2062///
2063/// This is what lets a caller fix up the **first chunk of an oversized
2064/// object** — an object too large to buffer is forwarded in chunks, and only
2065/// the first one carries the ID field. A caller that cannot guarantee the ID
2066/// field is whole in the chunk it passes gets
2067/// [`CodecError::InvalidField`] rather than a silent truncation.
2068///
2069/// **Do not add a completeness check.** A `raw.len() >= wire_len` assertion
2070/// would look defensive, would pass every test that feeds it whole objects,
2071/// and would refuse the prefix this function exists to accept.
2072///
2073/// # Errors
2074///
2075/// [`CodecError::InvalidField`] when `object_id` is not strictly greater
2076/// than `prev_forwarded`, when the recomputed delta exceeds the varint
2077/// range, or when `raw` does not begin with a decodable varint — which
2078/// includes a `raw` too short to hold the whole leading varint.
2079///
2080/// # Examples
2081///
2082/// ```
2083/// use moqtap_codec::dispatch::{reemit_subgroup_object, Reemit};
2084/// use moqtap_codec::version::DraftVersion;
2085///
2086/// // A draft-19 object that was encoded as the successor of ID 4 —
2087/// // leading delta 0 — re-emitted after ID 3 was the last one forwarded.
2088/// let raw = [0x00, 0x02, 0xca, 0xfe];
2089/// let mut out = Vec::new();
2090/// let what = reemit_subgroup_object(DraftVersion::Draft19, Some(3), 5, &raw, &mut out).unwrap();
2091/// assert_eq!(what, Reemit::Reencoded { id_bytes_before: 1, id_bytes_after: 1 });
2092/// assert_eq!(out, [0x01, 0x02, 0xca, 0xfe]);
2093/// ```
2094pub fn reemit_subgroup_object(
2095    draft: DraftVersion,
2096    prev_forwarded: Option<u64>,
2097    object_id: u64,
2098    raw: &[u8],
2099    out: &mut impl BufMut,
2100) -> Result<Reemit, CodecError> {
2101    if matches!(prev_forwarded, Some(prev) if object_id <= prev) {
2102        return Err(CodecError::InvalidField);
2103    }
2104
2105    // Measure the ID field. Every draft 07-20 puts it first and nothing past
2106    // it is decoded, so `raw` may stop anywhere after it. Which varint measures
2107    // it depends on the draft: 17 replaced the RFC 9000 encoding with MoQT's.
2108    let mut cursor: &[u8] = raw;
2109    draft.decode_varint(&mut cursor).map_err(|_| CodecError::InvalidField)?;
2110    let id_bytes_before = raw.len() - cursor.len();
2111
2112    // Drafts 07-13 write the ID absolutely: nothing about that field depends
2113    // on which objects were forwarded, so the bytes already say the truth.
2114    if !delta_encodes_object_ids(draft) {
2115        out.put_slice(raw);
2116        return Ok(Reemit::Verbatim);
2117    }
2118
2119    let delta = match prev_forwarded {
2120        None => object_id,
2121        Some(prev) => object_id
2122            .checked_sub(prev)
2123            .and_then(|v| v.checked_sub(1))
2124            .ok_or(CodecError::InvalidField)?,
2125    };
2126
2127    // The MoQT encoding reaches the full 64-bit range, so a delta a draft-17+
2128    // peer can legitimately send is not an error there.
2129    let field = if draft.uses_moqt_varint() {
2130        VarInt::from_u64_moqt(delta)
2131    } else {
2132        VarInt::from_u64(delta).map_err(|_| CodecError::InvalidField)?
2133    };
2134
2135    // Nine bytes: the MoQT encoding is one longer than RFC 9000 at the top.
2136    let mut encoded = [0u8; 9];
2137    let mut slot: &mut [u8] = &mut encoded;
2138    draft.encode_varint(field, &mut slot);
2139    let id_bytes_after = 9 - slot.len();
2140    let encoded = &encoded[..id_bytes_after];
2141
2142    if encoded == &raw[..id_bytes_before] {
2143        out.put_slice(raw);
2144        return Ok(Reemit::Verbatim);
2145    }
2146
2147    out.put_slice(encoded);
2148    out.put_slice(&raw[id_bytes_before..]);
2149    Ok(Reemit::Reencoded { id_bytes_before, id_bytes_after })
2150}
2151
2152/// `true` on the drafts whose subgroup objects encode the Object ID as
2153/// `id - prev - 1` rather than absolutely.
2154///
2155/// Needs no `#[cfg]`: [`DraftVersion`] is not feature-gated, so this answers
2156/// for a draft whose codec is not compiled in.
2157fn delta_encodes_object_ids(draft: DraftVersion) -> bool {
2158    matches!(
2159        draft,
2160        DraftVersion::Draft14
2161            | DraftVersion::Draft15
2162            | DraftVersion::Draft16
2163            | DraftVersion::Draft17
2164            | DraftVersion::Draft18
2165            | DraftVersion::Draft19
2166            | DraftVersion::Draft20
2167    )
2168}
2169
2170/// Enforce the strictly-increasing Object ID rule on the drafts that encode
2171/// IDs absolutely.
2172///
2173/// Drafts 14-20 get this for free: their delta is `id - prev - 1`, so a
2174/// repeated or decreasing ID underflows and the per-draft writer rejects it.
2175/// Drafts 07-13 write the ID verbatim and would happily emit a stream no
2176/// publisher can produce, so the check lives here. As on the delta drafts, the
2177/// state advances only once the object is actually written, which is what makes
2178/// elision *read it and do not write it*.
2179#[cfg(any(
2180    feature = "draft07",
2181    feature = "draft08",
2182    feature = "draft09",
2183    feature = "draft10",
2184    feature = "draft11",
2185    feature = "draft12",
2186    feature = "draft13"
2187))]
2188fn advance_absolute_id(
2189    prev_object_id: &mut Option<u64>,
2190    object: &AnySubgroupObject,
2191    write: impl FnOnce(&AnySubgroupObject) -> Result<(), CodecError>,
2192) -> Result<(), CodecError> {
2193    if matches!(*prev_object_id, Some(prev) if object.object_id <= prev) {
2194        return Err(CodecError::InvalidField);
2195    }
2196    write(object)?;
2197    *prev_object_id = Some(object.object_id);
2198    Ok(())
2199}
2200
2201/// A stream whose header says objects carry no extension block cannot encode
2202/// one, so refuse rather than drop the bytes.
2203#[cfg(any(
2204    feature = "draft14",
2205    feature = "draft15",
2206    feature = "draft16",
2207    feature = "draft17",
2208    feature = "draft18",
2209    feature = "draft19",
2210    feature = "draft20"
2211))]
2212fn reject_unrepresentable_extensions(
2213    extensions: bool,
2214    object: &AnySubgroupObject,
2215) -> Result<(), CodecError> {
2216    if !extensions && !object.extension_headers.is_empty() {
2217        return Err(CodecError::InvalidField);
2218    }
2219    Ok(())
2220}
2221
2222// ── Fetch object reader ─────────────────────────────────────
2223
2224/// Per-draft fetch reader state. Fetch objects are self-describing on drafts
2225/// 07-14, so those variants carry none; drafts 15-20 let an object take fields
2226/// from the one before it, so each owns the running state that resolves them.
2227#[derive(Debug, Clone)]
2228enum FetchReaderState {
2229    #[cfg(feature = "draft07")]
2230    Draft07,
2231    #[cfg(feature = "draft08")]
2232    Draft08,
2233    #[cfg(feature = "draft09")]
2234    Draft09,
2235    #[cfg(feature = "draft10")]
2236    Draft10,
2237    #[cfg(feature = "draft11")]
2238    Draft11,
2239    #[cfg(feature = "draft12")]
2240    Draft12,
2241    #[cfg(feature = "draft13")]
2242    Draft13,
2243    #[cfg(feature = "draft14")]
2244    Draft14,
2245    #[cfg(feature = "draft15")]
2246    Draft15(crate::draft15::data_stream::FetchObjectReader),
2247    #[cfg(feature = "draft16")]
2248    Draft16(crate::draft16::data_stream::FetchObjectReader),
2249    #[cfg(feature = "draft17")]
2250    Draft17(crate::draft17::data_stream::FetchObjectReader),
2251    #[cfg(feature = "draft18")]
2252    Draft18(crate::draft18::data_stream::FetchObjectReader),
2253    #[cfg(feature = "draft19")]
2254    Draft19(crate::draft19::data_stream::FetchObjectReader),
2255    #[cfg(feature = "draft20")]
2256    Draft20(crate::draft20::data_stream::FetchObjectReader),
2257}
2258
2259/// Stateful reader for the frames on a fetch data stream, for any enabled
2260/// draft.
2261///
2262/// Fetch objects are self-describing on drafts 07-14 and this reader carries no
2263/// state there. From draft-15 a Serialization Flags field decides which of an
2264/// object's Group ID, Subgroup ID, Object ID and Priority reach the wire at
2265/// all, and every field it omits is the object before it on the stream —
2266/// repeated, or stepped by one, or (from draft-18) counted from by a
2267/// difference. This reader holds that running state, so the values it produces
2268/// are absolute on every draft.
2269///
2270/// One reader belongs to one stream. Every draft counts "the prior Object"
2271/// along a single stream, so sharing a reader between streams, or restarting
2272/// one mid-stream, resolves later frames onto the wrong group, subgroup, ID or
2273/// priority — usually without an error anywhere.
2274///
2275/// The reader is [`Clone`] specifically so callers can probe a partial buffer
2276/// against a copy and commit only on success; see the module docs.
2277///
2278/// # Frames that are not objects
2279///
2280/// Drafts 16-20 add End of Range indicators, which state that a run of Objects
2281/// was not serialized. They arrive through the same calls as objects and are
2282/// told apart by [`AnyFetchObject::end_of_range`].
2283#[derive(Debug, Clone)]
2284pub struct AnyFetchObjectReader {
2285    state: FetchReaderState,
2286}
2287
2288impl AnyFetchObjectReader {
2289    /// Create a reader from the stream's fetch header and the Group Order the
2290    /// fetch was opened with.
2291    ///
2292    /// The order matters only on drafts 18 and 19, where an Object's Group ID
2293    /// is a difference from the previous Object's and the order decides its
2294    /// sign. Nothing on the data stream carries it — the FETCH settles it — and
2295    /// it is an argument rather than a default because a descending stream read
2296    /// as ascending does not fail: it decodes, under Group IDs walking the wrong
2297    /// way, and neither this crate nor the caller can tell afterwards. Both
2298    /// readings are legal streams.
2299    ///
2300    /// [`AnyControlMessage::fetch_group_order`](crate::dispatch::AnyControlMessage::fetch_group_order)
2301    /// answers it from the FETCH, including the case where the message names no
2302    /// GROUP_ORDER — draft-19 Section 10.2.8: "If omitted from FETCH, the
2303    /// receiver uses Ascending (0x1)". On drafts 07-17 the argument is ignored.
2304    ///
2305    /// Returns [`CodecError::UnsupportedDraft`] for drafts not compiled in.
2306    #[allow(unused_variables, unreachable_code)]
2307    pub fn new(
2308        header: &AnyFetchHeader,
2309        group_order: AnyFetchGroupOrder,
2310    ) -> Result<Self, CodecError> {
2311        let state = match header {
2312            #[cfg(feature = "draft07")]
2313            AnyFetchHeader::Draft07(_) => FetchReaderState::Draft07,
2314            #[cfg(feature = "draft08")]
2315            AnyFetchHeader::Draft08(_) => FetchReaderState::Draft08,
2316            #[cfg(feature = "draft09")]
2317            AnyFetchHeader::Draft09(_) => FetchReaderState::Draft09,
2318            #[cfg(feature = "draft10")]
2319            AnyFetchHeader::Draft10(_) => FetchReaderState::Draft10,
2320            #[cfg(feature = "draft11")]
2321            AnyFetchHeader::Draft11(_) => FetchReaderState::Draft11,
2322            #[cfg(feature = "draft12")]
2323            AnyFetchHeader::Draft12(_) => FetchReaderState::Draft12,
2324            #[cfg(feature = "draft13")]
2325            AnyFetchHeader::Draft13(_) => FetchReaderState::Draft13,
2326            #[cfg(feature = "draft14")]
2327            AnyFetchHeader::Draft14(_) => FetchReaderState::Draft14,
2328            // The header carries only a request id on drafts 15-20, so nothing
2329            // about it seeds the reader; the first object does.
2330            #[cfg(feature = "draft15")]
2331            AnyFetchHeader::Draft15(_) => {
2332                FetchReaderState::Draft15(crate::draft15::data_stream::FetchObjectReader::new())
2333            }
2334            #[cfg(feature = "draft16")]
2335            AnyFetchHeader::Draft16(_) => {
2336                FetchReaderState::Draft16(crate::draft16::data_stream::FetchObjectReader::new())
2337            }
2338            #[cfg(feature = "draft17")]
2339            AnyFetchHeader::Draft17(_) => {
2340                FetchReaderState::Draft17(crate::draft17::data_stream::FetchObjectReader::new())
2341            }
2342            #[cfg(feature = "draft18")]
2343            AnyFetchHeader::Draft18(_) => FetchReaderState::Draft18(
2344                crate::draft18::data_stream::FetchObjectReader::new(match group_order {
2345                    AnyFetchGroupOrder::Ascending => {
2346                        crate::draft18::data_stream::GroupOrder::Ascending
2347                    }
2348                    AnyFetchGroupOrder::Descending => {
2349                        crate::draft18::data_stream::GroupOrder::Descending
2350                    }
2351                }),
2352            ),
2353            #[cfg(feature = "draft19")]
2354            AnyFetchHeader::Draft19(_) => FetchReaderState::Draft19(
2355                crate::draft19::data_stream::FetchObjectReader::new(match group_order {
2356                    AnyFetchGroupOrder::Ascending => {
2357                        crate::draft19::data_stream::GroupOrder::Ascending
2358                    }
2359                    AnyFetchGroupOrder::Descending => {
2360                        crate::draft19::data_stream::GroupOrder::Descending
2361                    }
2362                }),
2363            ),
2364            #[cfg(feature = "draft20")]
2365            AnyFetchHeader::Draft20(_) => FetchReaderState::Draft20(
2366                crate::draft20::data_stream::FetchObjectReader::new(match group_order {
2367                    AnyFetchGroupOrder::Ascending => {
2368                        crate::draft20::data_stream::GroupOrder::Ascending
2369                    }
2370                    AnyFetchGroupOrder::Descending => {
2371                        crate::draft20::data_stream::GroupOrder::Descending
2372                    }
2373                }),
2374            ),
2375            #[allow(unreachable_patterns)]
2376            _ => {
2377                return Err(CodecError::UnsupportedDraft(format!(
2378                    "draft {:?} not enabled via feature flag",
2379                    header.draft()
2380                )));
2381            }
2382        };
2383        Ok(Self { state })
2384    }
2385
2386    /// The draft this reader decodes.
2387    #[allow(unreachable_code)]
2388    pub fn draft(&self) -> DraftVersion {
2389        match &self.state {
2390            #[cfg(feature = "draft07")]
2391            FetchReaderState::Draft07 => DraftVersion::Draft07,
2392            #[cfg(feature = "draft08")]
2393            FetchReaderState::Draft08 => DraftVersion::Draft08,
2394            #[cfg(feature = "draft09")]
2395            FetchReaderState::Draft09 => DraftVersion::Draft09,
2396            #[cfg(feature = "draft10")]
2397            FetchReaderState::Draft10 => DraftVersion::Draft10,
2398            #[cfg(feature = "draft11")]
2399            FetchReaderState::Draft11 => DraftVersion::Draft11,
2400            #[cfg(feature = "draft12")]
2401            FetchReaderState::Draft12 => DraftVersion::Draft12,
2402            #[cfg(feature = "draft13")]
2403            FetchReaderState::Draft13 => DraftVersion::Draft13,
2404            #[cfg(feature = "draft14")]
2405            FetchReaderState::Draft14 => DraftVersion::Draft14,
2406            #[cfg(feature = "draft15")]
2407            FetchReaderState::Draft15(_) => DraftVersion::Draft15,
2408            #[cfg(feature = "draft16")]
2409            FetchReaderState::Draft16(_) => DraftVersion::Draft16,
2410            #[cfg(feature = "draft17")]
2411            FetchReaderState::Draft17(_) => DraftVersion::Draft17,
2412            #[cfg(feature = "draft18")]
2413            FetchReaderState::Draft18(_) => DraftVersion::Draft18,
2414            #[cfg(feature = "draft19")]
2415            FetchReaderState::Draft19(_) => DraftVersion::Draft19,
2416            #[cfg(feature = "draft20")]
2417            FetchReaderState::Draft20(_) => DraftVersion::Draft20,
2418            #[allow(unreachable_patterns)]
2419            _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2420        }
2421    }
2422
2423    /// Decode the next fetch frame, including its payload.
2424    ///
2425    /// Returns [`CodecError::UnexpectedEnd`] when `buf` holds only part of a
2426    /// frame; the reader's state is unspecified after such an error, so callers
2427    /// that may be fed partial buffers must probe against a clone.
2428    ///
2429    /// Returns [`CodecError::InvalidField`] on drafts 15-20 when a frame takes
2430    /// a field from an object before it that does not exist — the first frame
2431    /// of a stream doing so is a protocol violation on every one of those
2432    /// drafts — and when a resolved Group ID, Subgroup ID or Object ID would
2433    /// leave the 64-bit range.
2434    #[allow(unused_variables, unreachable_code)]
2435    pub fn read_object(&mut self, buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
2436        match &mut self.state {
2437            #[cfg(feature = "draft07")]
2438            FetchReaderState::Draft07 => fo07::read_object(buf),
2439            #[cfg(feature = "draft08")]
2440            FetchReaderState::Draft08 => fo08::read_object(buf),
2441            #[cfg(feature = "draft09")]
2442            FetchReaderState::Draft09 => fo09::read_object(buf),
2443            #[cfg(feature = "draft10")]
2444            FetchReaderState::Draft10 => fo10::read_object(buf),
2445            #[cfg(feature = "draft11")]
2446            FetchReaderState::Draft11 => fo11::read_object(buf),
2447            #[cfg(feature = "draft12")]
2448            FetchReaderState::Draft12 => fo12::read_object(buf),
2449            #[cfg(feature = "draft13")]
2450            FetchReaderState::Draft13 => fo13::read_object(buf),
2451            #[cfg(feature = "draft14")]
2452            FetchReaderState::Draft14 => fo14::read_object(buf),
2453            #[cfg(feature = "draft15")]
2454            FetchReaderState::Draft15(inner) => fo15::read_object(inner, buf),
2455            #[cfg(feature = "draft16")]
2456            FetchReaderState::Draft16(inner) => fo16::read_object(inner, buf),
2457            #[cfg(feature = "draft17")]
2458            FetchReaderState::Draft17(inner) => fo17::read_object(inner, buf),
2459            #[cfg(feature = "draft18")]
2460            FetchReaderState::Draft18(inner) => fo18::read_object(inner, buf),
2461            #[cfg(feature = "draft19")]
2462            FetchReaderState::Draft19(inner) => fo19::read_object(inner, buf),
2463            #[cfg(feature = "draft20")]
2464            FetchReaderState::Draft20(inner) => fo20::read_object(inner, buf),
2465            #[allow(unreachable_patterns)]
2466            _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2467        }
2468    }
2469
2470    /// Decode the next fetch frame, keeping what re-encoding it later takes.
2471    ///
2472    /// Advances `buf` and this reader exactly as
2473    /// [`read_object_meta`](Self::read_object_meta) does, and reports the same
2474    /// framing in [`AnyFetchFrame::meta`]. What it additionally keeps is the
2475    /// shape the frame arrived in, which is the whole of what
2476    /// [`AnyFetchObjectWriter::reemit_object`] needs to write the frame back
2477    /// out against a different predecessor.
2478    ///
2479    /// Costs nothing over `read_object_meta`, which is itself defined over
2480    /// this: the per-draft header it keeps is one the decode produced and
2481    /// dropped.
2482    #[allow(unused_variables, unreachable_code)]
2483    pub fn read_object_frame(&mut self, buf: &mut impl Buf) -> Result<AnyFetchFrame, CodecError> {
2484        match &mut self.state {
2485            #[cfg(feature = "draft07")]
2486            FetchReaderState::Draft07 => fo07::read_object_meta(buf)
2487                .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft07, meta)),
2488            #[cfg(feature = "draft08")]
2489            FetchReaderState::Draft08 => fo08::read_object_meta(buf)
2490                .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft08, meta)),
2491            #[cfg(feature = "draft09")]
2492            FetchReaderState::Draft09 => fo09::read_object_meta(buf)
2493                .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft09, meta)),
2494            #[cfg(feature = "draft10")]
2495            FetchReaderState::Draft10 => fo10::read_object_meta(buf)
2496                .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft10, meta)),
2497            #[cfg(feature = "draft11")]
2498            FetchReaderState::Draft11 => fo11::read_object_meta(buf)
2499                .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft11, meta)),
2500            #[cfg(feature = "draft12")]
2501            FetchReaderState::Draft12 => fo12::read_object_meta(buf)
2502                .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft12, meta)),
2503            #[cfg(feature = "draft13")]
2504            FetchReaderState::Draft13 => fo13::read_object_meta(buf)
2505                .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft13, meta)),
2506            #[cfg(feature = "draft14")]
2507            FetchReaderState::Draft14 => fo14::read_object_meta(buf)
2508                .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft14, meta)),
2509            #[cfg(feature = "draft15")]
2510            FetchReaderState::Draft15(inner) => fo15::read_object_frame(inner, buf),
2511            #[cfg(feature = "draft16")]
2512            FetchReaderState::Draft16(inner) => fo16::read_object_frame(inner, buf),
2513            #[cfg(feature = "draft17")]
2514            FetchReaderState::Draft17(inner) => fo17::read_object_frame(inner, buf),
2515            #[cfg(feature = "draft18")]
2516            FetchReaderState::Draft18(inner) => fo18::read_object_frame(inner, buf),
2517            #[cfg(feature = "draft19")]
2518            FetchReaderState::Draft19(inner) => fo19::read_object_frame(inner, buf),
2519            #[cfg(feature = "draft20")]
2520            FetchReaderState::Draft20(inner) => fo20::read_object_frame(inner, buf),
2521            #[allow(unreachable_patterns)]
2522            _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2523        }
2524    }
2525
2526    /// Decode the next fetch frame's framing without copying its payload.
2527    ///
2528    /// Advances `buf` past the whole frame exactly as
2529    /// [`read_object`](Self::read_object) does, and advances the same reader
2530    /// state, so the two are interchangeable on one stream.
2531    #[allow(unused_variables, unreachable_code)]
2532    pub fn read_object_meta(
2533        &mut self,
2534        buf: &mut impl Buf,
2535    ) -> Result<AnyFetchObjectMeta, CodecError> {
2536        match &mut self.state {
2537            #[cfg(feature = "draft07")]
2538            FetchReaderState::Draft07 => fo07::read_object_meta(buf),
2539            #[cfg(feature = "draft08")]
2540            FetchReaderState::Draft08 => fo08::read_object_meta(buf),
2541            #[cfg(feature = "draft09")]
2542            FetchReaderState::Draft09 => fo09::read_object_meta(buf),
2543            #[cfg(feature = "draft10")]
2544            FetchReaderState::Draft10 => fo10::read_object_meta(buf),
2545            #[cfg(feature = "draft11")]
2546            FetchReaderState::Draft11 => fo11::read_object_meta(buf),
2547            #[cfg(feature = "draft12")]
2548            FetchReaderState::Draft12 => fo12::read_object_meta(buf),
2549            #[cfg(feature = "draft13")]
2550            FetchReaderState::Draft13 => fo13::read_object_meta(buf),
2551            #[cfg(feature = "draft14")]
2552            FetchReaderState::Draft14 => fo14::read_object_meta(buf),
2553            #[cfg(feature = "draft15")]
2554            FetchReaderState::Draft15(inner) => fo15::read_object_meta(inner, buf),
2555            #[cfg(feature = "draft16")]
2556            FetchReaderState::Draft16(inner) => fo16::read_object_meta(inner, buf),
2557            #[cfg(feature = "draft17")]
2558            FetchReaderState::Draft17(inner) => fo17::read_object_meta(inner, buf),
2559            #[cfg(feature = "draft18")]
2560            FetchReaderState::Draft18(inner) => fo18::read_object_meta(inner, buf),
2561            #[cfg(feature = "draft19")]
2562            FetchReaderState::Draft19(inner) => fo19::read_object_meta(inner, buf),
2563            #[cfg(feature = "draft20")]
2564            FetchReaderState::Draft20(inner) => fo20::read_object_meta(inner, buf),
2565            #[allow(unreachable_patterns)]
2566            _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2567        }
2568    }
2569}
2570
2571// ── Carrying a fetch frame from a reader to a writer ────────
2572
2573/// Per-draft capture of the shape one fetch frame arrived in.
2574///
2575/// Drafts 07-14 keep nothing: every field of a fetch object is on their wire
2576/// outright, so the bytes say the same thing whatever precedes them. Drafts
2577/// 15-20 keep the frame's own header, and draft-16 the resolved Location
2578/// beside it, because those two are exactly what each draft's
2579/// `FetchObjectWriter` is handed.
2580#[derive(Debug, Clone)]
2581enum FetchFrameShape {
2582    /// A frame whose fields are all absolute.
2583    #[cfg(any(
2584        feature = "draft07",
2585        feature = "draft08",
2586        feature = "draft09",
2587        feature = "draft10",
2588        feature = "draft11",
2589        feature = "draft12",
2590        feature = "draft13",
2591        feature = "draft14"
2592    ))]
2593    Absolute,
2594    #[cfg(feature = "draft15")]
2595    Draft15(crate::draft15::data_stream::FetchObjectHeader),
2596    #[cfg(feature = "draft16")]
2597    Draft16(
2598        crate::draft16::data_stream::FetchObjectHeader,
2599        crate::draft16::data_stream::FetchObjectLocation,
2600    ),
2601    #[cfg(feature = "draft17")]
2602    Draft17(crate::draft17::data_stream::FetchObject),
2603    #[cfg(feature = "draft18")]
2604    Draft18(crate::draft18::data_stream::FetchObject),
2605    #[cfg(feature = "draft19")]
2606    Draft19(crate::draft19::data_stream::FetchObject),
2607    #[cfg(feature = "draft20")]
2608    Draft20(crate::draft20::data_stream::FetchObject),
2609}
2610
2611/// One fetch frame, in the form re-encoding it takes.
2612///
2613/// Produced by [`AnyFetchObjectReader::read_object_frame`] and consumed by
2614/// [`AnyFetchObjectWriter::reemit_object`]. It is one value rather than two
2615/// because a frame's resolved identity and the shape it arrived in are only
2616/// meaningful together: the first says what the frame *is*, the second is the
2617/// encoding a writer keeps wherever it still says the same thing, and that is
2618/// what reproduces an untouched stream byte for byte.
2619#[derive(Debug, Clone)]
2620pub struct AnyFetchFrame {
2621    /// The framing, exactly as [`AnyFetchObjectReader::read_object_meta`]
2622    /// reports it.
2623    pub meta: AnyFetchObjectMeta,
2624    draft: DraftVersion,
2625    shape: FetchFrameShape,
2626}
2627
2628impl AnyFetchFrame {
2629    /// The draft whose stream this frame was read off.
2630    ///
2631    /// A writer refuses a frame from any other draft rather than re-encoding
2632    /// it: the two would agree on the resolved Location and disagree on
2633    /// everything the flags mean.
2634    #[must_use]
2635    pub fn draft(&self) -> DraftVersion {
2636        self.draft
2637    }
2638
2639    /// A frame from one of the drafts that keeps nothing.
2640    #[cfg(any(
2641        feature = "draft07",
2642        feature = "draft08",
2643        feature = "draft09",
2644        feature = "draft10",
2645        feature = "draft11",
2646        feature = "draft12",
2647        feature = "draft13",
2648        feature = "draft14"
2649    ))]
2650    fn absolute(draft: DraftVersion, meta: AnyFetchObjectMeta) -> Self {
2651        Self { meta, draft, shape: FetchFrameShape::Absolute }
2652    }
2653}
2654
2655// ── Fetch object writer ─────────────────────────────────────
2656
2657/// What [`AnyFetchObjectWriter::reemit_object`] had to do.
2658///
2659/// The counterpart of [`Reemit`], and deliberately not the same type. A
2660/// subgroup object's fix-up rewrites one leading varint and always writes the
2661/// whole object out; a fetch frame's is a re-encode of the whole header, and
2662/// the case worth having a shape for is the one where no re-encode is owed and
2663/// nothing is written at all.
2664#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2665pub enum FetchReemit {
2666    /// The frame's own bytes still encode it against the frame now in front of
2667    /// it, so **nothing was written** and the caller forwards `raw` untouched.
2668    ///
2669    /// This is every frame on a stream nothing was removed from, which is why
2670    /// it writes nothing: a relay that copied each frame through an output
2671    /// buffer to discover it would copy every payload it forwards.
2672    Unchanged,
2673    /// New framing was needed. The whole frame — the new framing followed by
2674    /// every byte of `raw` behind the old — was written to `out`, and the
2675    /// caller forwards that instead of `raw`.
2676    Reframed {
2677        /// Bytes the framing occupied in `raw`.
2678        framing_bytes_before: usize,
2679        /// Bytes it occupies in the output.
2680        framing_bytes_after: usize,
2681    },
2682}
2683
2684/// Per-draft writer state, mirroring [`FetchReaderState`].
2685#[derive(Debug, Clone)]
2686enum FetchWriterState {
2687    /// Drafts 07-14, which write every field of a fetch object outright and
2688    /// have nothing to write one *against*. The draft is carried so that a
2689    /// frame from another one is still refused.
2690    #[cfg(any(
2691        feature = "draft07",
2692        feature = "draft08",
2693        feature = "draft09",
2694        feature = "draft10",
2695        feature = "draft11",
2696        feature = "draft12",
2697        feature = "draft13",
2698        feature = "draft14"
2699    ))]
2700    Absolute(DraftVersion),
2701    #[cfg(feature = "draft15")]
2702    Draft15(crate::draft15::data_stream::FetchObjectWriter),
2703    #[cfg(feature = "draft16")]
2704    Draft16(crate::draft16::data_stream::FetchObjectWriter),
2705    #[cfg(feature = "draft17")]
2706    Draft17(crate::draft17::data_stream::FetchObjectWriter),
2707    #[cfg(feature = "draft18")]
2708    Draft18(crate::draft18::data_stream::FetchObjectWriter),
2709    #[cfg(feature = "draft19")]
2710    Draft19(crate::draft19::data_stream::FetchObjectWriter),
2711    #[cfg(feature = "draft20")]
2712    Draft20(crate::draft20::data_stream::FetchObjectWriter),
2713}
2714
2715/// Re-emitter for the frames of a fetch data stream, for any enabled draft.
2716///
2717/// The inverse of [`AnyFetchObjectReader`], and it exists for one caller: a
2718/// relay reading one fetch stream and writing another from the same frames,
2719/// having removed some of them. Removing a frame changes what the frames
2720/// behind it are encoded *against*, and on drafts 15-20 nearly every field of
2721/// a fetch object is defined against the frame before it — draft-17
2722/// Section 10.4.4.1, Table 7: "Object ID is the prior Object's ID plus one" —
2723/// so a survivor following a removed run cannot keep its original bytes.
2724///
2725/// # How it is driven
2726///
2727/// One writer belongs to one stream, and [`reemit_object`](Self::reemit_object)
2728/// is called **for every frame the caller forwards**, in wire order, whether or
2729/// not anything has been removed yet. That call is what moves the writer, so a
2730/// forwarded frame it never saw leaves it a frame behind and re-encodes the
2731/// next survivor against the wrong predecessor. A frame the caller *elides* is
2732/// the one it is not called for — that is the whole of eliding.
2733///
2734/// # What it costs
2735///
2736/// Nothing on drafts 07-14, and on drafts 15-20 one header re-derivation per
2737/// frame, which allocates only when the answer differs from the bytes that
2738/// arrived. A stream with nothing removed from it therefore forwards every
2739/// frame's own bytes and copies no payload.
2740///
2741/// # Drafts 18 and 19 need the Group Order
2742///
2743/// Their Group ID is a difference whose sign the fetch's Group Order decides,
2744/// exactly as for [`AnyFetchObjectReader`], and it is settled on the control
2745/// plane rather than on the data stream. [`new`](Self::new) takes it for that
2746/// reason: the wrong order re-encodes without error onto groups walking the
2747/// wrong way.
2748#[derive(Debug, Clone)]
2749pub struct AnyFetchObjectWriter {
2750    state: FetchWriterState,
2751}
2752
2753impl AnyFetchObjectWriter {
2754    /// Create a writer for a stream with the given header, whose groups are
2755    /// written in `group_order`.
2756    ///
2757    /// See the type's own documentation for what the order is for and why it
2758    /// cannot be read off the stream. On drafts 07-17 the argument is ignored.
2759    ///
2760    /// Returns [`CodecError::UnsupportedDraft`] for drafts not compiled in.
2761    #[allow(unused_variables, unreachable_code)]
2762    pub fn new(
2763        header: &AnyFetchHeader,
2764        group_order: AnyFetchGroupOrder,
2765    ) -> Result<Self, CodecError> {
2766        let state = match header {
2767            #[cfg(feature = "draft07")]
2768            AnyFetchHeader::Draft07(_) => FetchWriterState::Absolute(DraftVersion::Draft07),
2769            #[cfg(feature = "draft08")]
2770            AnyFetchHeader::Draft08(_) => FetchWriterState::Absolute(DraftVersion::Draft08),
2771            #[cfg(feature = "draft09")]
2772            AnyFetchHeader::Draft09(_) => FetchWriterState::Absolute(DraftVersion::Draft09),
2773            #[cfg(feature = "draft10")]
2774            AnyFetchHeader::Draft10(_) => FetchWriterState::Absolute(DraftVersion::Draft10),
2775            #[cfg(feature = "draft11")]
2776            AnyFetchHeader::Draft11(_) => FetchWriterState::Absolute(DraftVersion::Draft11),
2777            #[cfg(feature = "draft12")]
2778            AnyFetchHeader::Draft12(_) => FetchWriterState::Absolute(DraftVersion::Draft12),
2779            #[cfg(feature = "draft13")]
2780            AnyFetchHeader::Draft13(_) => FetchWriterState::Absolute(DraftVersion::Draft13),
2781            #[cfg(feature = "draft14")]
2782            AnyFetchHeader::Draft14(_) => FetchWriterState::Absolute(DraftVersion::Draft14),
2783            #[cfg(feature = "draft15")]
2784            AnyFetchHeader::Draft15(_) => {
2785                FetchWriterState::Draft15(crate::draft15::data_stream::FetchObjectWriter::new())
2786            }
2787            #[cfg(feature = "draft16")]
2788            AnyFetchHeader::Draft16(_) => {
2789                FetchWriterState::Draft16(crate::draft16::data_stream::FetchObjectWriter::new())
2790            }
2791            #[cfg(feature = "draft17")]
2792            AnyFetchHeader::Draft17(_) => {
2793                FetchWriterState::Draft17(crate::draft17::data_stream::FetchObjectWriter::new())
2794            }
2795            #[cfg(feature = "draft18")]
2796            AnyFetchHeader::Draft18(_) => FetchWriterState::Draft18(
2797                crate::draft18::data_stream::FetchObjectWriter::new(match group_order {
2798                    AnyFetchGroupOrder::Ascending => {
2799                        crate::draft18::data_stream::GroupOrder::Ascending
2800                    }
2801                    AnyFetchGroupOrder::Descending => {
2802                        crate::draft18::data_stream::GroupOrder::Descending
2803                    }
2804                }),
2805            ),
2806            #[cfg(feature = "draft19")]
2807            AnyFetchHeader::Draft19(_) => FetchWriterState::Draft19(
2808                crate::draft19::data_stream::FetchObjectWriter::new(match group_order {
2809                    AnyFetchGroupOrder::Ascending => {
2810                        crate::draft19::data_stream::GroupOrder::Ascending
2811                    }
2812                    AnyFetchGroupOrder::Descending => {
2813                        crate::draft19::data_stream::GroupOrder::Descending
2814                    }
2815                }),
2816            ),
2817            #[cfg(feature = "draft20")]
2818            AnyFetchHeader::Draft20(_) => FetchWriterState::Draft20(
2819                crate::draft20::data_stream::FetchObjectWriter::new(match group_order {
2820                    AnyFetchGroupOrder::Ascending => {
2821                        crate::draft20::data_stream::GroupOrder::Ascending
2822                    }
2823                    AnyFetchGroupOrder::Descending => {
2824                        crate::draft20::data_stream::GroupOrder::Descending
2825                    }
2826                }),
2827            ),
2828            #[allow(unreachable_patterns)]
2829            _ => {
2830                return Err(CodecError::UnsupportedDraft(format!(
2831                    "draft {:?} not enabled via feature flag",
2832                    header.draft()
2833                )));
2834            }
2835        };
2836        Ok(Self { state })
2837    }
2838
2839    /// The draft this writer encodes.
2840    #[must_use]
2841    #[allow(unreachable_code)]
2842    pub fn draft(&self) -> DraftVersion {
2843        match &self.state {
2844            #[cfg(any(
2845                feature = "draft07",
2846                feature = "draft08",
2847                feature = "draft09",
2848                feature = "draft10",
2849                feature = "draft11",
2850                feature = "draft12",
2851                feature = "draft13",
2852                feature = "draft14"
2853            ))]
2854            FetchWriterState::Absolute(draft) => *draft,
2855            #[cfg(feature = "draft15")]
2856            FetchWriterState::Draft15(_) => DraftVersion::Draft15,
2857            #[cfg(feature = "draft16")]
2858            FetchWriterState::Draft16(_) => DraftVersion::Draft16,
2859            #[cfg(feature = "draft17")]
2860            FetchWriterState::Draft17(_) => DraftVersion::Draft17,
2861            #[cfg(feature = "draft18")]
2862            FetchWriterState::Draft18(_) => DraftVersion::Draft18,
2863            #[cfg(feature = "draft19")]
2864            FetchWriterState::Draft19(_) => DraftVersion::Draft19,
2865            #[cfg(feature = "draft20")]
2866            FetchWriterState::Draft20(_) => DraftVersion::Draft20,
2867            #[allow(unreachable_patterns)]
2868            _ => unreachable!("AnyFetchObjectWriter has no enabled variants"),
2869        }
2870    }
2871
2872    /// Re-emit one forwarded fetch frame, re-encoding its framing against the
2873    /// frames actually forwarded before it, and advance.
2874    ///
2875    /// `frame` came from [`AnyFetchObjectReader::read_object_frame`] on the
2876    /// stream being read; `raw` is that frame's wire bytes. The return value
2877    /// says which bytes to forward, and the two answers are not symmetric:
2878    /// [`FetchReemit::Unchanged`] writes nothing and means `raw` is still
2879    /// correct, while [`FetchReemit::Reframed`] has written the whole frame to
2880    /// `out` and `raw` must not also be forwarded.
2881    ///
2882    /// # `raw` need not be a complete frame
2883    ///
2884    /// **Any prefix is legal provided the whole framing is present** — the
2885    /// framing being `meta.wire_len - meta.payload_length` bytes, which is a
2886    /// number the frame already carries. Everything behind it is copied
2887    /// byte-for-byte, however many bytes there are, and no length validation is
2888    /// performed. That is what lets a caller fix up the first chunk of a frame
2889    /// too large to buffer, where only the first chunk carries the framing at
2890    /// all.
2891    ///
2892    /// # Errors
2893    ///
2894    /// [`CodecError::UnsupportedDraft`] when `frame` was read off another
2895    /// draft's stream.
2896    ///
2897    /// [`CodecError::InvalidField`] when `raw` is shorter than the framing the
2898    /// frame declares, and when the frame has no encoding against the
2899    /// predecessor now in front of it — a Group ID that moves against the
2900    /// Group Order, an Object ID that does not advance, and the arithmetic
2901    /// overflows. The writer is left where it was in that case, so a caller
2902    /// that gives up on one frame and carries on is not also one frame out.
2903    #[allow(unused_variables)]
2904    pub fn reemit_object(
2905        &mut self,
2906        frame: &AnyFetchFrame,
2907        raw: &[u8],
2908        out: &mut impl BufMut,
2909    ) -> Result<FetchReemit, CodecError> {
2910        if frame.draft != self.draft() {
2911            return Err(CodecError::UnsupportedDraft(format!(
2912                "a draft {:?} fetch frame cannot be written onto a draft {:?} stream",
2913                frame.draft,
2914                self.draft()
2915            )));
2916        }
2917
2918        let framing_len = frame.meta.wire_len.saturating_sub(frame.meta.payload_length);
2919        let framing_len = usize::try_from(framing_len).map_err(|_| CodecError::InvalidField)?;
2920        if framing_len > raw.len() {
2921            return Err(CodecError::InvalidField);
2922        }
2923        let (framing, rest) = raw.split_at(framing_len);
2924
2925        match (&mut self.state, &frame.shape) {
2926            // Nothing on these drafts' wire is written against anything, so
2927            // the frame's own bytes are correct wherever it lands.
2928            #[cfg(any(
2929                feature = "draft07",
2930                feature = "draft08",
2931                feature = "draft09",
2932                feature = "draft10",
2933                feature = "draft11",
2934                feature = "draft12",
2935                feature = "draft13",
2936                feature = "draft14"
2937            ))]
2938            (FetchWriterState::Absolute(_), FetchFrameShape::Absolute) => {
2939                Ok(FetchReemit::Unchanged)
2940            }
2941            #[cfg(feature = "draft15")]
2942            (FetchWriterState::Draft15(writer), FetchFrameShape::Draft15(original)) => {
2943                let reframed = writer.header_for(original)?;
2944                if reframed == *original {
2945                    writer.advance(original);
2946                    return Ok(FetchReemit::Unchanged);
2947                }
2948                let mut encoded = Vec::with_capacity(framing.len() + 16);
2949                reframed.encode(&mut encoded)?;
2950                writer.advance(&reframed);
2951                Ok(put_reframed(&encoded, framing.len(), rest, out))
2952            }
2953            #[cfg(feature = "draft16")]
2954            (FetchWriterState::Draft16(writer), FetchFrameShape::Draft16(original, location)) => {
2955                let reframed = writer.header_for(original, location)?;
2956                if reframed == *original {
2957                    writer.advance(original, location);
2958                    return Ok(FetchReemit::Unchanged);
2959                }
2960                let mut encoded = Vec::with_capacity(framing.len() + 16);
2961                reframed.encode(&mut encoded)?;
2962                writer.advance(&reframed, location);
2963                Ok(put_reframed(&encoded, framing.len(), rest, out))
2964            }
2965            #[cfg(feature = "draft17")]
2966            (FetchWriterState::Draft17(writer), FetchFrameShape::Draft17(original)) => {
2967                let reframed = writer.header_for(original)?;
2968                if reframed == original.header {
2969                    writer.advance(original);
2970                    return Ok(FetchReemit::Unchanged);
2971                }
2972                let mut encoded = Vec::with_capacity(framing.len() + 16);
2973                reframed.encode(&mut encoded)?;
2974                writer.advance(original);
2975                Ok(put_reframed(&encoded, framing.len(), rest, out))
2976            }
2977            #[cfg(feature = "draft18")]
2978            (FetchWriterState::Draft18(writer), FetchFrameShape::Draft18(original)) => {
2979                let reframed = writer.header_for(original)?;
2980                if reframed == original.header {
2981                    writer.advance(original);
2982                    return Ok(FetchReemit::Unchanged);
2983                }
2984                let mut encoded = Vec::with_capacity(framing.len() + 16);
2985                reframed.encode(&mut encoded)?;
2986                writer.advance(original);
2987                Ok(put_reframed(&encoded, framing.len(), rest, out))
2988            }
2989            #[cfg(feature = "draft19")]
2990            (FetchWriterState::Draft19(writer), FetchFrameShape::Draft19(original)) => {
2991                let reframed = writer.header_for(original)?;
2992                if reframed == original.header {
2993                    writer.advance(original);
2994                    return Ok(FetchReemit::Unchanged);
2995                }
2996                let mut encoded = Vec::with_capacity(framing.len() + 16);
2997                reframed.encode(&mut encoded)?;
2998                writer.advance(original);
2999                Ok(put_reframed(&encoded, framing.len(), rest, out))
3000            }
3001            #[cfg(feature = "draft20")]
3002            (FetchWriterState::Draft20(writer), FetchFrameShape::Draft20(original)) => {
3003                let reframed = writer.header_for(original)?;
3004                if reframed == original.header {
3005                    writer.advance(original);
3006                    return Ok(FetchReemit::Unchanged);
3007                }
3008                let mut encoded = Vec::with_capacity(framing.len() + 16);
3009                reframed.encode(&mut encoded)?;
3010                writer.advance(original);
3011                Ok(put_reframed(&encoded, framing.len(), rest, out))
3012            }
3013            // Unreachable: the drafts were compared before this match, and one
3014            // draft has one state and one shape.
3015            #[allow(unreachable_patterns)]
3016            _ => Err(CodecError::UnsupportedDraft(format!(
3017                "no fetch writer for draft {:?}",
3018                frame.draft
3019            ))),
3020        }
3021    }
3022}
3023
3024/// Write a re-encoded frame out: the new framing, then every byte that stood
3025/// behind the old one.
3026#[cfg(any(
3027    feature = "draft15",
3028    feature = "draft16",
3029    feature = "draft17",
3030    feature = "draft18",
3031    feature = "draft19",
3032    feature = "draft20"
3033))]
3034fn put_reframed(
3035    encoded: &[u8],
3036    framing_bytes_before: usize,
3037    rest: &[u8],
3038    out: &mut impl BufMut,
3039) -> FetchReemit {
3040    out.put_slice(encoded);
3041    out.put_slice(rest);
3042    FetchReemit::Reframed { framing_bytes_before, framing_bytes_after: encoded.len() }
3043}
3044
3045// ── Header helpers for the stream-type-gated drafts ─────────
3046
3047/// Generates the subgroup stream-type check drafts 11-13 share: the header's
3048/// stream type must be a subgroup type, and it decides whether objects carry
3049/// an extension block.
3050macro_rules! subgroup_extensions_fn {
3051    ($name:ident, $feat:literal, $draft:ident) => {
3052        #[cfg(feature = $feat)]
3053        fn $name(header: &crate::$draft::data_stream::SubgroupHeader) -> Result<bool, CodecError> {
3054            if !header.stream_type.is_subgroup() {
3055                return Err(CodecError::InvalidField);
3056            }
3057            Ok(header.stream_type.has_extensions())
3058        }
3059    };
3060}
3061
3062subgroup_extensions_fn!(subgroup_extensions_11, "draft11", draft11);
3063subgroup_extensions_fn!(subgroup_extensions_12, "draft12", draft12);
3064subgroup_extensions_fn!(subgroup_extensions_13, "draft13", draft13);