Skip to main content

weida_protocol/
header.rs

1//! CBOR frame headers.
2//!
3//! The codecs are written by hand rather than derived. Three reasons:
4//!
5//! * the strictness rules (duplicate-key rejection, non-uint key rejection,
6//!   per-field string caps, depth-limited skipping) are protocol requirements,
7//!   not serialization defaults;
8//! * the exact byte layout is normative — see the golden vectors in
9//!   `docs/PROTOCOL.md` §8 — and a derive macro's field ordering is not a
10//!   contract we control;
11//! * every decoder here is a hostile-input boundary, so its allocation
12//!   behaviour has to be readable.
13//!
14//! All decode failures are protocol violations that close the connection.
15//!
16//! **No field of a DATA header is required by the decoder.** A header no longer
17//! carries the stream's role, so the decoder cannot know which fields the
18//! context demands; `endpoint`-on-initiating-streams is enforced by the
19//! transport's dispatch, which does know (`docs/PROTOCOL.md` §6.2).
20
21use std::convert::Infallible;
22
23use minicbor::data::Type;
24use minicbor::{Decoder, Encoder};
25use weida_core::Error;
26
27use crate::varint::{VarintError, decode_varint, encode_varint};
28
29/// Decoder limits. The string caps are normative
30/// (`docs/PROTOCOL.md` §6); the list and depth caps are defensive
31/// implementation limits documented in the same section.
32pub mod limits {
33    /// Cap for the DATA `endpoint` field.
34    pub const MAX_ENDPOINT_BYTES: usize = 512;
35    /// Cap for the DATA `content_type` field.
36    pub const MAX_CONTENT_TYPE_BYTES: usize = 256;
37    /// Cap for the DATA `traceparent` field.
38    pub const MAX_TRACEPARENT_BYTES: usize = 128;
39    /// Cap for the DATA `tracestate` field.
40    pub const MAX_TRACESTATE_BYTES: usize = 512;
41    /// Cap for the DATA `topic` field.
42    pub const MAX_TOPIC_BYTES: usize = 256;
43    /// Length of the DATA `producer` field: a raw 32-byte digest.
44    ///
45    /// Both a cap and an exact length. `docs/PROTOCOL.md` §6.2 defines the
46    /// value as "the raw 32-byte digest", so a longer one is a framing
47    /// violation and a shorter one names nothing this specification defines.
48    pub const PRODUCER_BYTES: usize = 32;
49    /// Cap for the SUBSCRIBE/UNSUBSCRIBE `filter` field.
50    pub const MAX_FILTER_BYTES: usize = 256;
51    /// Cap for the ERROR `message` field.
52    pub const MAX_MESSAGE_BYTES: usize = 1024;
53    /// Cap on the number of items in a HELLO list field.
54    ///
55    /// Without it, a hostile peer could pin `max_concurrent_uni_streams`
56    /// worth of large `Vec<u64>`s by opening many HELLO streams.
57    pub const MAX_LIST_ITEMS: usize = 64;
58    /// Cap on the number of levels a DATA header may order a report for
59    /// (`docs/PROTOCOL.md` §6.2, key `10`).
60    ///
61    /// A report order is a remote-controlled list, so it needs a cap for the
62    /// same reason [`MAX_LIST_ITEMS`] exists; 16 is more levels than the level
63    /// space defines below the application floor, so it constrains nothing a
64    /// sender legitimately wants.
65    pub const MAX_REPORT_LEVELS: usize = 16;
66    /// Nesting depth allowed when skipping an unknown field.
67    pub const MAX_SKIP_DEPTH: usize = 8;
68    /// Highest segment layer (`docs/PROTOCOL.md` §6.2 key `14`, §6.4 key
69    /// `3`).
70    pub const MAX_LAYER: u8 = 15;
71    /// Largest REPORT record map (`docs/PROTOCOL.md` §6.10).
72    pub const MAX_PATH_RECORD_BYTES: usize = 256;
73}
74
75/// HELLO keys.
76mod hello_key {
77    pub const VERSIONS: u64 = 0;
78    pub const MAX_HEADER_BYTES: u64 = 1;
79    pub const MAX_TRANSFERS: u64 = 2;
80    pub const CAPABILITIES: u64 = 3;
81    pub const REQUIRED_CAPABILITIES: u64 = 4;
82    pub const GUARANTEES_OFFERED: u64 = 5;
83    pub const GUARANTEES_REQUIRED: u64 = 6;
84}
85
86/// DATA keys.
87mod data_key {
88    pub const ENDPOINT: u64 = 0;
89    pub const CONTENT_LEN: u64 = 1;
90    pub const CONTENT_TYPE: u64 = 2;
91    pub const TRACEPARENT: u64 = 3;
92    pub const TRACESTATE: u64 = 4;
93    pub const TOPIC: u64 = 5;
94    pub const SEQUENCE: u64 = 6;
95    pub const PRODUCER: u64 = 7;
96    pub const ACHIEVED: u64 = 8;
97    pub const REPORT_ID: u64 = 9;
98    pub const REPORT: u64 = 10;
99    pub const REPORT_MODE: u64 = 11;
100    // Key 12 (`delivery_attempt`) is reserved in docs/PROTOCOL.md §6.2 and
101    // coded by the slice that writes it.
102    pub const SEGMENT: u64 = 13;
103    pub const LAYER: u64 = 14;
104}
105
106/// ERROR keys.
107mod error_key {
108    pub const CODE: u64 = 0;
109    pub const MESSAGE: u64 = 1;
110}
111
112/// SUBSCRIBE and UNSUBSCRIBE keys.
113mod subscription_key {
114    pub const ENDPOINT: u64 = 0;
115    pub const FILTER: u64 = 1;
116    pub const MAX_AGE_MS: u64 = 2;
117    pub const MAX_LAYER: u64 = 3;
118}
119
120/// CREDIT keys.
121mod credit_key {
122    pub const ENDPOINT: u64 = 0;
123    pub const FILTER: u64 = 1;
124    pub const LIMIT: u64 = 2;
125}
126
127/// CURSOR head-frame keys.
128mod cursor_key {
129    pub const REPORT_ID: u64 = 0;
130}
131
132/// FLOW keys (`docs/PROTOCOL.md` §6.8).
133mod flow_key {
134    pub const ENDPOINT: u64 = 0;
135    pub const FLOW: u64 = 1;
136    pub const CONTENT_TYPE: u64 = 2;
137    pub const TRACEPARENT: u64 = 3;
138    pub const TRACESTATE: u64 = 4;
139    pub const TOPIC: u64 = 5;
140}
141
142/// The topic filter grammar of `docs/PROTOCOL.md` §6.4.
143///
144/// A topic and a filter are byte strings split on [`filter::SEPARATOR`] into
145/// segments. The two wildcards are whole-segment tokens, and everything else
146/// is literal: there is no escape character, no normalization and no case
147/// folding ([decisions/0007](../../../docs/decisions/0007-topic-namespace.md)
148/// §4.2). Both halves of the grammar live here: [`filter::validate`], the rule
149/// that says which filters may exist at all — so an illegal one is refused at
150/// the codec boundary rather than reaching a matcher that would have to cope
151/// with it — and [`filter::matches`], the matcher itself. The matcher used to
152/// sit with the fan-out it served, in `weida::pubsub`, and moved when a second
153/// layer needed it: an L2 queue selects a consumer with the same grammar
154/// (B-202), and a second implementation of a wildcard language is exactly the
155/// kind of drift one definition exists to prevent.
156pub mod filter {
157    use super::HeaderError;
158
159    /// Segment separator: `.`, one byte.
160    pub const SEPARATOR: char = '.';
161    /// Matches exactly one whole segment.
162    pub const ONE_SEGMENT: &str = "*";
163    /// Matches zero or more trailing segments; legal only as the last segment.
164    pub const REST: &str = "#";
165
166    /// Does `topic` match `filter`?
167    ///
168    /// The segmented grammar of `docs/PROTOCOL.md` §6.4: segments split on `.`,
169    /// `*` for exactly one whole segment, a trailing `#` for zero or more, every
170    /// other byte literal, and the empty filter matching everything.
171    ///
172    /// The objection this function used to carry was that treating `*` as a
173    /// wildcard "would make topics with a literal `*` unaddressable and would put
174    /// a matching language in the hot path". Both halves were true and both are
175    /// accepted deliberately
176    /// ([decisions/0007](../../../../docs/decisions/0007-topic-namespace.md) §4.6): a
177    /// filter can no longer select a segment containing `.`, `*` or `#`
178    /// literally — there is no escape character, and no sheet reports a use for
179    /// one — while a byte prefix could not express a boundary at all, so
180    /// `sensors.temp` also selected `sensors.temperature`. The hot-path half is
181    /// answered by the shape rather than by the choice: `#` is legal only as the
182    /// final segment, so this is one left-to-right walk with no backtracking, no
183    /// allocation and work bounded by the 256 B filter cap.
184    ///
185    /// A `topic` is never a pattern: `*` and `#` in a published topic are literal
186    /// bytes here, exactly like any other.
187    pub fn matches(topic: &str, filter: &str) -> bool {
188        if filter.is_empty() {
189            return true;
190        }
191        let mut topic_segments = topic.split(SEPARATOR);
192        let mut filter_segments = filter.split(SEPARATOR);
193        loop {
194            let Some(pattern) = filter_segments.next() else {
195                // The filter is spent: it matches only if the topic is too.
196                return topic_segments.next().is_none();
197            };
198            // Only ever the final segment — `filter::validate` rejects anything
199            // else at the codec boundary — so everything left over matches.
200            if pattern == REST {
201                return true;
202            }
203            let Some(segment) = topic_segments.next() else {
204                return false;
205            };
206            if pattern != ONE_SEGMENT && pattern != segment {
207                return false;
208            }
209        }
210    }
211
212    /// Checks `filter` against the grammar.
213    ///
214    /// The empty filter is legal and matches every topic. One pass, no
215    /// allocation.
216    pub fn validate(filter: &str) -> Result<(), HeaderError> {
217        let mut segments = filter.split(SEPARATOR).peekable();
218        while let Some(segment) = segments.next() {
219            let is_last = segments.peek().is_none();
220            if segment.contains(ONE_SEGMENT) && segment != ONE_SEGMENT {
221                return Err(HeaderError::InvalidFilter(
222                    "`*` must occupy a whole segment",
223                ));
224            }
225            if segment.contains(REST) {
226                if segment != REST {
227                    return Err(HeaderError::InvalidFilter(
228                        "`#` must occupy a whole segment",
229                    ));
230                }
231                if !is_last {
232                    return Err(HeaderError::InvalidFilter("`#` must be the final segment"));
233                }
234            }
235        }
236        Ok(())
237    }
238}
239
240/// Guarantee set keys (`docs/PROTOCOL.md` §6.5).
241mod guarantee_key {
242    pub const DELIVERY: u64 = 0;
243    pub const ACKNOWLEDGEMENT: u64 = 1;
244    pub const DURABILITY: u64 = 2;
245    pub const REPLICAS: u64 = 3;
246    pub const ORDERING: u64 = 4;
247    pub const DEDUPLICATION: u64 = 5;
248    pub const DEDUP_WINDOW_MS: u64 = 6;
249    pub const BACKPRESSURE: u64 = 7;
250    pub const PRODUCER_NAMING: u64 = 8;
251    pub const CONTROL_ISOLATED: u64 = 9;
252}
253
254/// Declares an enum whose wire form is a small `uint`, with the `core` level
255/// first so that `Default` and "absent means core" agree by construction.
256///
257/// The derived `Ord` ranks by declaration order while `to_wire` reads
258/// explicit literals, and the ladder comparisons rest on the two agreeing —
259/// `GuaranteeSet::intersect` picks the weaker level with `.min()`,
260/// `GuaranteeSet::reaches` compares with `<`, and `CursorLevel`'s derived
261/// order inherits the same coincidence — so the macro asserts the agreement
262/// at compile time rather than letting a variant inserted mid-block with a
263/// higher literal silently rank a stronger guarantee below a weaker one.
264macro_rules! wire_enum {
265    ($(#[$meta:meta])* $name:ident { $($(#[$vmeta:meta])* $variant:ident = $value:literal),+ $(,)? }) => {
266        $(#[$meta])*
267        #[derive(Clone, Copy, Debug, Default, PartialEq, Eq, PartialOrd, Ord, Hash)]
268        pub enum $name {
269            $($(#[$vmeta])* $variant,)+
270        }
271
272        impl $name {
273            /// The wire value of `docs/PROTOCOL.md` §6.5.
274            pub fn to_wire(self) -> u64 {
275                match self {
276                    $($name::$variant => $value,)+
277                }
278            }
279
280            /// The level a wire value names, or `None` if the value is not
281            /// one this version defines.
282            pub fn from_wire(value: u64) -> Option<$name> {
283                match value {
284                    $($value => Some($name::$variant),)+
285                    _ => None,
286                }
287            }
288        }
289
290        // The invariant the ladder comparisons depend on. A build failure is
291        // the only acceptable outcome: at run time the mis-ranking is
292        // invisible — every value still encodes and decodes — and shows up
293        // only as a negotiated guarantee weaker than the one reported.
294        const _: () = {
295            let values = [$($value as u64),+];
296            let mut i = 1;
297            while i < values.len() {
298                assert!(
299                    values[i - 1] < values[i],
300                    concat!(
301                        stringify!($name),
302                        ": wire values must ascend with declaration order, ",
303                        "because the derived Ord is the ladder"
304                    )
305                );
306                i += 1;
307            }
308        };
309    };
310}
311
312wire_enum! {
313    /// Delivery dimension ([`GUARANTEES.md`] §3). A ladder: later is stronger.
314    ///
315    /// [`GUARANTEES.md`]: https://git.doodleshnookie.net/tuco86/weida/blob/main/docs/GUARANTEES.md
316    Delivery {
317        /// v0: no retries, losses reported.
318        #[default]
319        BestEffort = 0,
320        /// Reserved.
321        AtMostOnce = 1,
322        /// Reserved.
323        AtLeastOnce = 2,
324    }
325}
326
327wire_enum! {
328    /// Acknowledgement/completion dimension. A ladder; the durability axes of
329    /// [`Durability`] and `replicas` are *not* part of it.
330    Acknowledgement {
331        /// Nothing is reported.
332        None = 0,
333        /// v0: QUIC's fin-acknowledgement.
334        #[default]
335        TransportReceipt = 1,
336        /// Reserved for the L2 broker.
337        Accepted = 2,
338        /// Reserved for the L2 broker.
339        Stored = 3,
340        /// Reserved for the L2 broker.
341        Replicated = 4,
342        /// Reserved for the L2 broker.
343        Processed = 5,
344    }
345}
346
347wire_enum! {
348    /// Persistence axis of `Stored`/`Replicated`
349    /// ([decisions/0004](../../../docs/decisions/0004-durability-levels.md) §4.1).
350    Durability {
351        /// Survives the broker process.
352        #[default]
353        Written = 0,
354        /// Survives loss of power on that node.
355        Flushed = 1,
356    }
357}
358
359wire_enum! {
360    /// Ordering dimension. A ladder: later is stronger.
361    OrderingMode {
362        /// v0.
363        #[default]
364        None = 0,
365        /// Report gaps, deliver as messages arrive.
366        PerProducerDetect = 1,
367        /// Hold messages back up to a bounded buffer.
368        PerProducerReassemble = 2,
369        /// L2 only.
370        PerKey = 3,
371        /// Reserved.
372        Total = 4,
373    }
374}
375
376wire_enum! {
377    /// Deduplication dimension. A ladder: later is stronger.
378    Deduplication {
379        /// v0.
380        #[default]
381        None = 0,
382        /// Suppressed within a time window.
383        Bounded = 1,
384        /// L2 only.
385        Durable = 2,
386    }
387}
388
389wire_enum! {
390    /// Backpressure dimension. **Not ordered**: these are behaviours, not
391    /// strengths, so two peers state the same one or fail to negotiate.
392    Backpressure {
393        /// v0 for Req/Rep and Push/Pull.
394        #[default]
395        Block = 0,
396        /// Refuse past a cap.
397        Reject = 1,
398        /// v0 for fan-out.
399        Drop = 2,
400        /// Reserved.
401        Spill = 3,
402        /// Reserved.
403        Coalesce = 4,
404    }
405}
406
407wire_enum! {
408    /// How a producer is named for the sequence field of
409    /// [decisions/0001](../../../docs/decisions/0001-sequence-field.md) §7.3.
410    /// **Not ordered**: two peers state the same one or fail.
411    ProducerNaming {
412        /// The proved connection fingerprint; the counter restarts with the
413        /// connection (v0 default,
414        /// [decisions/0008](../../../docs/decisions/0008-session-identity.md) §4.3).
415        #[default]
416        Fingerprint = 0,
417        /// A name supplied above L0, carried in DATA key `7`.
418        Stable = 1,
419    }
420}
421
422wire_enum! {
423    /// How often a reporter emits a record (`docs/PROTOCOL.md` §6.2, key
424    /// `11`).
425    ///
426    /// **Not ordered**: these are two shapes of the same report, not two
427    /// strengths. A cursor is never load-bearing, so neither mode is a
428    /// guarantee and neither is negotiated
429    /// ([decisions/0023](../../../docs/decisions/0023-completion-is-a-cursor.md)
430    /// §4.5).
431    ReportMode {
432        /// Records as the level advances, coalesced at the reporter's own
433        /// granularity.
434        #[default]
435        Progress = 0,
436        /// One record per level, at the end.
437        FinalOnly = 1,
438    }
439}
440
441/// A level a cursor can name: one weida defines, or one the application does.
442///
443/// The level space is **open**
444/// ([decisions/0023](../../../docs/decisions/0023-completion-is-a-cursor.md)
445/// §4.4): values below [`CursorLevel::APPLICATION_FLOOR`] are weida's own
446/// ladder, [`Acknowledgement`], and everything at or above it is an
447/// application stage weida carries and orders but never interprets. An
448/// undefined value *below* the floor is a protocol violation rather than an
449/// application level, because the reserved range is where a later version of
450/// this specification will put its own stages.
451#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
452pub enum CursorLevel {
453    /// A level this version of the protocol defines.
454    Known(Acknowledgement),
455    /// An application stage, at or above the floor.
456    Application(u64),
457}
458
459impl CursorLevel {
460    /// First wire value an application may name.
461    pub const APPLICATION_FLOOR: u64 = 16;
462
463    /// The wire value.
464    pub fn to_wire(self) -> u64 {
465        match self {
466            CursorLevel::Known(level) => level.to_wire(),
467            CursorLevel::Application(value) => value,
468        }
469    }
470
471    /// Interprets a wire value, or `None` if it is an undefined value in the
472    /// reserved range.
473    pub fn from_wire(value: u64) -> Option<CursorLevel> {
474        if value >= CursorLevel::APPLICATION_FLOOR {
475            Some(CursorLevel::Application(value))
476        } else {
477            Acknowledgement::from_wire(value).map(CursorLevel::Known)
478        }
479    }
480
481    /// An application stage, or `None` below the floor: the reserved range is
482    /// not an application's to name.
483    pub fn application(value: u64) -> Option<CursorLevel> {
484        (value >= CursorLevel::APPLICATION_FLOOR).then_some(CursorLevel::Application(value))
485    }
486}
487
488/// One level per guarantee dimension, as declared in HELLO keys `5` and `6`
489/// (`docs/PROTOCOL.md` §6.5).
490///
491/// [`GuaranteeSet::CORE`] is the default set and is exactly what v0 does, so
492/// an absent HELLO key, an empty map and `CORE` are the same statement
493/// ([decisions/0006](../../../docs/decisions/0006-guarantee-sets.md) §4.2).
494/// Every field is a small `Copy` value: a set costs no allocation, which is
495/// what lets it ride a header a peer controls.
496#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
497pub struct GuaranteeSet {
498    /// Delivery dimension.
499    pub delivery: Delivery,
500    /// Acknowledgement/completion dimension.
501    pub acknowledgement: Acknowledgement,
502    /// Persistence axis; legal only with `Stored` or `Replicated`.
503    pub durability: Option<Durability>,
504    /// Replica count, leader included; legal only with `Replicated`, and ≥ 2.
505    pub replicas: Option<u64>,
506    /// Ordering dimension.
507    pub ordering: OrderingMode,
508    /// Deduplication dimension.
509    pub deduplication: Deduplication,
510    /// Dedup window; required with `Bounded`, forbidden otherwise.
511    pub dedup_window_ms: Option<u64>,
512    /// Backpressure behaviour. Not ordered.
513    pub backpressure: Backpressure,
514    /// How the producer of a sequenced transfer is named. Not ordered.
515    pub producer_naming: ProducerNaming,
516    /// Whether control traffic is isolated from bulk traffic
517    /// ([decisions/0002](../../../docs/decisions/0002-control-and-bulk-separation.md)
518    /// §6.1). Ordered: `true` is strictly stronger.
519    pub control_isolated: bool,
520}
521
522impl GuaranteeSet {
523    /// The default set: what v0 offers and requires.
524    pub const CORE: GuaranteeSet = GuaranteeSet {
525        delivery: Delivery::BestEffort,
526        acknowledgement: Acknowledgement::TransportReceipt,
527        durability: None,
528        replicas: None,
529        ordering: OrderingMode::None,
530        deduplication: Deduplication::None,
531        dedup_window_ms: None,
532        backpressure: Backpressure::Block,
533        producer_naming: ProducerNaming::Fingerprint,
534        control_isolated: false,
535    };
536
537    /// Is this the default set? A `core` declaration is never written: it is
538    /// what an absent key already means.
539    pub fn is_core(&self) -> bool {
540        *self == GuaranteeSet::CORE
541    }
542
543    /// Checks the dimension combinations §6.5 forbids.
544    fn validate(&self) -> Result<(), HeaderError> {
545        let stored_or_replicated = matches!(
546            self.acknowledgement,
547            Acknowledgement::Stored | Acknowledgement::Replicated
548        );
549        if self.durability.is_some() && !stored_or_replicated {
550            return Err(HeaderError::InvalidGuarantees(
551                "durability without Stored or Replicated",
552            ));
553        }
554        match self.replicas {
555            Some(_) if self.acknowledgement != Acknowledgement::Replicated => {
556                return Err(HeaderError::InvalidGuarantees(
557                    "replicas without Replicated",
558                ));
559            }
560            Some(n) if n < 2 => {
561                return Err(HeaderError::InvalidGuarantees(
562                    "a replica count below 2 is not a replication",
563                ));
564            }
565            _ => {}
566        }
567        match (self.deduplication, self.dedup_window_ms) {
568            (Deduplication::Bounded, None) => {
569                return Err(HeaderError::InvalidGuarantees(
570                    "Bounded deduplication without a window",
571                ));
572            }
573            (level, Some(_)) if level != Deduplication::Bounded => {
574                return Err(HeaderError::InvalidGuarantees(
575                    "a dedup window without Bounded deduplication",
576                ));
577            }
578            _ => {}
579        }
580        Ok(())
581    }
582
583    /// Writes the set as a CBOR map, omitting every dimension left at `core`.
584    fn encode_into(
585        &self,
586        e: &mut Encoder<Vec<u8>>,
587    ) -> Result<(), minicbor::encode::Error<Infallible>> {
588        let core = GuaranteeSet::CORE;
589        let count = u64::from(self.delivery != core.delivery)
590            + u64::from(self.acknowledgement != core.acknowledgement)
591            + u64::from(self.durability.is_some())
592            + u64::from(self.replicas.is_some())
593            + u64::from(self.ordering != core.ordering)
594            + u64::from(self.deduplication != core.deduplication)
595            + u64::from(self.dedup_window_ms.is_some())
596            + u64::from(self.backpressure != core.backpressure)
597            + u64::from(self.producer_naming != core.producer_naming)
598            + u64::from(self.control_isolated != core.control_isolated);
599        e.map(count)?;
600        if self.delivery != core.delivery {
601            e.u64(guarantee_key::DELIVERY)?
602                .u64(self.delivery.to_wire())?;
603        }
604        if self.acknowledgement != core.acknowledgement {
605            e.u64(guarantee_key::ACKNOWLEDGEMENT)?
606                .u64(self.acknowledgement.to_wire())?;
607        }
608        if let Some(durability) = self.durability {
609            e.u64(guarantee_key::DURABILITY)?
610                .u64(durability.to_wire())?;
611        }
612        if let Some(replicas) = self.replicas {
613            e.u64(guarantee_key::REPLICAS)?.u64(replicas)?;
614        }
615        if self.ordering != core.ordering {
616            e.u64(guarantee_key::ORDERING)?
617                .u64(self.ordering.to_wire())?;
618        }
619        if self.deduplication != core.deduplication {
620            e.u64(guarantee_key::DEDUPLICATION)?
621                .u64(self.deduplication.to_wire())?;
622        }
623        if let Some(window) = self.dedup_window_ms {
624            e.u64(guarantee_key::DEDUP_WINDOW_MS)?.u64(window)?;
625        }
626        if self.backpressure != core.backpressure {
627            e.u64(guarantee_key::BACKPRESSURE)?
628                .u64(self.backpressure.to_wire())?;
629        }
630        if self.producer_naming != core.producer_naming {
631            e.u64(guarantee_key::PRODUCER_NAMING)?
632                .u64(self.producer_naming.to_wire())?;
633        }
634        if self.control_isolated != core.control_isolated {
635            e.u64(guarantee_key::CONTROL_ISOLATED)?
636                .u64(u64::from(self.control_isolated))?;
637        }
638        Ok(())
639    }
640
641    /// Reads a set from the nested map at the decoder's position.
642    ///
643    /// The nesting is one level deep by specification (§5), and an unknown
644    /// dimension is skipped exactly like an unknown top-level key.
645    fn decode_from(m: &mut MapReader<'_, '_>) -> Result<GuaranteeSet, HeaderError> {
646        let mut set = GuaranteeSet::CORE;
647        let mut inner = MapReader::new(m.d)?;
648        while let Some(key) = inner.next_key()? {
649            match key {
650                guarantee_key::DELIVERY => set.delivery = level(inner.u64()?, "delivery")?,
651                guarantee_key::ACKNOWLEDGEMENT => {
652                    set.acknowledgement = level(inner.u64()?, "acknowledgement")?;
653                }
654                guarantee_key::DURABILITY => {
655                    set.durability = Some(level(inner.u64()?, "durability")?);
656                }
657                guarantee_key::REPLICAS => set.replicas = Some(inner.u64()?),
658                guarantee_key::ORDERING => set.ordering = level(inner.u64()?, "ordering")?,
659                guarantee_key::DEDUPLICATION => {
660                    set.deduplication = level(inner.u64()?, "deduplication")?;
661                }
662                guarantee_key::DEDUP_WINDOW_MS => set.dedup_window_ms = Some(inner.u64()?),
663                guarantee_key::BACKPRESSURE => {
664                    set.backpressure = level(inner.u64()?, "backpressure")?;
665                }
666                guarantee_key::PRODUCER_NAMING => {
667                    set.producer_naming = level(inner.u64()?, "producer naming")?;
668                }
669                guarantee_key::CONTROL_ISOLATED => {
670                    set.control_isolated = match inner.u64()? {
671                        0 => false,
672                        1 => true,
673                        _ => {
674                            return Err(HeaderError::InvalidGuarantees(
675                                "control_isolated is 0 or 1",
676                            ));
677                        }
678                    };
679                }
680                _ => inner.skip()?,
681            }
682        }
683        set.validate()?;
684        Ok(set)
685    }
686
687    /// The weaker of two offered sets, dimension by dimension
688    /// (`docs/PROTOCOL.md` §2.3 step 5).
689    ///
690    /// Ladders take the minimum. Dimensions that are **not** ordered —
691    /// backpressure, producer naming, and the two independent axes of a
692    /// durability level — have no "weaker", so the two declarations must be
693    /// equal; the name of the dimension comes back as the error so a peer can
694    /// be told which one disagreed.
695    pub fn intersect(&self, other: &GuaranteeSet) -> Result<GuaranteeSet, &'static str> {
696        if self.backpressure != other.backpressure {
697            return Err("backpressure");
698        }
699        if self.producer_naming != other.producer_naming {
700            return Err("producer naming");
701        }
702        if self.durability.is_some()
703            && other.durability.is_some()
704            && self.durability != other.durability
705        {
706            return Err("durability");
707        }
708        if self.replicas.is_some() && other.replicas.is_some() && self.replicas != other.replicas {
709            return Err("replicas");
710        }
711
712        let acknowledgement = self.acknowledgement.min(other.acknowledgement);
713        let keeps_durability = matches!(
714            acknowledgement,
715            Acknowledgement::Stored | Acknowledgement::Replicated
716        );
717        let deduplication = self.deduplication.min(other.deduplication);
718        let mut merged = GuaranteeSet {
719            delivery: self.delivery.min(other.delivery),
720            acknowledgement,
721            // A dimension the weakened acknowledgement can no longer carry is
722            // dropped rather than kept: dropping is what "weaker" means here,
723            // and keeping it would produce a set §6.5 forbids.
724            durability: keeps_durability
725                .then_some(self.durability.or(other.durability))
726                .flatten(),
727            replicas: (acknowledgement == Acknowledgement::Replicated)
728                .then_some(self.replicas.or(other.replicas))
729                .flatten(),
730            ordering: self.ordering.min(other.ordering),
731            deduplication,
732            // A shorter window is the weaker promise.
733            dedup_window_ms: None,
734            backpressure: self.backpressure,
735            producer_naming: self.producer_naming,
736            control_isolated: self.control_isolated && other.control_isolated,
737        };
738        if deduplication == Deduplication::Bounded {
739            merged.dedup_window_ms = match (self.dedup_window_ms, other.dedup_window_ms) {
740                (Some(a), Some(b)) => Some(a.min(b)),
741                (Some(a), None) | (None, Some(a)) => Some(a),
742                (None, None) => None,
743            };
744        }
745        Ok(merged)
746    }
747
748    /// Does this set reach `required` on every dimension?
749    ///
750    /// Ladders compare by level, the durability axes compare per axis, and the
751    /// unordered dimensions must match exactly. A longer dedup window is the
752    /// stronger promise.
753    pub fn reaches(&self, required: &GuaranteeSet) -> bool {
754        if self.delivery < required.delivery
755            || self.acknowledgement < required.acknowledgement
756            || self.ordering < required.ordering
757            || self.deduplication < required.deduplication
758        {
759            return false;
760        }
761        if self.backpressure != required.backpressure
762            || self.producer_naming != required.producer_naming
763        {
764            return false;
765        }
766        if !self.control_isolated && required.control_isolated {
767            return false;
768        }
769        match (self.durability, required.durability) {
770            (_, None) => {}
771            (Some(have), Some(want)) if have >= want => {}
772            _ => return false,
773        }
774        match (self.replicas, required.replicas) {
775            (_, None) => {}
776            (Some(have), Some(want)) if have >= want => {}
777            _ => return false,
778        }
779        match (self.dedup_window_ms, required.dedup_window_ms) {
780            (_, None) => {}
781            (Some(have), Some(want)) if have >= want => {}
782            _ => return false,
783        }
784        true
785    }
786}
787
788/// Maps a wire value to a level, naming the dimension when it is unknown.
789fn level<T: WireLevel>(value: u64, dimension: &'static str) -> Result<T, HeaderError> {
790    T::from_wire_value(value).ok_or(HeaderError::UnknownLevel { dimension, value })
791}
792
793/// Maps a wire value to a segment layer, refusing one above
794/// [`limits::MAX_LAYER`] with `reason`.
795fn layer(value: u64, reason: &'static str) -> Result<u8, HeaderError> {
796    u8::try_from(value)
797        .ok()
798        .filter(|layer| *layer <= limits::MAX_LAYER)
799        .ok_or(HeaderError::InvalidLayer(reason))
800}
801
802/// Lets [`level`] work for every dimension enum without a macro per call.
803trait WireLevel: Sized {
804    fn from_wire_value(value: u64) -> Option<Self>;
805}
806
807macro_rules! impl_wire_level {
808    ($($name:ident),+ $(,)?) => {
809        $(impl WireLevel for $name {
810            fn from_wire_value(value: u64) -> Option<$name> {
811                $name::from_wire(value)
812            }
813        })+
814    };
815}
816
817impl_wire_level!(
818    Delivery,
819    Acknowledgement,
820    Durability,
821    OrderingMode,
822    Deduplication,
823    Backpressure,
824    ProducerNaming,
825    ReportMode,
826);
827
828/// Why a header was rejected. Every variant is a protocol violation.
829#[derive(Clone, Debug, PartialEq, Eq)]
830pub enum HeaderError {
831    /// The bytes are not well-formed CBOR, or a value had the wrong type.
832    Malformed(&'static str),
833    /// An indefinite-length item was used where the protocol forbids it.
834    Indefinite,
835    /// A map key appeared twice.
836    DuplicateKey(u64),
837    /// A map key was not greater than the preceding one; keys must ascend.
838    UnorderedKey(u64),
839    /// A map key was not an unsigned integer.
840    NonUintKey,
841    /// A required key was absent.
842    MissingKey(u64),
843    /// A text field exceeded its cap.
844    StringTooLong {
845        /// Key of the offending field.
846        key: u64,
847        /// Length found.
848        len: usize,
849        /// Cap for this field.
850        max: usize,
851    },
852    /// A list field declared more items than the decoder accepts.
853    ListTooLong {
854        /// Key of the offending field.
855        key: u64,
856        /// Declared item count.
857        len: u64,
858        /// Cap for list fields.
859        max: usize,
860    },
861    /// An unknown field nested deeper than [`limits::MAX_SKIP_DEPTH`].
862    DepthExceeded,
863    /// Bytes remained after the header map.
864    TrailingBytes,
865    /// A guarantee set, or a DATA header's achieved level, named a level this
866    /// version does not define.
867    UnknownLevel {
868        /// Dimension whose value was unknown.
869        dimension: &'static str,
870        /// The value found.
871        value: u64,
872    },
873    /// A guarantee set's dimension combination is one §6.5 forbids, or a
874    /// HELLO requires more than it offers (§6.1).
875    InvalidGuarantees(&'static str),
876    /// A topic filter violated the grammar of `docs/PROTOCOL.md` §6.4.
877    InvalidFilter(&'static str),
878    /// A DATA header's report order is malformed: the levels do not ascend,
879    /// there are too many of them, or the order and its id disagree
880    /// (`docs/PROTOCOL.md` §6.2, keys `9`-`11`).
881    InvalidReport(&'static str),
882    /// A segment layer is above [`limits::MAX_LAYER`], or a DATA header
883    /// carries key `14` without key `13` (`docs/PROTOCOL.md` §6.2 key `14`,
884    /// §6.4 key `3`).
885    InvalidLayer(&'static str),
886    /// A REPORT record is longer than [`limits::MAX_PATH_RECORD_BYTES`]
887    /// (`docs/PROTOCOL.md` §6.10).
888    InvalidPathReport(&'static str),
889}
890
891impl std::fmt::Display for HeaderError {
892    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
893        match self {
894            HeaderError::Malformed(what) => write!(f, "malformed header: {what}"),
895            HeaderError::Indefinite => f.write_str("indefinite-length items are not allowed"),
896            HeaderError::DuplicateKey(k) => write!(f, "duplicate header key {k}"),
897            HeaderError::UnorderedKey(k) => {
898                write!(f, "header key {k} is out of ascending order")
899            }
900            HeaderError::NonUintKey => f.write_str("header key is not an unsigned integer"),
901            HeaderError::MissingKey(k) => write!(f, "required header key {k} is missing"),
902            HeaderError::StringTooLong { key, len, max } => {
903                write!(
904                    f,
905                    "key {key}: text of {len} bytes exceeds the {max} byte cap"
906                )
907            }
908            HeaderError::ListTooLong { key, len, max } => {
909                write!(
910                    f,
911                    "key {key}: list of {len} items exceeds the {max} item cap"
912                )
913            }
914            HeaderError::DepthExceeded => f.write_str("unknown field nested too deeply"),
915            HeaderError::TrailingBytes => f.write_str("trailing bytes after the header"),
916            HeaderError::UnknownLevel { dimension, value } => {
917                write!(f, "unknown {dimension} level {value}")
918            }
919            HeaderError::InvalidGuarantees(why) => write!(f, "invalid guarantee set: {why}"),
920            HeaderError::InvalidFilter(why) => write!(f, "invalid topic filter: {why}"),
921            HeaderError::InvalidReport(reason) => write!(f, "invalid report: {reason}"),
922            HeaderError::InvalidLayer(reason) => write!(f, "invalid layer: {reason}"),
923            HeaderError::InvalidPathReport(reason) => write!(f, "invalid path report: {reason}"),
924        }
925    }
926}
927
928impl std::error::Error for HeaderError {}
929
930impl From<HeaderError> for Error {
931    fn from(e: HeaderError) -> Error {
932        Error::Protocol(e.to_string())
933    }
934}
935
936/// Appends an encoded header to a buffer the caller owns.
937///
938/// Every `encode` in this file is a thin wrapper over an `encode_into` that
939/// goes through here, so a hot send path can reuse one buffer and the
940/// canonical form has exactly one implementation — the property that matters,
941/// because two encoders that can disagree about key order would be a wire
942/// divergence rather than an optimisation (B-250).
943fn encode_into_with(
944    out: &mut Vec<u8>,
945    f: impl FnOnce(&mut Encoder<Vec<u8>>) -> Result<(), minicbor::encode::Error<Infallible>>,
946) {
947    // `Encoder` owns its writer, so the buffer is handed over and taken back.
948    // `std::mem::take` keeps the caller's allocation: the `Vec` that comes
949    // back is the same one, grown at most by this header.
950    let mut e = Encoder::new(std::mem::take(out));
951    f(&mut e).expect("encoding into a Vec is infallible");
952    *out = e.into_writer();
953}
954
955/// Skips one CBOR value iteratively, refusing to recurse and refusing to nest
956/// deeper than `max_depth`.
957///
958/// `minicbor` has its own `skip`, but the protocol requires a specific,
959/// auditable bound on hostile nesting; a value that nests deeper is rejected
960/// rather than tolerated.
961fn skip_value(d: &mut Decoder<'_>, max_depth: usize) -> Result<(), HeaderError> {
962    // `stack` holds the outstanding item counts of enclosing containers; its
963    // length is the current nesting depth and is bounded by `max_depth`.
964    let mut stack: Vec<u64> = Vec::new();
965    let mut remaining: u64 = 1;
966
967    loop {
968        if remaining == 0 {
969            match stack.pop() {
970                Some(outer) => {
971                    remaining = outer;
972                    continue;
973                }
974                None => return Ok(()),
975            }
976        }
977        remaining -= 1;
978
979        let ty = d
980            .datatype()
981            .map_err(|_| HeaderError::Malformed("truncated value"))?;
982        let nested = match ty {
983            Type::Bool => {
984                d.bool().map_err(|_| HeaderError::Malformed("bool"))?;
985                None
986            }
987            Type::Null => {
988                d.null().map_err(|_| HeaderError::Malformed("null"))?;
989                None
990            }
991            Type::Undefined => {
992                d.undefined()
993                    .map_err(|_| HeaderError::Malformed("undefined"))?;
994                None
995            }
996            Type::U8
997            | Type::U16
998            | Type::U32
999            | Type::U64
1000            | Type::I8
1001            | Type::I16
1002            | Type::I32
1003            | Type::I64
1004            | Type::Int => {
1005                d.int().map_err(|_| HeaderError::Malformed("integer"))?;
1006                None
1007            }
1008            Type::F32 | Type::F64 => {
1009                d.f64().map_err(|_| HeaderError::Malformed("float"))?;
1010                None
1011            }
1012            Type::Bytes => {
1013                d.bytes()
1014                    .map_err(|_| HeaderError::Malformed("byte string"))?;
1015                None
1016            }
1017            Type::String => {
1018                d.str().map_err(|_| HeaderError::Malformed("text string"))?;
1019                None
1020            }
1021            Type::Array => Some(
1022                d.array()
1023                    .map_err(|_| HeaderError::Malformed("array"))?
1024                    .ok_or(HeaderError::Indefinite)?,
1025            ),
1026            Type::Map => {
1027                let pairs = d
1028                    .map()
1029                    .map_err(|_| HeaderError::Malformed("map"))?
1030                    .ok_or(HeaderError::Indefinite)?;
1031                Some(
1032                    pairs
1033                        .checked_mul(2)
1034                        .ok_or(HeaderError::Malformed("map length overflow"))?,
1035                )
1036            }
1037            Type::BytesIndef | Type::StringIndef | Type::ArrayIndef | Type::MapIndef => {
1038                return Err(HeaderError::Indefinite);
1039            }
1040            Type::Break => return Err(HeaderError::Malformed("unexpected break")),
1041            // Tags, half-floats and other simple values carry no meaning in
1042            // weida headers; extensions must use plain data items.
1043            Type::Tag => return Err(HeaderError::Malformed("tags are not allowed")),
1044            Type::F16 => return Err(HeaderError::Malformed("half floats are not allowed")),
1045            Type::Simple => return Err(HeaderError::Malformed("simple values are not allowed")),
1046            Type::Unknown(_) => return Err(HeaderError::Malformed("unknown major type")),
1047        };
1048
1049        if let Some(count) = nested
1050            && count > 0
1051        {
1052            if stack.len() >= max_depth {
1053                return Err(HeaderError::DepthExceeded);
1054            }
1055            stack.push(remaining);
1056            remaining = count;
1057        }
1058    }
1059}
1060
1061/// Reader for one strict header map.
1062struct MapReader<'a, 'b> {
1063    d: &'a mut Decoder<'b>,
1064    remaining: u64,
1065    /// Bitmask of seen keys `0..=63`, for the required-key checks. The
1066    /// specification reserves that range, so presence needs no allocation.
1067    seen: u64,
1068    /// Previously read key.
1069    ///
1070    /// Keys are required to ascend strictly, which makes duplicate detection
1071    /// complete for *every* key — including extension keys the decoder skips —
1072    /// in constant space. A set of seen extension keys would be exactly the
1073    /// remote-controlled allocation the invariants forbid.
1074    last: Option<u64>,
1075}
1076
1077impl<'a, 'b> MapReader<'a, 'b> {
1078    fn new(d: &'a mut Decoder<'b>) -> Result<MapReader<'a, 'b>, HeaderError> {
1079        let len = d
1080            .map()
1081            .map_err(|_| HeaderError::Malformed("header is not a map"))?
1082            .ok_or(HeaderError::Indefinite)?;
1083        Ok(MapReader {
1084            d,
1085            remaining: len,
1086            seen: 0,
1087            last: None,
1088        })
1089    }
1090
1091    fn next_key(&mut self) -> Result<Option<u64>, HeaderError> {
1092        if self.remaining == 0 {
1093            return Ok(None);
1094        }
1095        self.remaining -= 1;
1096        match self.d.datatype() {
1097            Ok(Type::U8 | Type::U16 | Type::U32 | Type::U64) => {}
1098            Ok(_) => return Err(HeaderError::NonUintKey),
1099            Err(_) => return Err(HeaderError::Malformed("truncated key")),
1100        }
1101        let key = self.d.u64().map_err(|_| HeaderError::NonUintKey)?;
1102        if let Some(prev) = self.last {
1103            if key == prev {
1104                return Err(HeaderError::DuplicateKey(key));
1105            }
1106            if key < prev {
1107                return Err(HeaderError::UnorderedKey(key));
1108            }
1109        }
1110        self.last = Some(key);
1111        if key < 64 {
1112            self.seen |= 1u64 << key;
1113        }
1114        Ok(Some(key))
1115    }
1116
1117    fn saw(&self, key: u64) -> bool {
1118        key < 64 && self.seen & (1u64 << key) != 0
1119    }
1120
1121    fn require(&self, key: u64) -> Result<(), HeaderError> {
1122        if self.saw(key) {
1123            Ok(())
1124        } else {
1125            Err(HeaderError::MissingKey(key))
1126        }
1127    }
1128
1129    fn u64(&mut self) -> Result<u64, HeaderError> {
1130        self.d
1131            .u64()
1132            .map_err(|_| HeaderError::Malformed("expected an unsigned integer"))
1133    }
1134
1135    fn text(&mut self, key: u64, max: usize) -> Result<String, HeaderError> {
1136        let s = self
1137            .d
1138            .str()
1139            .map_err(|_| HeaderError::Malformed("expected a text string"))?;
1140        if s.len() > max {
1141            return Err(HeaderError::StringTooLong {
1142                key,
1143                len: s.len(),
1144                max,
1145            });
1146        }
1147        Ok(s.to_owned())
1148    }
1149
1150    /// Reads a byte string of exactly `N` bytes.
1151    ///
1152    /// The cap is checked before the length is trusted for anything, and a
1153    /// shorter value is rejected rather than padded: `docs/PROTOCOL.md` §6.2
1154    /// defines the one field that uses this as a raw 32-byte digest, and half
1155    /// a digest identifies nobody.
1156    fn byte_array<const N: usize>(&mut self, key: u64) -> Result<[u8; N], HeaderError> {
1157        let bytes = self
1158            .d
1159            .bytes()
1160            .map_err(|_| HeaderError::Malformed("expected a byte string"))?;
1161        if bytes.len() > N {
1162            return Err(HeaderError::StringTooLong {
1163                key,
1164                len: bytes.len(),
1165                max: N,
1166            });
1167        }
1168        bytes
1169            .try_into()
1170            .map_err(|_| HeaderError::Malformed("byte string has the wrong length"))
1171    }
1172
1173    fn uint_list(&mut self, key: u64) -> Result<Vec<u64>, HeaderError> {
1174        let len = self
1175            .d
1176            .array()
1177            .map_err(|_| HeaderError::Malformed("expected an array"))?
1178            .ok_or(HeaderError::Indefinite)?;
1179        if len > limits::MAX_LIST_ITEMS as u64 {
1180            return Err(HeaderError::ListTooLong {
1181                key,
1182                len,
1183                max: limits::MAX_LIST_ITEMS,
1184            });
1185        }
1186        // `len` is now bounded by MAX_LIST_ITEMS, so reserving is safe.
1187        let mut out = Vec::with_capacity(len as usize);
1188        for _ in 0..len {
1189            out.push(self.u64()?);
1190        }
1191        Ok(out)
1192    }
1193
1194    /// Reads a report order: a definite-length array of strictly ascending
1195    /// cursor levels, capped at [`limits::MAX_REPORT_LEVELS`].
1196    ///
1197    /// Ascent is checked here rather than after the fact for the same reason
1198    /// map keys are: it makes duplicate detection complete in constant space,
1199    /// and it makes the wire form canonical, so two peers ordering the same
1200    /// levels send the same bytes.
1201    fn report_levels(&mut self) -> Result<Vec<CursorLevel>, HeaderError> {
1202        let len = self
1203            .d
1204            .array()
1205            .map_err(|_| HeaderError::Malformed("expected an array"))?
1206            .ok_or(HeaderError::Indefinite)?;
1207        if len > limits::MAX_REPORT_LEVELS as u64 {
1208            return Err(HeaderError::InvalidReport("too many report levels"));
1209        }
1210        // `len` is now bounded by MAX_REPORT_LEVELS, so reserving is safe.
1211        let mut out: Vec<CursorLevel> = Vec::with_capacity(len as usize);
1212        let mut last: Option<u64> = None;
1213        for _ in 0..len {
1214            let value = self.u64()?;
1215            if let Some(prev) = last
1216                && value <= prev
1217            {
1218                return Err(HeaderError::InvalidReport("report levels must ascend"));
1219            }
1220            last = Some(value);
1221            out.push(
1222                CursorLevel::from_wire(value).ok_or(HeaderError::UnknownLevel {
1223                    dimension: "report",
1224                    value,
1225                })?,
1226            );
1227        }
1228        Ok(out)
1229    }
1230
1231    fn skip(&mut self) -> Result<(), HeaderError> {
1232        skip_value(self.d, limits::MAX_SKIP_DEPTH)
1233    }
1234}
1235
1236/// Rejects trailing bytes after a header map.
1237fn finish(d: &Decoder<'_>) -> Result<(), HeaderError> {
1238    if d.position() == d.input().len() {
1239        Ok(())
1240    } else {
1241        Err(HeaderError::TrailingBytes)
1242    }
1243}
1244
1245/// HELLO header: connection negotiation input.
1246#[derive(Clone, Debug, PartialEq, Eq)]
1247pub struct Hello {
1248    /// Wire protocol versions the sender supports.
1249    pub versions: Vec<u64>,
1250    /// Largest header the sender is willing to receive.
1251    pub max_header_bytes: u64,
1252    /// Advisory concurrent inbound transfer count.
1253    pub max_transfers: u64,
1254    /// Capability codes the sender supports.
1255    pub capabilities: Vec<u64>,
1256    /// Capability codes the sender requires the peer to support.
1257    pub required_capabilities: Vec<u64>,
1258    /// Guarantee set the sender can honour (key `5`).
1259    ///
1260    /// `None` means the default set: an absent key and
1261    /// [`GuaranteeSet::CORE`] are the same declaration, which is why a v0
1262    /// HELLO is unchanged on the wire.
1263    pub guarantees_offered: Option<GuaranteeSet>,
1264    /// Guarantee set the sender requires of the peer (key `6`).
1265    ///
1266    /// MUST be reachable by `guarantees_offered` on every dimension: requiring
1267    /// what you cannot honour yourself is a configuration error
1268    /// (`docs/PROTOCOL.md` §6.1), and a decoder rejects it.
1269    pub guarantees_required: Option<GuaranteeSet>,
1270}
1271
1272impl Hello {
1273    /// The HELLO a v0 implementation sends: no guarantee declarations, so
1274    /// `core` offered and `core` required.
1275    pub fn v0(max_header_bytes: u64, max_transfers: u64) -> Hello {
1276        Hello {
1277            versions: vec![crate::VERSION],
1278            max_header_bytes,
1279            max_transfers,
1280            capabilities: Vec::new(),
1281            required_capabilities: Vec::new(),
1282            guarantees_offered: None,
1283            guarantees_required: None,
1284        }
1285    }
1286
1287    /// The set this HELLO offers; an absent declaration means `core`.
1288    pub fn offered(&self) -> GuaranteeSet {
1289        self.guarantees_offered.unwrap_or(GuaranteeSet::CORE)
1290    }
1291
1292    /// The set this HELLO requires; an absent declaration means `core`.
1293    pub fn required(&self) -> GuaranteeSet {
1294        self.guarantees_required.unwrap_or(GuaranteeSet::CORE)
1295    }
1296
1297    /// Encodes the header.
1298    pub fn encode(&self) -> Vec<u8> {
1299        let mut out = Vec::new();
1300        self.encode_into(&mut out);
1301        out
1302    }
1303
1304    /// Appends the encoded header to `out`, for a send path that reuses a
1305    /// buffer (B-250). The canonical form has one implementation and this is
1306    /// it; [`Self::encode`] is a wrapper.
1307    pub fn encode_into(&self, out: &mut Vec<u8>) {
1308        encode_into_with(out, |e| {
1309            // A `core` declaration is never written: an absent key already
1310            // says it, and a v0 HELLO must stay byte-identical (§6.1).
1311            let offered = self.guarantees_offered.filter(|s| !s.is_core());
1312            let required = self.guarantees_required.filter(|s| !s.is_core());
1313            e.map(5 + u64::from(offered.is_some()) + u64::from(required.is_some()))?;
1314            e.u64(hello_key::VERSIONS)?
1315                .array(self.versions.len() as u64)?;
1316            for v in &self.versions {
1317                e.u64(*v)?;
1318            }
1319            e.u64(hello_key::MAX_HEADER_BYTES)?
1320                .u64(self.max_header_bytes)?;
1321            e.u64(hello_key::MAX_TRANSFERS)?.u64(self.max_transfers)?;
1322            e.u64(hello_key::CAPABILITIES)?
1323                .array(self.capabilities.len() as u64)?;
1324            for c in &self.capabilities {
1325                e.u64(*c)?;
1326            }
1327            e.u64(hello_key::REQUIRED_CAPABILITIES)?
1328                .array(self.required_capabilities.len() as u64)?;
1329            for c in &self.required_capabilities {
1330                e.u64(*c)?;
1331            }
1332            if let Some(set) = offered {
1333                e.u64(hello_key::GUARANTEES_OFFERED)?;
1334                set.encode_into(e)?;
1335            }
1336            if let Some(set) = required {
1337                e.u64(hello_key::GUARANTEES_REQUIRED)?;
1338                set.encode_into(e)?;
1339            }
1340            Ok(())
1341        })
1342    }
1343
1344    /// Decodes the header.
1345    pub fn decode(bytes: &[u8]) -> Result<Hello, HeaderError> {
1346        let mut d = Decoder::new(bytes);
1347        let mut versions = Vec::new();
1348        let mut max_header_bytes = 0;
1349        let mut max_transfers = 0;
1350        let mut capabilities = Vec::new();
1351        let mut required_capabilities = Vec::new();
1352        let mut guarantees_offered = None;
1353        let mut guarantees_required = None;
1354        {
1355            let mut m = MapReader::new(&mut d)?;
1356            while let Some(key) = m.next_key()? {
1357                match key {
1358                    hello_key::VERSIONS => versions = m.uint_list(key)?,
1359                    hello_key::MAX_HEADER_BYTES => max_header_bytes = m.u64()?,
1360                    hello_key::MAX_TRANSFERS => max_transfers = m.u64()?,
1361                    hello_key::CAPABILITIES => capabilities = m.uint_list(key)?,
1362                    hello_key::REQUIRED_CAPABILITIES => required_capabilities = m.uint_list(key)?,
1363                    hello_key::GUARANTEES_OFFERED => {
1364                        guarantees_offered = Some(GuaranteeSet::decode_from(&mut m)?);
1365                    }
1366                    hello_key::GUARANTEES_REQUIRED => {
1367                        guarantees_required = Some(GuaranteeSet::decode_from(&mut m)?);
1368                    }
1369                    _ => m.skip()?,
1370                }
1371            }
1372            for key in [
1373                hello_key::VERSIONS,
1374                hello_key::MAX_HEADER_BYTES,
1375                hello_key::MAX_TRANSFERS,
1376                hello_key::CAPABILITIES,
1377                hello_key::REQUIRED_CAPABILITIES,
1378            ] {
1379                m.require(key)?;
1380            }
1381        }
1382        finish(&d)?;
1383        let hello = Hello {
1384            versions,
1385            max_header_bytes,
1386            max_transfers,
1387            capabilities,
1388            required_capabilities,
1389            guarantees_offered,
1390            guarantees_required,
1391        };
1392        // §6.1: requiring more than you offer is a configuration error, and
1393        // one a decoder can see in a single header.
1394        if !hello.offered().reaches(&hello.required()) {
1395            return Err(HeaderError::InvalidGuarantees(
1396                "guarantees_required is not covered by guarantees_offered",
1397            ));
1398        }
1399        Ok(hello)
1400    }
1401}
1402
1403/// DATA header: one transfer.
1404///
1405/// Every field is optional at the decoder. Which of them the *context*
1406/// requires is a dispatch question: an initiating stream must name an endpoint
1407/// and the reply half of an exchange must not, but the decoder sees bytes, not
1408/// streams (`docs/PROTOCOL.md` §6.2).
1409#[derive(Clone, Debug, Default, PartialEq, Eq)]
1410pub struct DataHeader {
1411    /// Endpoint path. Required on an initiating stream, ignored on a reply.
1412    pub endpoint: Option<String>,
1413    /// Advisory payload length.
1414    pub content_len: Option<u64>,
1415    /// Opaque content type label.
1416    pub content_type: Option<String>,
1417    /// W3C `traceparent`.
1418    pub traceparent: Option<String>,
1419    /// W3C `tracestate`, opaque passthrough.
1420    pub tracestate: Option<String>,
1421    /// Pub/Sub topic; opaque bytes, selected by the filter grammar of
1422    /// `docs/PROTOCOL.md` §6.4. Only meaningful on transfers fanned out by a
1423    /// publisher.
1424    pub topic: Option<String>,
1425    /// Per-producer sequence number, for ordering and gap detection
1426    /// (`docs/PROTOCOL.md` §6.2, key `6`).
1427    ///
1428    /// Written by a publisher whose connection negotiated `PerProducer`
1429    /// ordering, and by nothing under `core`: the number is assigned once per
1430    /// published message, before fan-out, so a copy a subscriber lost shows up
1431    /// as a hole in its own sequence. It is not a transfer identifier and
1432    /// correlates nothing — an exchange is correlated by its stream.
1433    pub sequence: Option<u64>,
1434    /// Producer identity: the raw 32-byte digest (`docs/PROTOCOL.md` §6.2,
1435    /// key `7`).
1436    ///
1437    /// **Specified ahead of code**, and absent in the default case by design:
1438    /// the receiver already knows the sending peer's proved fingerprint from
1439    /// the handshake, so this names a producer only where it is *not* the
1440    /// connection peer — a relay, or a name an L2 subscription supplies
1441    /// ([decisions/0008](../../../docs/decisions/0008-session-identity.md)
1442    /// §4.4). The `sha256:<64 hex>` spelling is presentation only and never
1443    /// goes on the wire.
1444    pub producer: Option<[u8; limits::PRODUCER_BYTES]>,
1445    /// The completion level the sender **achieved** for the message it is
1446    /// answering (`docs/PROTOCOL.md` §6.2, key `8`).
1447    ///
1448    /// This is the L2 confirm, and it is a statement about one hop: a broker
1449    /// that has taken responsibility for a message in memory writes
1450    /// [`Acknowledgement::Accepted`] on the reply half of the producer's
1451    /// exchange, which is what makes the reply a publisher confirm without a
1452    /// frame kind of its own
1453    /// ([decisions/0018](../../../docs/decisions/0018-minimal-broker.md)
1454    /// §4.6). It is *achieved*, never requested — a level a peer wants is
1455    /// negotiated in HELLO and refused there if it cannot be reached
1456    /// ([0006](../../../docs/decisions/0006-guarantee-sets.md) §4.4) — and it
1457    /// is never relayed: the producer's confirm says nothing about what a
1458    /// consumer later does with the message
1459    /// ([GUARANTEES.md] §2).
1460    ///
1461    /// A v0 sender leaves it absent, and an absent key is not
1462    /// `Acknowledgement::None`: it says this hop makes no claim beyond the
1463    /// transport receipt QUIC already gave.
1464    ///
1465    /// [GUARANTEES.md]: https://git.doodleshnookie.net/tuco86/weida/blob/main/docs/GUARANTEES.md
1466    pub achieved: Option<Acknowledgement>,
1467    /// Identifier the sender assigns to the report it orders (key `9`).
1468    ///
1469    /// Present exactly when [`DataHeader::report`] is non-empty. It names the
1470    /// CURSOR stream that will report on *this* transfer, and it is scoped to
1471    /// the connection and to the direction that allocated it: a peer reports
1472    /// only on transfers it received, so the two directions cannot collide.
1473    pub report_id: Option<u64>,
1474    /// Levels the sender asks to be reported, strictly ascending (key `10`).
1475    ///
1476    /// An **order**, not a guarantee: a receiver that cannot reach a level
1477    /// simply does not report it, and the transfer does not fail for it. A
1478    /// level a peer must reach is the negotiated `acknowledgement` dimension
1479    /// of HELLO instead
1480    /// ([decisions/0006](../../../docs/decisions/0006-guarantee-sets.md)
1481    /// §4.4).
1482    pub report: Vec<CursorLevel>,
1483    /// How often the reporter should emit a record (key `11`).
1484    ///
1485    /// [`ReportMode::Progress`] is the default and is never written.
1486    pub report_mode: ReportMode,
1487    /// Segment number per `(sender, path, topic)`, from 0 (key `13`).
1488    ///
1489    /// Written by a radio or by `Peer::segment`
1490    /// ([decisions/0034](../../../docs/decisions/0034-late-is-lost.md)
1491    /// §4.6, [decisions/0037](../../../docs/decisions/0037-layered-segments.md)
1492    /// §4.2). It is not a `PerProducer` sequence: it is written whatever
1493    /// ordering was negotiated, and a dish uses it to discard a segment older
1494    /// than the newest it delivered.
1495    pub segment: Option<u64>,
1496    /// Segment layer, `0..=15` (key `14`); absent means 0; written only
1497    /// beside `segment`
1498    /// ([decisions/0037](../../../docs/decisions/0037-layered-segments.md)
1499    /// §4.3). The encoder writes what it is given, so a sender passes `None`
1500    /// for layer 0.
1501    pub layer: Option<u8>,
1502}
1503
1504impl DataHeader {
1505    /// A header addressing `endpoint`, for the initiating half of a stream.
1506    pub fn addressed(endpoint: impl Into<String>) -> DataHeader {
1507        DataHeader {
1508            endpoint: Some(endpoint.into()),
1509            ..DataHeader::default()
1510        }
1511    }
1512
1513    /// A header for the reply half of an exchange: no endpoint, no topic.
1514    ///
1515    /// The stream is the correlation, so a reply carries no identifier of the
1516    /// request it answers.
1517    pub fn reply() -> DataHeader {
1518        DataHeader::default()
1519    }
1520
1521    /// Encodes the header.
1522    ///
1523    /// Key `10` goes out in the canonical form §6.2 makes normative —
1524    /// strictly ascending by wire value, no repeats — whatever order
1525    /// [`DataHeader::report`] happens to hold. That rule is enforced here
1526    /// because encoding cannot fail: a vector in any other order would
1527    /// otherwise produce bytes that close the connection at every conformant
1528    /// peer, and there would be no way to tell the caller so.
1529    ///
1530    /// The report's other two rules stay the caller's for exactly that
1531    /// reason — both need an error, and this function has none to give. At
1532    /// most [`limits::MAX_REPORT_LEVELS`] distinct levels, and key `9`
1533    /// present exactly when key `10` is: `weida`'s `data_header` refuses an
1534    /// oversized order with `Error::LimitExceeded` and allocates the report
1535    /// id alongside the order, so no caller reaches this encoder with either
1536    /// mistake.
1537    pub fn encode(&self) -> Vec<u8> {
1538        let mut out = Vec::new();
1539        self.encode_into(&mut out);
1540        out
1541    }
1542
1543    /// Appends the encoded header to `out`, for a send path that reuses a
1544    /// buffer (B-250). The canonical form has one implementation and this is
1545    /// it; [`Self::encode`] is a wrapper.
1546    pub fn encode_into(&self, out: &mut Vec<u8>) {
1547        // Sorted and deduplicated by wire value, not by the enum's derived
1548        // order, because the wire value is what ascends on the wire.
1549        let mut report: Vec<u64> = self.report.iter().map(|level| level.to_wire()).collect();
1550        report.sort_unstable();
1551        report.dedup();
1552        let count = u64::from(self.endpoint.is_some())
1553            + u64::from(self.content_len.is_some())
1554            + u64::from(self.content_type.is_some())
1555            + u64::from(self.traceparent.is_some())
1556            + u64::from(self.tracestate.is_some())
1557            + u64::from(self.topic.is_some())
1558            + u64::from(self.sequence.is_some())
1559            + u64::from(self.producer.is_some())
1560            + u64::from(self.achieved.is_some())
1561            + u64::from(self.report_id.is_some())
1562            + u64::from(!report.is_empty())
1563            + u64::from(self.report_mode != ReportMode::default())
1564            + u64::from(self.segment.is_some())
1565            + u64::from(self.layer.is_some());
1566        encode_into_with(out, |e| {
1567            e.map(count)?;
1568            if let Some(endpoint) = &self.endpoint {
1569                e.u64(data_key::ENDPOINT)?.str(endpoint)?;
1570            }
1571            if let Some(len) = self.content_len {
1572                e.u64(data_key::CONTENT_LEN)?.u64(len)?;
1573            }
1574            if let Some(ct) = &self.content_type {
1575                e.u64(data_key::CONTENT_TYPE)?.str(ct)?;
1576            }
1577            if let Some(tp) = &self.traceparent {
1578                e.u64(data_key::TRACEPARENT)?.str(tp)?;
1579            }
1580            if let Some(ts) = &self.tracestate {
1581                e.u64(data_key::TRACESTATE)?.str(ts)?;
1582            }
1583            if let Some(topic) = &self.topic {
1584                e.u64(data_key::TOPIC)?.str(topic)?;
1585            }
1586            // Keys 6 and 7 are written only when set, which for every v0
1587            // sender means never: nothing in `weida` populates them yet.
1588            if let Some(sequence) = self.sequence {
1589                e.u64(data_key::SEQUENCE)?.u64(sequence)?;
1590            }
1591            if let Some(producer) = &self.producer {
1592                e.u64(data_key::PRODUCER)?.bytes(producer)?;
1593            }
1594            if let Some(achieved) = self.achieved {
1595                e.u64(data_key::ACHIEVED)?.u64(achieved.to_wire())?;
1596            }
1597            if let Some(report_id) = self.report_id {
1598                e.u64(data_key::REPORT_ID)?.u64(report_id)?;
1599            }
1600            if !report.is_empty() {
1601                e.u64(data_key::REPORT)?.array(report.len() as u64)?;
1602                for value in &report {
1603                    e.u64(*value)?;
1604                }
1605            }
1606            // `Progress` is never written: an absent key already says it, so
1607            // a header that orders a report in the default mode stays as
1608            // short as the mode is uninteresting (§6.5's rule for levels).
1609            if self.report_mode != ReportMode::default() {
1610                e.u64(data_key::REPORT_MODE)?
1611                    .u64(self.report_mode.to_wire())?;
1612            }
1613            if let Some(segment) = self.segment {
1614                e.u64(data_key::SEGMENT)?.u64(segment)?;
1615            }
1616            if let Some(layer) = self.layer {
1617                e.u64(data_key::LAYER)?.u64(u64::from(layer))?;
1618            }
1619            Ok(())
1620        })
1621    }
1622
1623    /// Decodes the header.
1624    pub fn decode(bytes: &[u8]) -> Result<DataHeader, HeaderError> {
1625        let mut d = Decoder::new(bytes);
1626        let mut header = DataHeader::default();
1627        {
1628            let mut m = MapReader::new(&mut d)?;
1629            while let Some(key) = m.next_key()? {
1630                match key {
1631                    data_key::ENDPOINT => {
1632                        header.endpoint = Some(m.text(key, limits::MAX_ENDPOINT_BYTES)?)
1633                    }
1634                    data_key::CONTENT_LEN => header.content_len = Some(m.u64()?),
1635                    data_key::CONTENT_TYPE => {
1636                        header.content_type = Some(m.text(key, limits::MAX_CONTENT_TYPE_BYTES)?)
1637                    }
1638                    data_key::TRACEPARENT => {
1639                        header.traceparent = Some(m.text(key, limits::MAX_TRACEPARENT_BYTES)?)
1640                    }
1641                    data_key::TRACESTATE => {
1642                        header.tracestate = Some(m.text(key, limits::MAX_TRACESTATE_BYTES)?)
1643                    }
1644                    data_key::TOPIC => header.topic = Some(m.text(key, limits::MAX_TOPIC_BYTES)?),
1645                    data_key::SEQUENCE => header.sequence = Some(m.u64()?),
1646                    data_key::PRODUCER => header.producer = Some(m.byte_array(key)?),
1647                    // An unknown level is not a level: a peer naming one this
1648                    // version does not define is refused rather than silently
1649                    // read as the weakest, because the value decides what a
1650                    // producer believes about its message.
1651                    data_key::ACHIEVED => {
1652                        let value = m.u64()?;
1653                        header.achieved = Some(Acknowledgement::from_wire(value).ok_or(
1654                            HeaderError::UnknownLevel {
1655                                dimension: "achieved",
1656                                value,
1657                            },
1658                        )?);
1659                    }
1660                    data_key::REPORT_ID => header.report_id = Some(m.u64()?),
1661                    data_key::REPORT => header.report = m.report_levels()?,
1662                    data_key::REPORT_MODE => {
1663                        header.report_mode = level(m.u64()?, "report_mode")?;
1664                    }
1665                    data_key::SEGMENT => header.segment = Some(m.u64()?),
1666                    data_key::LAYER => header.layer = Some(layer(m.u64()?, "layer above 15")?),
1667                    _ => m.skip()?,
1668                }
1669            }
1670        }
1671        finish(&d)?;
1672        // A layer belongs to a segment: key 14 alone names a layer of
1673        // nothing.
1674        if header.layer.is_some() && header.segment.is_none() {
1675            return Err(HeaderError::InvalidLayer("layer without segment"));
1676        }
1677        // Keys 9 and 10 are one statement in two halves: an order with no
1678        // stream to report on, or a stream with nothing to report, names a
1679        // report nobody can serve.
1680        if !header.report.is_empty() && header.report_id.is_none() {
1681            return Err(HeaderError::InvalidReport("report without report_id"));
1682        }
1683        if header.report_id.is_some() && header.report.is_empty() {
1684            return Err(HeaderError::InvalidReport("report_id without report"));
1685        }
1686        Ok(header)
1687    }
1688}
1689
1690/// ERROR header.
1691///
1692/// Legal only on the reply half of a bidirectional stream: an ERROR is the
1693/// alternative to a reply, so it needs no reference to what it answers.
1694#[derive(Clone, Debug, PartialEq, Eq)]
1695pub struct ErrorHeader {
1696    /// Raw error code.
1697    pub code: u64,
1698    /// Human-readable detail; never machine-interpreted.
1699    pub message: Option<String>,
1700}
1701
1702impl ErrorHeader {
1703    /// Builds a header for a known error code.
1704    pub fn new(code: weida_core::ErrorCode) -> ErrorHeader {
1705        ErrorHeader {
1706            code: code.to_wire(),
1707            message: None,
1708        }
1709    }
1710
1711    /// The error code, or `None` for an unknown one.
1712    pub fn error_code(&self) -> Option<weida_core::ErrorCode> {
1713        weida_core::ErrorCode::from_wire(self.code)
1714    }
1715
1716    /// Encodes the header.
1717    pub fn encode(&self) -> Vec<u8> {
1718        let mut out = Vec::new();
1719        self.encode_into(&mut out);
1720        out
1721    }
1722
1723    /// Appends the encoded header to `out`, for a send path that reuses a
1724    /// buffer (B-250). The canonical form has one implementation and this is
1725    /// it; [`Self::encode`] is a wrapper.
1726    pub fn encode_into(&self, out: &mut Vec<u8>) {
1727        let count = 1 + u64::from(self.message.is_some());
1728        encode_into_with(out, |e| {
1729            e.map(count)?;
1730            e.u64(error_key::CODE)?.u64(self.code)?;
1731            if let Some(msg) = &self.message {
1732                e.u64(error_key::MESSAGE)?.str(msg)?;
1733            }
1734            Ok(())
1735        })
1736    }
1737
1738    /// Decodes the header.
1739    pub fn decode(bytes: &[u8]) -> Result<ErrorHeader, HeaderError> {
1740        let mut d = Decoder::new(bytes);
1741        let mut code = 0;
1742        let mut message = None;
1743        {
1744            let mut m = MapReader::new(&mut d)?;
1745            while let Some(key) = m.next_key()? {
1746                match key {
1747                    error_key::CODE => code = m.u64()?,
1748                    error_key::MESSAGE => message = Some(m.text(key, limits::MAX_MESSAGE_BYTES)?),
1749                    _ => m.skip()?,
1750                }
1751            }
1752            m.require(error_key::CODE)?;
1753        }
1754        finish(&d)?;
1755        Ok(ErrorHeader { code, message })
1756    }
1757}
1758
1759/// SUBSCRIBE and UNSUBSCRIBE header.
1760///
1761/// Both frames carry the same two keys: the publisher path to (un)subscribe on
1762/// and the topic filter. The filter is a segmented pattern, not a byte prefix
1763/// ([`filter`]): the empty filter matches every topic, `*` matches one whole
1764/// segment and a trailing `#` matches zero or more.
1765#[derive(Clone, Debug, PartialEq, Eq)]
1766pub struct SubscriptionHeader {
1767    /// Publisher endpoint path.
1768    pub endpoint: String,
1769    /// Topic filter; the empty string matches everything. A decoded header's
1770    /// filter has passed [`filter::validate`].
1771    pub filter: String,
1772    /// The dish's latency budget in milliseconds (key `2`), meaningful on a
1773    /// RADIO path only
1774    /// ([decisions/0034](../../../docs/decisions/0034-late-is-lost.md)
1775    /// §4.6). Encoded only when present.
1776    pub max_age_ms: Option<u64>,
1777    /// The highest layer the dish wants (key `3`), `0..=15`, meaningful on a
1778    /// RADIO path only; absent means no cap
1779    /// ([decisions/0037](../../../docs/decisions/0037-layered-segments.md)
1780    /// §4.3). Encoded only when present.
1781    pub max_layer: Option<u8>,
1782}
1783
1784impl SubscriptionHeader {
1785    /// A header for `endpoint` and `filter`.
1786    pub fn new(endpoint: impl Into<String>, filter: impl Into<String>) -> SubscriptionHeader {
1787        SubscriptionHeader {
1788            endpoint: endpoint.into(),
1789            filter: filter.into(),
1790            max_age_ms: None,
1791            max_layer: None,
1792        }
1793    }
1794
1795    /// Encodes the header.
1796    pub fn encode(&self) -> Vec<u8> {
1797        let mut out = Vec::new();
1798        self.encode_into(&mut out);
1799        out
1800    }
1801
1802    /// Appends the encoded header to `out`, for a send path that reuses a
1803    /// buffer (B-250). The canonical form has one implementation and this is
1804    /// it; [`Self::encode`] is a wrapper.
1805    pub fn encode_into(&self, out: &mut Vec<u8>) {
1806        // Keys 0 and 1 are required, so neither is elided: an absent filter
1807        // and an empty filter would otherwise be indistinguishable on the
1808        // wire, and the empty filter is the "everything" subscription. Keys 2
1809        // and 3 are optional and written only when set.
1810        encode_into_with(out, |e| {
1811            e.map(2 + u64::from(self.max_age_ms.is_some()) + u64::from(self.max_layer.is_some()))?;
1812            e.u64(subscription_key::ENDPOINT)?.str(&self.endpoint)?;
1813            e.u64(subscription_key::FILTER)?.str(&self.filter)?;
1814            if let Some(max_age_ms) = self.max_age_ms {
1815                e.u64(subscription_key::MAX_AGE_MS)?.u64(max_age_ms)?;
1816            }
1817            if let Some(max_layer) = self.max_layer {
1818                e.u64(subscription_key::MAX_LAYER)?
1819                    .u64(u64::from(max_layer))?;
1820            }
1821            Ok(())
1822        })
1823    }
1824
1825    /// Decodes the header.
1826    pub fn decode(bytes: &[u8]) -> Result<SubscriptionHeader, HeaderError> {
1827        let mut d = Decoder::new(bytes);
1828        let mut endpoint = None;
1829        let mut filter = None;
1830        let mut max_age_ms = None;
1831        let mut max_layer = None;
1832        {
1833            let mut m = MapReader::new(&mut d)?;
1834            while let Some(key) = m.next_key()? {
1835                match key {
1836                    subscription_key::ENDPOINT => {
1837                        endpoint = Some(m.text(key, limits::MAX_ENDPOINT_BYTES)?)
1838                    }
1839                    subscription_key::FILTER => {
1840                        filter = Some(m.text(key, limits::MAX_FILTER_BYTES)?)
1841                    }
1842                    subscription_key::MAX_AGE_MS => max_age_ms = Some(m.u64()?),
1843                    subscription_key::MAX_LAYER => {
1844                        max_layer = Some(layer(m.u64()?, "max_layer above 15")?);
1845                    }
1846                    _ => m.skip()?,
1847                }
1848            }
1849            m.require(subscription_key::ENDPOINT)?;
1850            m.require(subscription_key::FILTER)?;
1851        }
1852        // The grammar is checked here, at the codec boundary, so no matcher
1853        // ever sees a filter it would have to interpret twice; an illegal one
1854        // closes the connection with `PROTOCOL_VIOLATION`
1855        // (`docs/PROTOCOL.md` §6.4).
1856        let filter = filter.expect("presence checked above");
1857        filter::validate(&filter)?;
1858        finish(&d)?;
1859        Ok(SubscriptionHeader {
1860            endpoint: endpoint.expect("presence checked above"),
1861            filter,
1862            max_age_ms,
1863            max_layer,
1864        })
1865    }
1866}
1867
1868/// FLOW header (kind `7`).
1869///
1870/// The registration of a datagram flow: the stream it opens is the flow's
1871/// lifetime, and the flow id prefixes every datagram the flow carries
1872/// ([decisions/0034](../../../docs/decisions/0034-late-is-lost.md) §4.2).
1873/// The optional keys reuse DATA's meaning and caps.
1874#[derive(Clone, Debug, PartialEq, Eq)]
1875pub struct FlowHeader {
1876    /// Endpoint path the flow is addressed to (key `0`, required).
1877    pub endpoint: String,
1878    /// Flow id, chosen by the sender and unique per connection and
1879    /// direction (key `1`, required).
1880    pub flow: u64,
1881    /// Opaque media type label (key `2`).
1882    pub content_type: Option<String>,
1883    /// W3C Trace Context `traceparent` (key `3`).
1884    pub traceparent: Option<String>,
1885    /// W3C Trace Context `tracestate` (key `4`).
1886    pub tracestate: Option<String>,
1887    /// Topic the flow carries (key `5`).
1888    pub topic: Option<String>,
1889}
1890
1891impl FlowHeader {
1892    /// A header for flow `flow` addressed to `endpoint`, with no optional key.
1893    pub fn new(endpoint: impl Into<String>, flow: u64) -> FlowHeader {
1894        FlowHeader {
1895            endpoint: endpoint.into(),
1896            flow,
1897            content_type: None,
1898            traceparent: None,
1899            tracestate: None,
1900            topic: None,
1901        }
1902    }
1903
1904    /// Encodes the header.
1905    pub fn encode(&self) -> Vec<u8> {
1906        let mut out = Vec::new();
1907        self.encode_into(&mut out);
1908        out
1909    }
1910
1911    /// Appends the encoded header to `out`. The canonical form has one
1912    /// implementation and this is it; [`Self::encode`] is a wrapper.
1913    pub fn encode_into(&self, out: &mut Vec<u8>) {
1914        let count = 2
1915            + u64::from(self.content_type.is_some())
1916            + u64::from(self.traceparent.is_some())
1917            + u64::from(self.tracestate.is_some())
1918            + u64::from(self.topic.is_some());
1919        encode_into_with(out, |e| {
1920            e.map(count)?;
1921            e.u64(flow_key::ENDPOINT)?.str(&self.endpoint)?;
1922            e.u64(flow_key::FLOW)?.u64(self.flow)?;
1923            if let Some(ct) = &self.content_type {
1924                e.u64(flow_key::CONTENT_TYPE)?.str(ct)?;
1925            }
1926            if let Some(tp) = &self.traceparent {
1927                e.u64(flow_key::TRACEPARENT)?.str(tp)?;
1928            }
1929            if let Some(ts) = &self.tracestate {
1930                e.u64(flow_key::TRACESTATE)?.str(ts)?;
1931            }
1932            if let Some(topic) = &self.topic {
1933                e.u64(flow_key::TOPIC)?.str(topic)?;
1934            }
1935            Ok(())
1936        })
1937    }
1938
1939    /// Decodes the header.
1940    pub fn decode(bytes: &[u8]) -> Result<FlowHeader, HeaderError> {
1941        let mut d = Decoder::new(bytes);
1942        let mut endpoint = None;
1943        let mut flow = None;
1944        let mut content_type = None;
1945        let mut traceparent = None;
1946        let mut tracestate = None;
1947        let mut topic = None;
1948        {
1949            let mut m = MapReader::new(&mut d)?;
1950            while let Some(key) = m.next_key()? {
1951                match key {
1952                    flow_key::ENDPOINT => endpoint = Some(m.text(key, limits::MAX_ENDPOINT_BYTES)?),
1953                    flow_key::FLOW => flow = Some(m.u64()?),
1954                    flow_key::CONTENT_TYPE => {
1955                        content_type = Some(m.text(key, limits::MAX_CONTENT_TYPE_BYTES)?)
1956                    }
1957                    flow_key::TRACEPARENT => {
1958                        traceparent = Some(m.text(key, limits::MAX_TRACEPARENT_BYTES)?)
1959                    }
1960                    flow_key::TRACESTATE => {
1961                        tracestate = Some(m.text(key, limits::MAX_TRACESTATE_BYTES)?)
1962                    }
1963                    flow_key::TOPIC => topic = Some(m.text(key, limits::MAX_TOPIC_BYTES)?),
1964                    _ => m.skip()?,
1965                }
1966            }
1967            m.require(flow_key::ENDPOINT)?;
1968            m.require(flow_key::FLOW)?;
1969        }
1970        finish(&d)?;
1971        Ok(FlowHeader {
1972            endpoint: endpoint.expect("presence checked above"),
1973            flow: flow.expect("presence checked above"),
1974            content_type,
1975            traceparent,
1976            tracestate,
1977            topic,
1978        })
1979    }
1980}
1981
1982/// CREDIT header (kind `5`).
1983///
1984/// The L2 credit of
1985/// [decisions/0003](../../../docs/decisions/0003-credit-unit.md) §4.2-§4.3:
1986/// which subscription, and how many messages that subscription will accept in
1987/// total. Three keys, all required — a subscription is `(endpoint, filter)`
1988/// on the connection the frame arrives on, and an absent limit would be
1989/// indistinguishable from a limit of zero, which is the pause.
1990///
1991/// **The limit is absolute and cumulative, not a delta.** It counts messages
1992/// delivered on that subscription since it was created, so a lost frame costs
1993/// nothing and a duplicated one changes nothing. It is also **monotone at the
1994/// receiver**: a broker keeps the highest limit it has seen, which is what
1995/// makes a reordered frame harmless on a transport that does not order the
1996/// streams control frames ride. Monotone is the whole rule: a receiver
1997/// ignores any limit that is not strictly greater than the standing one, so
1998/// restating a number already delivered changes nothing unless the
1999/// subscription had exhausted its credit anyway. **v0 offers no way to lower
2000/// a standing limit.** The only pause is the `0` a fresh subscription starts
2001/// at, so a consumer that wants to stay in control grants in increments it
2002/// is willing to receive. AMQP 1.0 can shrink `link-credit` against an
2003/// absolute baseline [amqp10 §5.1]; this frame cannot, and a peer that reads
2004/// it as if it could would wait for a stop no broker can deliver.
2005#[derive(Clone, Debug, PartialEq, Eq)]
2006pub struct CreditHeader {
2007    /// Endpoint path of the queue the subscription is on.
2008    pub endpoint: String,
2009    /// Topic filter of the subscription. A decoded header's filter has passed
2010    /// [`filter::validate`].
2011    pub filter: String,
2012    /// Messages this subscription will accept in total, counted from its
2013    /// creation.
2014    pub limit: u64,
2015}
2016
2017impl CreditHeader {
2018    /// A header granting `limit` to the subscription `(endpoint, filter)`.
2019    pub fn new(endpoint: impl Into<String>, filter: impl Into<String>, limit: u64) -> CreditHeader {
2020        CreditHeader {
2021            endpoint: endpoint.into(),
2022            filter: filter.into(),
2023            limit,
2024        }
2025    }
2026
2027    /// Encodes the header.
2028    pub fn encode(&self) -> Vec<u8> {
2029        let mut out = Vec::new();
2030        self.encode_into(&mut out);
2031        out
2032    }
2033
2034    /// Appends the encoded header to `out`, for a send path that reuses a
2035    /// buffer (B-250). The canonical form has one implementation and this is
2036    /// it; [`Self::encode`] is a wrapper.
2037    pub fn encode_into(&self, out: &mut Vec<u8>) {
2038        encode_into_with(out, |e| {
2039            e.map(3)?;
2040            e.u64(credit_key::ENDPOINT)?.str(&self.endpoint)?;
2041            e.u64(credit_key::FILTER)?.str(&self.filter)?;
2042            e.u64(credit_key::LIMIT)?.u64(self.limit)?;
2043            Ok(())
2044        })
2045    }
2046
2047    /// Decodes the header.
2048    pub fn decode(bytes: &[u8]) -> Result<CreditHeader, HeaderError> {
2049        let mut d = Decoder::new(bytes);
2050        let mut endpoint = None;
2051        let mut filter = None;
2052        let mut limit = None;
2053        {
2054            let mut m = MapReader::new(&mut d)?;
2055            while let Some(key) = m.next_key()? {
2056                match key {
2057                    credit_key::ENDPOINT => {
2058                        endpoint = Some(m.text(key, limits::MAX_ENDPOINT_BYTES)?)
2059                    }
2060                    credit_key::FILTER => filter = Some(m.text(key, limits::MAX_FILTER_BYTES)?),
2061                    credit_key::LIMIT => limit = Some(m.u64()?),
2062                    _ => m.skip()?,
2063                }
2064            }
2065            m.require(credit_key::ENDPOINT)?;
2066            m.require(credit_key::FILTER)?;
2067            m.require(credit_key::LIMIT)?;
2068        }
2069        // Same boundary as SUBSCRIBE: a filter that does not name a
2070        // subscription cannot grant credit to one.
2071        let filter = filter.expect("presence checked above");
2072        filter::validate(&filter)?;
2073        finish(&d)?;
2074        Ok(CreditHeader {
2075            endpoint: endpoint.expect("presence checked above"),
2076            filter,
2077            limit: limit.expect("presence checked above"),
2078        })
2079    }
2080}
2081
2082/// CURSOR head frame (kind `6`).
2083///
2084/// One key, required: which report this stream carries. The stream then
2085/// carries `(level, offset)` records until FIN and **never any payload**
2086/// ([decisions/0024](../../../docs/decisions/0024-three-families-one-back-channel.md)
2087/// §4.4) — which is why the head frame needs no length, no endpoint and no
2088/// correlation beyond the id the DATA header allocated.
2089#[derive(Clone, Copy, Debug, PartialEq, Eq)]
2090pub struct CursorHeader {
2091    /// The `report_id` of the DATA header that ordered this report.
2092    pub report_id: u64,
2093}
2094
2095impl CursorHeader {
2096    /// Encodes the header.
2097    pub fn encode(&self) -> Vec<u8> {
2098        let mut out = Vec::new();
2099        self.encode_into(&mut out);
2100        out
2101    }
2102
2103    /// Appends the encoded header to `out`, for a send path that reuses a
2104    /// buffer (B-250). The canonical form has one implementation and this is
2105    /// it; [`Self::encode`] is a wrapper.
2106    pub fn encode_into(&self, out: &mut Vec<u8>) {
2107        encode_into_with(out, |e| {
2108            e.map(1)?;
2109            e.u64(cursor_key::REPORT_ID)?.u64(self.report_id)?;
2110            Ok(())
2111        })
2112    }
2113
2114    /// Decodes the header.
2115    pub fn decode(bytes: &[u8]) -> Result<CursorHeader, HeaderError> {
2116        let mut d = Decoder::new(bytes);
2117        let mut report_id = None;
2118        {
2119            let mut m = MapReader::new(&mut d)?;
2120            while let Some(key) = m.next_key()? {
2121                match key {
2122                    cursor_key::REPORT_ID => report_id = Some(m.u64()?),
2123                    _ => m.skip()?,
2124                }
2125            }
2126            m.require(cursor_key::REPORT_ID)?;
2127        }
2128        finish(&d)?;
2129        Ok(CursorHeader {
2130            report_id: report_id.expect("presence checked above"),
2131        })
2132    }
2133}
2134
2135/// Longest possible cursor record: two 8-byte QUIC varints.
2136pub const MAX_CURSOR_RECORD_LEN: usize = 2 * crate::varint::MAX_ENCODED_LEN;
2137
2138/// Appends one `(level, offset)` record to `out`.
2139///
2140/// Records are QUIC varint pairs rather than CBOR: a record is hot-path and
2141/// self-delimiting, and a CBOR map per record would cost a map header per
2142/// reported byte range for no gain — the head frame already carries every
2143/// field a record could need to name.
2144pub fn encode_cursor_record(
2145    level: CursorLevel,
2146    offset: u64,
2147    out: &mut Vec<u8>,
2148) -> Result<(), VarintError> {
2149    encode_varint(level.to_wire(), out)?;
2150    encode_varint(offset, out)
2151}
2152
2153/// Decodes one record from the front of `input`.
2154///
2155/// `Ok(None)` means the input ends inside a record: read more bytes and
2156/// retry. That is **not** a violation — a reader sees whatever slice the
2157/// transport handed it, and a record is at most
2158/// [`MAX_CURSOR_RECORD_LEN`] bytes, so the retry is bounded. An undefined
2159/// level in the reserved range *is* a violation: the value decides what the
2160/// receiver believes about its own transfer.
2161pub fn decode_cursor_record(
2162    input: &[u8],
2163) -> Result<Option<(CursorLevel, u64, usize)>, HeaderError> {
2164    // `decode_varint` accepts the whole representable range, so its only
2165    // failure is a value the input ended inside of.
2166    let Ok((raw_level, level_len)) = decode_varint(input) else {
2167        return Ok(None);
2168    };
2169    let Ok((offset, offset_len)) = decode_varint(&input[level_len..]) else {
2170        return Ok(None);
2171    };
2172    let level = CursorLevel::from_wire(raw_level).ok_or(HeaderError::UnknownLevel {
2173        dimension: "cursor",
2174        value: raw_level,
2175    })?;
2176    Ok(Some((level, offset, level_len + offset_len)))
2177}
2178
2179/// REPORT head frame (kind `8`): the empty map, every key reserved
2180/// (`docs/PROTOCOL.md` §6.10). Records follow it until FIN.
2181#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
2182pub struct ReportHead;
2183
2184impl ReportHead {
2185    /// Encodes the head frame: `A0`.
2186    pub fn encode(&self) -> Vec<u8> {
2187        let mut out = Vec::new();
2188        encode_into_with(&mut out, |e| {
2189            e.map(0)?;
2190            Ok(())
2191        });
2192        out
2193    }
2194
2195    /// Decodes the head frame; unknown keys are skipped by §5's rule.
2196    pub fn decode(bytes: &[u8]) -> Result<ReportHead, HeaderError> {
2197        let mut d = Decoder::new(bytes);
2198        {
2199            let mut m = MapReader::new(&mut d)?;
2200            while m.next_key()?.is_some() {
2201                m.skip()?;
2202            }
2203        }
2204        finish(&d)?;
2205        Ok(ReportHead)
2206    }
2207}
2208
2209/// REPORT record keys (`docs/PROTOCOL.md` §6.10).
2210mod path_key {
2211    pub const RTT_US: u64 = 0;
2212    pub const MIN_RTT_US: u64 = 1;
2213    pub const CWND: u64 = 2;
2214    pub const CONGESTION_EVENTS: u64 = 3;
2215    pub const LOST_PACKETS: u64 = 4;
2216    pub const LOST_BYTES: u64 = 5;
2217    pub const SENT_PACKETS: u64 = 6;
2218    pub const CURRENT_MTU: u64 = 7;
2219    pub const TX_DATAGRAMS: u64 = 8;
2220    pub const TX_BYTES: u64 = 9;
2221    pub const RX_DATAGRAMS: u64 = 10;
2222    pub const RX_BYTES: u64 = 11;
2223}
2224
2225/// One REPORT record: the sender's view of the QUIC path under the
2226/// connection (`docs/PROTOCOL.md` §6.10,
2227/// [decisions/0036](../../../docs/decisions/0036-connection-statistics.md)
2228/// §4.5). Numbers only; no address travels. A key the peer did not send
2229/// reads as `0`, and the encoder writes only the non-zero ones.
2230#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
2231pub struct PathRecord {
2232    /// Smoothed round-trip time, µs (key `0`).
2233    pub rtt_us: u64,
2234    /// Smallest round-trip time seen, µs (key `1`).
2235    pub min_rtt_us: u64,
2236    /// Congestion window, bytes (key `2`).
2237    pub cwnd: u64,
2238    /// Congestion events the sender's controller reacted to (key `3`).
2239    pub congestion_events: u64,
2240    /// Packets the sender sent and declared lost (key `4`).
2241    pub lost_packets: u64,
2242    /// Bytes the sender sent and declared lost (key `5`).
2243    pub lost_bytes: u64,
2244    /// Packets the sender sent (key `6`).
2245    pub sent_packets: u64,
2246    /// The largest UDP payload the path carries now (key `7`).
2247    pub current_mtu: u64,
2248    /// UDP datagrams the sender sent (key `8`).
2249    pub tx_datagrams: u64,
2250    /// Bytes in them (key `9`).
2251    pub tx_bytes: u64,
2252    /// UDP datagrams the sender received (key `10`).
2253    pub rx_datagrams: u64,
2254    /// Bytes in them (key `11`).
2255    pub rx_bytes: u64,
2256}
2257
2258impl PathRecord {
2259    fn fields(&self) -> [(u64, u64); 12] {
2260        [
2261            (path_key::RTT_US, self.rtt_us),
2262            (path_key::MIN_RTT_US, self.min_rtt_us),
2263            (path_key::CWND, self.cwnd),
2264            (path_key::CONGESTION_EVENTS, self.congestion_events),
2265            (path_key::LOST_PACKETS, self.lost_packets),
2266            (path_key::LOST_BYTES, self.lost_bytes),
2267            (path_key::SENT_PACKETS, self.sent_packets),
2268            (path_key::CURRENT_MTU, self.current_mtu),
2269            (path_key::TX_DATAGRAMS, self.tx_datagrams),
2270            (path_key::TX_BYTES, self.tx_bytes),
2271            (path_key::RX_DATAGRAMS, self.rx_datagrams),
2272            (path_key::RX_BYTES, self.rx_bytes),
2273        ]
2274    }
2275
2276    /// Appends `[varint length][CBOR map]` to `out`. Twelve keys of at most
2277    /// nine bytes each stay far below [`limits::MAX_PATH_RECORD_BYTES`].
2278    pub fn encode_into(&self, out: &mut Vec<u8>) {
2279        let fields = self.fields();
2280        let present = fields.iter().filter(|(_, v)| *v != 0).count() as u64;
2281        let mut map = Vec::new();
2282        encode_into_with(&mut map, |e| {
2283            e.map(present)?;
2284            for (key, value) in fields {
2285                if value != 0 {
2286                    e.u64(key)?.u64(value)?;
2287                }
2288            }
2289            Ok(())
2290        });
2291        encode_varint(map.len() as u64, out).expect("a record is far below 2^62 bytes");
2292        out.extend_from_slice(&map);
2293    }
2294
2295    /// Decodes one record from the front of `input`.
2296    ///
2297    /// `Ok(None)` means the input ends inside a record: read more bytes and
2298    /// retry, which is bounded because a record is at most
2299    /// [`limits::MAX_PATH_RECORD_BYTES`] bytes plus its prefix; at FIN the
2300    /// same answer is a truncated record and a violation. A length prefix
2301    /// above the cap is refused before any of the map is read.
2302    pub fn decode(input: &[u8]) -> Result<Option<(PathRecord, usize)>, HeaderError> {
2303        let Ok((len, prefix)) = decode_varint(input) else {
2304            return Ok(None);
2305        };
2306        if len > limits::MAX_PATH_RECORD_BYTES as u64 {
2307            return Err(HeaderError::InvalidPathReport("record above 256 bytes"));
2308        }
2309        let len = len as usize;
2310        let Some(map) = input.get(prefix..prefix + len) else {
2311            return Ok(None);
2312        };
2313        let mut record = PathRecord::default();
2314        let mut d = Decoder::new(map);
2315        {
2316            let mut m = MapReader::new(&mut d)?;
2317            while let Some(key) = m.next_key()? {
2318                let slot = match key {
2319                    path_key::RTT_US => &mut record.rtt_us,
2320                    path_key::MIN_RTT_US => &mut record.min_rtt_us,
2321                    path_key::CWND => &mut record.cwnd,
2322                    path_key::CONGESTION_EVENTS => &mut record.congestion_events,
2323                    path_key::LOST_PACKETS => &mut record.lost_packets,
2324                    path_key::LOST_BYTES => &mut record.lost_bytes,
2325                    path_key::SENT_PACKETS => &mut record.sent_packets,
2326                    path_key::CURRENT_MTU => &mut record.current_mtu,
2327                    path_key::TX_DATAGRAMS => &mut record.tx_datagrams,
2328                    path_key::TX_BYTES => &mut record.tx_bytes,
2329                    path_key::RX_DATAGRAMS => &mut record.rx_datagrams,
2330                    path_key::RX_BYTES => &mut record.rx_bytes,
2331                    _ => {
2332                        m.skip()?;
2333                        continue;
2334                    }
2335                };
2336                *slot = m.u64()?;
2337            }
2338        }
2339        finish(&d)?;
2340        Ok(Some((record, prefix + len)))
2341    }
2342}
2343
2344#[cfg(test)]
2345mod tests {
2346    /// The tests build headers by hand; production code goes through each
2347    /// type's `encode_into`.
2348    fn encode_with(
2349        f: impl FnOnce(
2350            &mut minicbor::Encoder<Vec<u8>>,
2351        ) -> Result<(), minicbor::encode::Error<std::convert::Infallible>>,
2352    ) -> Vec<u8> {
2353        let mut out = Vec::new();
2354        super::encode_into_with(&mut out, f);
2355        out
2356    }
2357
2358    use super::*;
2359    use weida_core::ErrorCode;
2360
2361    // --- golden vectors, docs/PROTOCOL.md §8 ------------------------------
2362
2363    #[test]
2364    fn golden_data_request_header() {
2365        let h = DataHeader::addressed("/t");
2366        let bytes = h.encode();
2367        assert_eq!(bytes, vec![0xA1, 0x00, 0x62, 0x2F, 0x74]);
2368        assert_eq!(bytes.len(), 0x05);
2369        assert_eq!(DataHeader::decode(&bytes).unwrap(), h);
2370    }
2371
2372    #[test]
2373    fn golden_data_reply_header() {
2374        // The stream is the correlation, so a reply header is an empty map.
2375        let h = DataHeader::reply();
2376        let bytes = h.encode();
2377        assert_eq!(bytes, vec![0xA0]);
2378        assert_eq!(bytes.len(), 0x01);
2379        assert_eq!(DataHeader::decode(&bytes).unwrap(), h);
2380    }
2381
2382    #[test]
2383    fn golden_hello_header() {
2384        let h = Hello::v0(16384, 1024);
2385        let bytes = h.encode();
2386        assert_eq!(
2387            bytes,
2388            vec![
2389                0xA5, 0x00, 0x81, 0x00, 0x01, 0x19, 0x40, 0x00, 0x02, 0x19, 0x04, 0x00, 0x03, 0x80,
2390                0x04, 0x80
2391            ]
2392        );
2393        assert_eq!(bytes.len(), 0x10);
2394        assert_eq!(Hello::decode(&bytes).unwrap(), h);
2395    }
2396
2397    #[test]
2398    fn golden_error_header() {
2399        let h = ErrorHeader::new(ErrorCode::NoReply);
2400        let bytes = h.encode();
2401        assert_eq!(bytes, vec![0xA1, 0x00, 0x05]);
2402        assert_eq!(bytes.len(), 0x03);
2403        assert_eq!(ErrorHeader::decode(&bytes).unwrap(), h);
2404    }
2405
2406    #[test]
2407    fn golden_pub_copy_data_header() {
2408        let mut h = DataHeader::addressed("/md");
2409        h.topic = Some("px.eur".into());
2410        let bytes = h.encode();
2411        assert_eq!(
2412            bytes,
2413            vec![
2414                0xA2, 0x00, 0x63, 0x2F, 0x6D, 0x64, 0x05, 0x66, 0x70, 0x78, 0x2E, 0x65, 0x75, 0x72
2415            ]
2416        );
2417        assert_eq!(bytes.len(), 0x0E);
2418        assert_eq!(DataHeader::decode(&bytes).unwrap(), h);
2419    }
2420
2421    /// The digest of the §8 vectors: SHA-256 of `"test"`, the value the
2422    /// address examples in `docs/PROTOCOL.md` already use.
2423    const VECTOR_PRODUCER: [u8; limits::PRODUCER_BYTES] = [
2424        0x9F, 0x86, 0xD0, 0x81, 0x88, 0x4C, 0x7D, 0x65, 0x9A, 0x2F, 0xEA, 0xA0, 0xC5, 0x5A, 0xD0,
2425        0x15, 0xA3, 0xBF, 0x4F, 0x1B, 0x2B, 0x0B, 0x82, 0x2C, 0xD1, 0x5D, 0x6C, 0x15, 0xB0, 0xF0,
2426        0x0A, 0x08,
2427    ];
2428
2429    #[test]
2430    fn golden_sequenced_data_header() {
2431        let mut h = DataHeader::addressed("/t");
2432        h.sequence = Some(1);
2433        let bytes = h.encode();
2434        assert_eq!(bytes, vec![0xA2, 0x00, 0x62, 0x2F, 0x74, 0x06, 0x01]);
2435        assert_eq!(bytes.len(), 0x07);
2436        assert_eq!(DataHeader::decode(&bytes).unwrap(), h);
2437    }
2438
2439    #[test]
2440    fn golden_relayed_data_header() {
2441        let mut h = DataHeader::addressed("/t");
2442        h.sequence = Some(1);
2443        h.producer = Some(VECTOR_PRODUCER);
2444        let bytes = h.encode();
2445        let mut expected = vec![0xA3, 0x00, 0x62, 0x2F, 0x74, 0x06, 0x01, 0x07, 0x58, 0x20];
2446        expected.extend_from_slice(&VECTOR_PRODUCER);
2447        assert_eq!(bytes, expected);
2448        assert_eq!(bytes.len(), 0x2A);
2449        assert_eq!(DataHeader::decode(&bytes).unwrap(), h);
2450    }
2451
2452    #[test]
2453    fn a_producer_longer_than_the_cap_is_rejected() {
2454        let bytes = encode_with(|e| {
2455            e.map(1)?;
2456            e.u64(data_key::PRODUCER)?
2457                .bytes(&[0u8; limits::PRODUCER_BYTES + 1])?;
2458            Ok(())
2459        });
2460        assert_eq!(
2461            DataHeader::decode(&bytes),
2462            Err(HeaderError::StringTooLong {
2463                key: data_key::PRODUCER,
2464                len: limits::PRODUCER_BYTES + 1,
2465                max: limits::PRODUCER_BYTES,
2466            })
2467        );
2468    }
2469
2470    #[test]
2471    fn a_producer_shorter_than_a_digest_is_rejected() {
2472        // Half a digest identifies nobody, so it is a framing violation
2473        // rather than a value to carry (`docs/PROTOCOL.md` §6.2).
2474        let bytes = encode_with(|e| {
2475            e.map(1)?;
2476            e.u64(data_key::PRODUCER)?.bytes(&[0u8; 16])?;
2477            Ok(())
2478        });
2479        assert!(matches!(
2480            DataHeader::decode(&bytes),
2481            Err(HeaderError::Malformed(_))
2482        ));
2483    }
2484
2485    #[test]
2486    fn the_new_keys_reject_the_wrong_cbor_type() {
2487        let sequence_as_text = encode_with(|e| {
2488            e.map(1)?;
2489            e.u64(data_key::SEQUENCE)?.str("7")?;
2490            Ok(())
2491        });
2492        assert!(DataHeader::decode(&sequence_as_text).is_err());
2493
2494        let producer_as_text = encode_with(|e| {
2495            e.map(1)?;
2496            e.u64(data_key::PRODUCER)?.str("sha256:…")?;
2497            Ok(())
2498        });
2499        assert!(DataHeader::decode(&producer_as_text).is_err());
2500    }
2501
2502    #[test]
2503    fn a_v0_header_carries_neither_new_key() {
2504        // What the runtime writes on a `core` connection: neither key 6 —
2505        // which needs negotiated `PerProducer` ordering — nor key 7, which
2506        // nothing in this repository sets.
2507        let mut h = DataHeader::addressed("/t");
2508        h.traceparent = Some("00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01".into());
2509        let bytes = h.encode();
2510        let mut d = Decoder::new(&bytes);
2511        let pairs = d.map().unwrap().unwrap();
2512        let keys: Vec<u64> = (0..pairs)
2513            .map(|_| {
2514                let key = d.u64().unwrap();
2515                d.skip().unwrap();
2516                key
2517            })
2518            .collect();
2519        assert_eq!(keys, vec![data_key::ENDPOINT, data_key::TRACEPARENT]);
2520    }
2521
2522    #[test]
2523    fn golden_subscription_headers() {
2524        let h = SubscriptionHeader::new("/md", "px.");
2525        let bytes = h.encode();
2526        assert_eq!(
2527            bytes,
2528            vec![
2529                0xA2, 0x00, 0x63, 0x2F, 0x6D, 0x64, 0x01, 0x63, 0x70, 0x78, 0x2E
2530            ]
2531        );
2532        assert_eq!(bytes.len(), 0x0B);
2533        // One header layout serves both kinds; only the kind byte differs, and
2534        // that byte belongs to the preamble (see `tests/golden_vectors.rs`).
2535        assert_eq!(SubscriptionHeader::decode(&bytes).unwrap(), h);
2536    }
2537
2538    // --- roundtrips -------------------------------------------------------
2539
2540    #[test]
2541    fn data_header_roundtrip_with_every_field() {
2542        let h = DataHeader {
2543            endpoint: Some("/transform".into()),
2544            content_len: Some(1 << 40),
2545            content_type: Some("application/octet-stream".into()),
2546            traceparent: Some("00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01".into()),
2547            tracestate: Some("vendor=value".into()),
2548            topic: Some("px.eur".into()),
2549            sequence: Some(u64::MAX),
2550            producer: Some([0x5A; limits::PRODUCER_BYTES]),
2551            achieved: Some(Acknowledgement::Processed),
2552            report_id: Some(7),
2553            report: vec![
2554                CursorLevel::Known(Acknowledgement::Accepted),
2555                CursorLevel::Known(Acknowledgement::Processed),
2556                CursorLevel::Application(CursorLevel::APPLICATION_FLOOR),
2557            ],
2558            report_mode: ReportMode::FinalOnly,
2559            segment: Some(u64::MAX),
2560            layer: Some(limits::MAX_LAYER),
2561        };
2562        assert_eq!(DataHeader::decode(&h.encode()).unwrap(), h);
2563    }
2564
2565    #[test]
2566    fn error_header_roundtrip_with_and_without_message() {
2567        let bare = ErrorHeader::new(ErrorCode::UnknownEndpoint);
2568        assert_eq!(ErrorHeader::decode(&bare.encode()).unwrap(), bare);
2569        assert_eq!(bare.error_code(), Some(ErrorCode::UnknownEndpoint));
2570
2571        let with_msg = ErrorHeader {
2572            code: 4,
2573            message: Some("handler panicked".into()),
2574        };
2575        assert_eq!(ErrorHeader::decode(&with_msg.encode()).unwrap(), with_msg);
2576    }
2577
2578    #[test]
2579    fn keys_are_emitted_in_ascending_order() {
2580        let h = DataHeader {
2581            endpoint: Some("/x".into()),
2582            content_len: Some(1),
2583            content_type: Some("t".into()),
2584            traceparent: Some("p".into()),
2585            tracestate: Some("s".into()),
2586            topic: Some("k".into()),
2587            sequence: Some(9),
2588            producer: Some([0u8; limits::PRODUCER_BYTES]),
2589            achieved: Some(Acknowledgement::Accepted),
2590            report_id: Some(1),
2591            report: vec![CursorLevel::Known(Acknowledgement::Stored)],
2592            report_mode: ReportMode::FinalOnly,
2593            segment: Some(3),
2594            layer: Some(2),
2595        };
2596        let bytes = h.encode();
2597        let mut d = Decoder::new(&bytes);
2598        let n = d.map().unwrap().unwrap();
2599        let mut last = None;
2600        for _ in 0..n {
2601            let key = d.u64().unwrap();
2602            if let Some(prev) = last {
2603                assert!(key > prev, "keys must ascend: {prev} then {key}");
2604            }
2605            last = Some(key);
2606            d.skip().unwrap();
2607        }
2608    }
2609
2610    #[test]
2611    fn an_unsorted_report_with_repeats_is_emitted_as_the_canonical_ascending_set() {
2612        // A caller hands over the levels in the order it thought of them.
2613        // Decoding is the proof: the decoder refuses a non-ascending or
2614        // repeated array, so a header that survives its own encoder was
2615        // canonicalized on the way out.
2616        let h = DataHeader {
2617            report_id: Some(1),
2618            report: vec![
2619                CursorLevel::Application(CursorLevel::APPLICATION_FLOOR),
2620                CursorLevel::Known(Acknowledgement::Stored),
2621                CursorLevel::Application(CursorLevel::APPLICATION_FLOOR),
2622                CursorLevel::Known(Acknowledgement::Accepted),
2623                CursorLevel::Known(Acknowledgement::Stored),
2624            ],
2625            ..DataHeader::reply()
2626        };
2627        let decoded = DataHeader::decode(&h.encode()).expect("encoder emits the canonical form");
2628        assert_eq!(
2629            decoded.report,
2630            vec![
2631                CursorLevel::Known(Acknowledgement::Accepted),
2632                CursorLevel::Known(Acknowledgement::Stored),
2633                CursorLevel::Application(CursorLevel::APPLICATION_FLOOR),
2634            ]
2635        );
2636    }
2637
2638    // --- optional fields --------------------------------------------------
2639
2640    #[test]
2641    fn every_data_field_is_optional_at_the_decoder() {
2642        // The endpoint requirement lives in dispatch, not here: a reply half
2643        // legitimately carries none, and the decoder cannot tell the halves
2644        // apart.
2645        assert_eq!(DataHeader::decode(&[0xA0]).unwrap(), DataHeader::default());
2646
2647        let only_topic = encode_with(|e| {
2648            e.map(1)?;
2649            e.u64(data_key::TOPIC)?.str("px.eur")?;
2650            Ok(())
2651        });
2652        let h = DataHeader::decode(&only_topic).unwrap();
2653        assert_eq!(h.topic.as_deref(), Some("px.eur"));
2654        assert_eq!(h.endpoint, None);
2655    }
2656
2657    #[test]
2658    fn absent_fields_are_omitted_by_the_encoder() {
2659        let h = DataHeader::addressed("/t");
2660        assert_eq!(h.encode(), vec![0xA1, 0x00, 0x62, 0x2F, 0x74]);
2661    }
2662
2663    // --- forward compatibility -------------------------------------------
2664
2665    #[test]
2666    fn unknown_keys_are_skipped() {
2667        // Re-encode the golden DATA header with an extra key 63 holding a
2668        // nested structure, and check it still decodes to the same value.
2669        let h = DataHeader::addressed("/t");
2670        let extended = encode_with(|e| {
2671            e.map(2)?;
2672            e.u64(0)?.str("/t")?;
2673            e.u64(63)?.array(2)?.u64(7)?.map(1)?.u64(1)?.bool(true)?;
2674            Ok(())
2675        });
2676        assert_eq!(DataHeader::decode(&extended).unwrap(), h);
2677    }
2678
2679    #[test]
2680    fn unknown_keys_above_the_reserved_range_are_skipped() {
2681        let extended = encode_with(|e| {
2682            e.map(2)?;
2683            e.u64(1)?.u64(5)?;
2684            e.u64(1000)?.str("future")?;
2685            Ok(())
2686        });
2687        let h = DataHeader::decode(&extended).unwrap();
2688        assert_eq!(h.content_len, Some(5));
2689    }
2690
2691    #[test]
2692    fn skipping_tolerates_nesting_up_to_the_depth_limit() {
2693        for depth in [1usize, limits::MAX_SKIP_DEPTH] {
2694            let bytes = encode_with(|e| {
2695                e.map(2)?;
2696                e.u64(data_key::CONTENT_LEN)?.u64(1)?;
2697                e.u64(50)?;
2698                for _ in 0..depth {
2699                    e.array(1)?;
2700                }
2701                e.u64(1)?;
2702                Ok(())
2703            });
2704            let h = DataHeader::decode(&bytes).unwrap_or_else(|e| panic!("depth {depth}: {e}"));
2705            assert_eq!(h.content_len, Some(1), "depth {depth}");
2706        }
2707    }
2708
2709    #[test]
2710    fn skipping_rejects_nesting_beyond_the_depth_limit() {
2711        let bytes = encode_with(|e| {
2712            e.map(1)?;
2713            e.u64(50)?;
2714            for _ in 0..(limits::MAX_SKIP_DEPTH + 1) {
2715                e.array(1)?;
2716            }
2717            e.u64(1)?;
2718            Ok(())
2719        });
2720        assert_eq!(
2721            DataHeader::decode(&bytes).unwrap_err(),
2722            HeaderError::DepthExceeded
2723        );
2724    }
2725
2726    #[test]
2727    fn skipping_a_wide_shallow_structure_is_fine() {
2728        let bytes = encode_with(|e| {
2729            e.map(2)?;
2730            e.u64(data_key::CONTENT_LEN)?.u64(1)?;
2731            e.u64(40)?.array(64)?;
2732            for i in 0..64u64 {
2733                e.u64(i)?;
2734            }
2735            Ok(())
2736        });
2737        assert_eq!(DataHeader::decode(&bytes).unwrap().content_len, Some(1));
2738    }
2739
2740    // --- strictness -------------------------------------------------------
2741
2742    #[test]
2743    fn duplicate_keys_are_rejected() {
2744        let bytes = encode_with(|e| {
2745            e.map(2)?;
2746            e.u64(1)?.u64(1)?;
2747            e.u64(1)?.u64(2)?;
2748            Ok(())
2749        });
2750        assert_eq!(
2751            DataHeader::decode(&bytes).unwrap_err(),
2752            HeaderError::DuplicateKey(1)
2753        );
2754    }
2755
2756    #[test]
2757    fn non_uint_keys_are_rejected() {
2758        let bytes = encode_with(|e| {
2759            e.map(1)?;
2760            e.str("endpoint")?.str("/t")?;
2761            Ok(())
2762        });
2763        assert_eq!(
2764            DataHeader::decode(&bytes).unwrap_err(),
2765            HeaderError::NonUintKey
2766        );
2767
2768        let negative = encode_with(|e| {
2769            e.map(1)?;
2770            e.i64(-1)?.u64(1)?;
2771            Ok(())
2772        });
2773        assert_eq!(
2774            DataHeader::decode(&negative).unwrap_err(),
2775            HeaderError::NonUintKey
2776        );
2777    }
2778
2779    #[test]
2780    fn indefinite_maps_are_rejected() {
2781        let bytes = encode_with(|e| {
2782            e.begin_map()?;
2783            e.u64(1)?.u64(1)?;
2784            e.end()?;
2785            Ok(())
2786        });
2787        assert_eq!(
2788            DataHeader::decode(&bytes).unwrap_err(),
2789            HeaderError::Indefinite
2790        );
2791    }
2792
2793    #[test]
2794    fn indefinite_arrays_are_rejected() {
2795        let bytes = encode_with(|e| {
2796            e.map(5)?;
2797            e.u64(0)?.begin_array()?.u64(0)?.end()?;
2798            e.u64(1)?.u64(1)?;
2799            e.u64(2)?.u64(1)?;
2800            e.u64(3)?.array(0)?;
2801            e.u64(4)?.array(0)?;
2802            Ok(())
2803        });
2804        assert_eq!(Hello::decode(&bytes).unwrap_err(), HeaderError::Indefinite);
2805    }
2806
2807    #[test]
2808    fn value_type_mismatches_are_rejected() {
2809        let bytes = encode_with(|e| {
2810            e.map(1)?;
2811            e.u64(data_key::CONTENT_LEN)?.str("not a number")?;
2812            Ok(())
2813        });
2814        assert!(matches!(
2815            DataHeader::decode(&bytes).unwrap_err(),
2816            HeaderError::Malformed(_)
2817        ));
2818    }
2819
2820    #[test]
2821    fn missing_required_keys_are_rejected() {
2822        // ERROR without a code.
2823        let bytes = encode_with(|e| {
2824            e.map(1)?;
2825            e.u64(error_key::MESSAGE)?.str("why")?;
2826            Ok(())
2827        });
2828        assert_eq!(
2829            ErrorHeader::decode(&bytes).unwrap_err(),
2830            HeaderError::MissingKey(error_key::CODE)
2831        );
2832
2833        // HELLO missing capabilities.
2834        let bytes = encode_with(|e| {
2835            e.map(4)?;
2836            e.u64(0)?.array(1)?.u64(0)?;
2837            e.u64(1)?.u64(16384)?;
2838            e.u64(2)?.u64(16)?;
2839            e.u64(4)?.array(0)?;
2840            Ok(())
2841        });
2842        assert_eq!(
2843            Hello::decode(&bytes).unwrap_err(),
2844            HeaderError::MissingKey(hello_key::CAPABILITIES)
2845        );
2846    }
2847
2848    // --- subscription headers ---------------------------------------------
2849
2850    #[test]
2851    fn an_empty_filter_is_legal_and_survives_the_roundtrip() {
2852        let h = SubscriptionHeader::new("/md", "");
2853        let bytes = h.encode();
2854        assert_eq!(SubscriptionHeader::decode(&bytes).unwrap(), h);
2855        // The key is written even though the value is empty: absent and empty
2856        // must stay distinguishable, and empty means "every topic".
2857        assert!(
2858            bytes.contains(&0x60),
2859            "the empty filter is encoded: {bytes:?}"
2860        );
2861    }
2862
2863    #[test]
2864    fn the_filter_grammar_accepts_what_docs_protocol_6_4_permits() {
2865        for ok in [
2866            "",
2867            "#",
2868            "px",
2869            "px.eur",
2870            "px.*",
2871            "*.eur",
2872            "sensors.*.temp",
2873            "px.#",
2874            "px.",
2875            "a..b",
2876        ] {
2877            assert_eq!(filter::validate(ok), Ok(()), "{ok:?} must be legal");
2878        }
2879    }
2880
2881    #[test]
2882    fn the_filter_grammar_rejects_partial_and_misplaced_wildcards() {
2883        for bad in [
2884            "px*", "*px", "p*x.eur", "px.e*ur", "#.px", "px.#.eur", "px#",
2885        ] {
2886            assert!(
2887                matches!(filter::validate(bad), Err(HeaderError::InvalidFilter(_))),
2888                "{bad:?} must be rejected"
2889            );
2890        }
2891    }
2892
2893    #[test]
2894    fn an_illegal_filter_is_rejected_at_the_codec_boundary() {
2895        // Encoding does not validate — a test may build any bytes — but
2896        // decoding does, which is what makes the grammar enforceable against
2897        // a peer (`docs/PROTOCOL.md` §6.4).
2898        let bytes = SubscriptionHeader::new("/md", "px.#.eur").encode();
2899        assert!(matches!(
2900            SubscriptionHeader::decode(&bytes),
2901            Err(HeaderError::InvalidFilter(_))
2902        ));
2903        let e: Error = HeaderError::InvalidFilter("`#` must be the final segment").into();
2904        assert!(e.to_string().contains("invalid topic filter"));
2905    }
2906
2907    #[test]
2908    fn subscription_strings_are_capped() {
2909        for (key, max) in [
2910            (subscription_key::ENDPOINT, limits::MAX_ENDPOINT_BYTES),
2911            (subscription_key::FILTER, limits::MAX_FILTER_BYTES),
2912        ] {
2913            let build = |len: usize| {
2914                let text = "a".repeat(len);
2915                let mut h = SubscriptionHeader::new("/md", "px.");
2916                if key == subscription_key::ENDPOINT {
2917                    h.endpoint = text;
2918                } else {
2919                    h.filter = text;
2920                }
2921                h.encode()
2922            };
2923            assert_eq!(
2924                SubscriptionHeader::decode(&build(max + 1)).unwrap_err(),
2925                HeaderError::StringTooLong {
2926                    key,
2927                    len: max + 1,
2928                    max
2929                },
2930                "key {key}"
2931            );
2932            assert!(
2933                SubscriptionHeader::decode(&build(max)).is_ok(),
2934                "key {key} at cap"
2935            );
2936        }
2937    }
2938
2939    #[test]
2940    fn subscription_headers_require_both_keys() {
2941        let only = |key: u64| {
2942            encode_with(|e| {
2943                e.map(1)?;
2944                e.u64(key)?.str("/md")?;
2945                Ok(())
2946            })
2947        };
2948        assert_eq!(
2949            SubscriptionHeader::decode(&only(subscription_key::ENDPOINT)).unwrap_err(),
2950            HeaderError::MissingKey(subscription_key::FILTER)
2951        );
2952        assert_eq!(
2953            SubscriptionHeader::decode(&only(subscription_key::FILTER)).unwrap_err(),
2954            HeaderError::MissingKey(subscription_key::ENDPOINT)
2955        );
2956    }
2957
2958    #[test]
2959    fn subscription_headers_reject_malformed_input() {
2960        assert!(SubscriptionHeader::decode(&[]).is_err());
2961        // Trailing bytes.
2962        let mut bytes = SubscriptionHeader::new("/md", "px.").encode();
2963        bytes.push(0xff);
2964        assert_eq!(
2965            SubscriptionHeader::decode(&bytes).unwrap_err(),
2966            HeaderError::TrailingBytes
2967        );
2968        // Unknown keys are skipped, like every other header.
2969        let extended = encode_with(|e| {
2970            e.map(3)?;
2971            e.u64(0)?.str("/md")?;
2972            e.u64(1)?.str("px.")?;
2973            e.u64(40)?.array(2)?.u64(1)?.u64(2)?;
2974            Ok(())
2975        });
2976        assert_eq!(
2977            SubscriptionHeader::decode(&extended).unwrap(),
2978            SubscriptionHeader::new("/md", "px.")
2979        );
2980    }
2981
2982    #[test]
2983    fn oversized_strings_are_rejected_per_field() {
2984        // Built through the encoder, which emits keys in ascending order.
2985        let with_text = |key: u64, text: String| -> Vec<u8> {
2986            let mut h = DataHeader::reply();
2987            match key {
2988                data_key::ENDPOINT => h.endpoint = Some(text),
2989                data_key::CONTENT_TYPE => h.content_type = Some(text),
2990                data_key::TRACEPARENT => h.traceparent = Some(text),
2991                data_key::TRACESTATE => h.tracestate = Some(text),
2992                data_key::TOPIC => h.topic = Some(text),
2993                other => panic!("key {other} is not a text field"),
2994            }
2995            h.encode()
2996        };
2997        let cases: [(u64, usize); 5] = [
2998            (data_key::ENDPOINT, limits::MAX_ENDPOINT_BYTES),
2999            (data_key::CONTENT_TYPE, limits::MAX_CONTENT_TYPE_BYTES),
3000            (data_key::TRACEPARENT, limits::MAX_TRACEPARENT_BYTES),
3001            (data_key::TRACESTATE, limits::MAX_TRACESTATE_BYTES),
3002            (data_key::TOPIC, limits::MAX_TOPIC_BYTES),
3003        ];
3004        for (key, max) in cases {
3005            assert_eq!(
3006                DataHeader::decode(&with_text(key, "a".repeat(max + 1))).unwrap_err(),
3007                HeaderError::StringTooLong {
3008                    key,
3009                    len: max + 1,
3010                    max
3011                },
3012                "key {key}"
3013            );
3014            assert!(
3015                DataHeader::decode(&with_text(key, "a".repeat(max))).is_ok(),
3016                "key {key} at cap"
3017            );
3018        }
3019    }
3020
3021    #[test]
3022    fn unordered_keys_are_rejected() {
3023        // Descending keys break the ascending-order rule, which is what makes
3024        // duplicate detection complete for extension keys.
3025        let bytes = encode_with(|e| {
3026            e.map(3)?;
3027            e.u64(2)?.str("t")?;
3028            e.u64(1)?.u64(1)?;
3029            e.u64(3)?.str("p")?;
3030            Ok(())
3031        });
3032        assert_eq!(
3033            DataHeader::decode(&bytes).unwrap_err(),
3034            HeaderError::UnorderedKey(1)
3035        );
3036    }
3037
3038    #[test]
3039    fn duplicate_extension_keys_are_rejected() {
3040        let bytes = encode_with(|e| {
3041            e.map(3)?;
3042            e.u64(1)?.u64(1)?;
3043            e.u64(1000)?.u64(1)?;
3044            e.u64(1000)?.u64(2)?;
3045            Ok(())
3046        });
3047        assert_eq!(
3048            DataHeader::decode(&bytes).unwrap_err(),
3049            HeaderError::DuplicateKey(1000)
3050        );
3051    }
3052
3053    #[test]
3054    fn oversized_error_messages_are_rejected() {
3055        let big = "m".repeat(limits::MAX_MESSAGE_BYTES + 1);
3056        let bytes = encode_with(|e| {
3057            e.map(2)?;
3058            e.u64(error_key::CODE)?.u64(2)?;
3059            e.u64(error_key::MESSAGE)?.str(&big)?;
3060            Ok(())
3061        });
3062        assert_eq!(
3063            ErrorHeader::decode(&bytes).unwrap_err(),
3064            HeaderError::StringTooLong {
3065                key: error_key::MESSAGE,
3066                len: limits::MAX_MESSAGE_BYTES + 1,
3067                max: limits::MAX_MESSAGE_BYTES
3068            }
3069        );
3070    }
3071
3072    #[test]
3073    fn oversized_lists_are_rejected_without_allocating() {
3074        // A one-byte array header claiming 2^32 items must not reserve memory.
3075        let bytes = encode_with(|e| {
3076            e.map(1)?;
3077            e.u64(0)?.array(u64::from(u32::MAX))?;
3078            Ok(())
3079        });
3080        assert_eq!(
3081            Hello::decode(&bytes).unwrap_err(),
3082            HeaderError::ListTooLong {
3083                key: hello_key::VERSIONS,
3084                len: u64::from(u32::MAX),
3085                max: limits::MAX_LIST_ITEMS
3086            }
3087        );
3088    }
3089
3090    #[test]
3091    fn lists_exactly_at_the_cap_are_accepted() {
3092        let bytes = encode_with(|e| {
3093            e.map(5)?;
3094            e.u64(0)?.array(limits::MAX_LIST_ITEMS as u64)?;
3095            for i in 0..limits::MAX_LIST_ITEMS as u64 {
3096                e.u64(i)?;
3097            }
3098            e.u64(1)?.u64(16384)?;
3099            e.u64(2)?.u64(16)?;
3100            e.u64(3)?.array(0)?;
3101            e.u64(4)?.array(0)?;
3102            Ok(())
3103        });
3104        assert_eq!(
3105            Hello::decode(&bytes).unwrap().versions.len(),
3106            limits::MAX_LIST_ITEMS
3107        );
3108    }
3109
3110    #[test]
3111    fn trailing_bytes_are_rejected() {
3112        let mut bytes = ErrorHeader::new(ErrorCode::Rejected).encode();
3113        bytes.push(0xff);
3114        assert_eq!(
3115            ErrorHeader::decode(&bytes).unwrap_err(),
3116            HeaderError::TrailingBytes
3117        );
3118    }
3119
3120    #[test]
3121    fn truncated_headers_are_rejected() {
3122        let full = DataHeader::addressed("/t").encode();
3123        for cut in 0..full.len() {
3124            assert!(
3125                DataHeader::decode(&full[..cut]).is_err(),
3126                "prefix of {cut} bytes must not decode"
3127            );
3128        }
3129    }
3130
3131    #[test]
3132    fn empty_input_is_rejected_for_every_header() {
3133        assert!(Hello::decode(&[]).is_err());
3134        assert!(DataHeader::decode(&[]).is_err());
3135        assert!(ErrorHeader::decode(&[]).is_err());
3136        assert!(SubscriptionHeader::decode(&[]).is_err());
3137    }
3138
3139    #[test]
3140    fn tags_are_rejected() {
3141        // Key 50 is unknown, so the value goes through `skip_value`, which is
3142        // where the tag rule lives.
3143        let bytes = encode_with(|e| {
3144            e.map(1)?;
3145            e.u64(50)?.tag(minicbor::data::IanaTag::Cbor)?.u64(1)?;
3146            Ok(())
3147        });
3148        assert_eq!(
3149            DataHeader::decode(&bytes).unwrap_err(),
3150            HeaderError::Malformed("tags are not allowed")
3151        );
3152    }
3153
3154    // --- reserved value passthrough ---------------------------------------
3155
3156    #[test]
3157    fn unknown_error_codes_survive_decoding() {
3158        let err = ErrorHeader {
3159            code: 99,
3160            message: None,
3161        };
3162        let decoded = ErrorHeader::decode(&err.encode()).unwrap();
3163        assert_eq!(decoded.code, 99);
3164        assert_eq!(decoded.error_code(), None);
3165    }
3166
3167    #[test]
3168    fn header_errors_become_protocol_errors() {
3169        let e: Error = HeaderError::DuplicateKey(3).into();
3170        assert!(matches!(e, Error::Protocol(_)));
3171        assert!(e.to_string().contains("duplicate header key 3"));
3172    }
3173}