Skip to main content

moqtap_codec/draft16/
fields.rs

1use crate::draft16::message::ControlMessage;
2use crate::fields::{FieldMap as Map, FieldValue as Value};
3use crate::kvp::{KeyValuePair, KvpValue};
4use crate::types::*;
5use crate::varint::VarInt;
6
7fn vi(v: u64) -> Value {
8    Value::Uint(v)
9}
10
11fn ns_to_json(ns: &TrackNamespace) -> Value {
12    Value::Array(
13        ns.0.iter().map(|e| Value::Text(String::from_utf8_lossy(e).into_owned())).collect(),
14    )
15}
16
17fn d16_setup_param_name(key: u64) -> Option<&'static str> {
18    match key {
19        0x01 => Some("path"),
20        0x02 => Some("max_request_id"),
21        0x03 => Some("authorization_token"),
22        0x04 => Some("max_auth_token_cache_size"),
23        0x05 => Some("authority"),
24        0x07 => Some("moqt_implementation"),
25        _ => None,
26    }
27}
28
29fn d16_msg_param_name(key: u64) -> Option<&'static str> {
30    match key {
31        0x02 => Some("delivery_timeout"),
32        0x03 => Some("authorization_token"),
33        0x04 => Some("max_cache_duration"),
34        0x08 => Some("expires"),
35        0x09 => Some("largest_object"),
36        0x0e => Some("publisher_priority"),
37        0x10 => Some("forward"),
38        0x20 => Some("subscriber_priority"),
39        0x21 => Some("subscription_filter"),
40        0x22 => Some("group_order"),
41        0x30 => Some("dynamic_groups"),
42        0x32 => Some("new_group_request"),
43        _ => None,
44    }
45}
46
47fn auth_token_to_json_d16(bytes: &[u8]) -> Value {
48    let mut buf = bytes;
49    let alias_type = match VarInt::decode(&mut buf) {
50        Ok(v) => v,
51        Err(_) => return Value::Bytes(bytes.to_vec()),
52    };
53    let at = alias_type.into_inner();
54    let mut o = Map::new();
55    o.insert("alias_type".into(), vi(at));
56    match at {
57        0 | 2 => {
58            if let Ok(ta) = VarInt::decode(&mut buf) {
59                o.insert("token_alias".into(), vi(ta.into_inner()));
60            }
61        }
62        1 => {
63            if let Ok(ta) = VarInt::decode(&mut buf) {
64                o.insert("token_alias".into(), vi(ta.into_inner()));
65            }
66            if let Ok(tt) = VarInt::decode(&mut buf) {
67                o.insert("token_type".into(), vi(tt.into_inner()));
68            }
69            o.insert("token_value".into(), Value::Bytes(buf.to_vec()));
70        }
71        _ => {
72            if let Ok(tt) = VarInt::decode(&mut buf) {
73                o.insert("token_type".into(), vi(tt.into_inner()));
74            }
75            o.insert("token_value".into(), Value::Bytes(buf.to_vec()));
76        }
77    }
78    Value::Map(o)
79}
80
81fn decode_subscription_filter(bytes: &[u8]) -> Value {
82    let mut buf = bytes;
83    let filter_type = VarInt::decode(&mut buf).unwrap().into_inner();
84    let mut obj = Map::new();
85    obj.insert("filter_type".into(), vi(filter_type));
86    match filter_type {
87        3 => {
88            let start_group = VarInt::decode(&mut buf).unwrap().into_inner();
89            let start_object = VarInt::decode(&mut buf).unwrap().into_inner();
90            obj.insert("start_group".into(), vi(start_group));
91            obj.insert("start_object".into(), vi(start_object));
92        }
93        4 => {
94            let start_group = VarInt::decode(&mut buf).unwrap().into_inner();
95            let start_object = VarInt::decode(&mut buf).unwrap().into_inner();
96            let end_group = VarInt::decode(&mut buf).unwrap().into_inner();
97            obj.insert("start_group".into(), vi(start_group));
98            obj.insert("start_object".into(), vi(start_object));
99            obj.insert("end_group".into(), vi(end_group));
100        }
101        _ => {}
102    }
103    Value::Map(obj)
104}
105
106/// Render a draft-16 LARGEST_OBJECT (0x09) parameter value: a Group and an
107/// Object, as two varints.
108///
109/// # Nothing has checked that the value is two varints
110///
111/// Drafts 17 and later give 0x09 a `Location` encoding: their decoders read the
112/// two varints and re-serialise them into the stored value, so what reaches
113/// their extractor is two varints by construction. Draft-16 has no such table.
114/// 0x09 is an odd Type, so `KeyValuePair::decode` keeps whatever
115/// length-prefixed bytes arrived, and `decode_parameters_in` checks duplicates,
116/// authorization tokens, varint value ranges and subscription filters — none of
117/// which looks at 0x09. `KNOWN_MESSAGE_PARAMETERS` admits it and
118/// `check_parameter_scope` permits it on SUBSCRIBE_OK, so a short SUBSCRIBE_OK
119/// carrying `0x09` with an empty value reaches here.
120///
121/// An empty value fails the first read and a single `0x00` fails the *second*,
122/// which is the nastier of the two: the first varint decodes cleanly and the
123/// value looks well formed right up to the point where it is not.
124///
125/// # What a value it cannot read renders as
126///
127/// The raw bytes, as `fields::params` and this file's own
128/// `auth_token_to_json_d16` do. Field extraction runs on a message that has
129/// already decoded, so it has no refusal to give: what a peer sent is what
130/// there is to show.
131fn decode_largest_object(bytes: &[u8]) -> Value {
132    let mut buf = bytes;
133    let Ok(group) = VarInt::decode(&mut buf) else {
134        return Value::Bytes(bytes.to_vec());
135    };
136    let Ok(object) = VarInt::decode(&mut buf) else {
137        return Value::Bytes(bytes.to_vec());
138    };
139    let mut obj = Map::new();
140    obj.insert("group".into(), vi(group.into_inner()));
141    obj.insert("object".into(), vi(object.into_inner()));
142    Value::Map(obj)
143}
144
145fn kvp_to_json_d16_inner(
146    params: &[KeyValuePair],
147    name_fn: fn(u64) -> Option<&'static str>,
148) -> Value {
149    crate::fields::kvp_entries(params, |key, value| {
150        let Some(name) = name_fn(key) else {
151            return (None, None);
152        };
153        let rendered = match (value, key) {
154            (KvpValue::Bytes(b), 0x21) => decode_subscription_filter(b),
155            (KvpValue::Bytes(b), 0x09) => decode_largest_object(b),
156            (KvpValue::Bytes(b), _) if name == "authorization_token" => auth_token_to_json_d16(b),
157            (KvpValue::Varint(v), _) => vi(v.into_inner()),
158            (KvpValue::Bytes(b), _) => Value::Text(String::from_utf8_lossy(b).into_owned()),
159        };
160        (Some(name), Some(rendered))
161    })
162}
163
164fn kvp_to_json_d16(params: &[KeyValuePair]) -> Value {
165    kvp_to_json_d16_inner(params, d16_msg_param_name)
166}
167
168fn kvp_to_json_d16_setup(params: &[KeyValuePair]) -> Value {
169    kvp_to_json_d16_inner(params, d16_setup_param_name)
170}
171
172/// Track-extension parameter names. `default_publisher_priority` shares key
173/// 0x0e with the message-level `publisher_priority`, so a separate table is
174/// used when rendering `track_extensions` blocks.
175fn d16_track_ext_name(key: u64) -> Option<&'static str> {
176    match key {
177        0x02 => Some("delivery_timeout"),
178        0x04 => Some("max_cache_duration"),
179        0x0e => Some("default_publisher_priority"),
180        _ => None,
181    }
182}
183
184fn kvp_to_json_d16_track_ext(params: &[KeyValuePair]) -> Value {
185    kvp_to_json_d16_inner(params, d16_track_ext_name)
186}
187
188/// This draft's field names for a decoded control message.
189///
190/// Keys are the names this draft gives its fields, in the order it defines
191/// them. An optional field the message did not carry is absent rather than
192/// zero.
193pub fn message_fields(msg: &ControlMessage) -> Map {
194    let obj = match msg {
195        ControlMessage::ClientSetup(m) => {
196            let mut o = Map::new();
197            o.insert("parameters".into(), kvp_to_json_d16_setup(&m.parameters));
198            o
199        }
200        ControlMessage::ServerSetup(m) => {
201            let mut o = Map::new();
202            o.insert("parameters".into(), kvp_to_json_d16_setup(&m.parameters));
203            o
204        }
205        ControlMessage::GoAway(m) => {
206            let mut o = Map::new();
207            o.insert(
208                "new_session_uri".into(),
209                Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
210            );
211            o
212        }
213        ControlMessage::MaxRequestId(m) => {
214            let mut o = Map::new();
215            o.insert("max_request_id".into(), vi(m.request_id.into_inner()));
216            o
217        }
218        ControlMessage::RequestsBlocked(m) => {
219            let mut o = Map::new();
220            o.insert("maximum_request_id".into(), vi(m.maximum_request_id.into_inner()));
221            o
222        }
223        ControlMessage::RequestOk(m) => {
224            let mut o = Map::new();
225            o.insert("request_id".into(), vi(m.request_id.into_inner()));
226            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
227            o
228        }
229        ControlMessage::RequestError(m) => {
230            let mut o = Map::new();
231            o.insert("request_id".into(), vi(m.request_id.into_inner()));
232            o.insert("error_code".into(), vi(m.error_code.into_inner()));
233            o.insert("retry_interval".into(), vi(m.retry_interval.into_inner()));
234            o.insert(
235                "reason_phrase".into(),
236                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
237            );
238            o
239        }
240        ControlMessage::Subscribe(m) => {
241            let mut o = Map::new();
242            o.insert("request_id".into(), vi(m.request_id.into_inner()));
243            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
244            o.insert(
245                "track_name".into(),
246                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
247            );
248            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
249            o
250        }
251        ControlMessage::SubscribeOk(m) => {
252            let mut o = Map::new();
253            o.insert("request_id".into(), vi(m.request_id.into_inner()));
254            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
255            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
256            if !m.track_extensions.is_empty() {
257                o.insert("track_extensions".into(), kvp_to_json_d16_track_ext(&m.track_extensions));
258            }
259            o
260        }
261        ControlMessage::RequestUpdate(m) => {
262            let mut o = Map::new();
263            o.insert("request_id".into(), vi(m.request_id.into_inner()));
264            o.insert("existing_request_id".into(), vi(m.existing_request_id.into_inner()));
265            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
266            o
267        }
268        ControlMessage::Unsubscribe(m) => {
269            let mut o = Map::new();
270            o.insert("request_id".into(), vi(m.request_id.into_inner()));
271            o
272        }
273        ControlMessage::Publish(m) => {
274            let mut o = Map::new();
275            o.insert("request_id".into(), vi(m.request_id.into_inner()));
276            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
277            o.insert(
278                "track_name".into(),
279                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
280            );
281            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
282            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
283            if !m.track_extensions.is_empty() {
284                o.insert("track_extensions".into(), kvp_to_json_d16_track_ext(&m.track_extensions));
285            }
286            o
287        }
288        ControlMessage::PublishOk(m) => {
289            let mut o = Map::new();
290            o.insert("request_id".into(), vi(m.request_id.into_inner()));
291            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
292            o
293        }
294        ControlMessage::PublishDone(m) => {
295            let mut o = Map::new();
296            o.insert("request_id".into(), vi(m.request_id.into_inner()));
297            o.insert("status_code".into(), vi(m.status_code.into_inner()));
298            o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
299            o.insert(
300                "reason_phrase".into(),
301                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
302            );
303            o
304        }
305        ControlMessage::PublishNamespace(m) => {
306            let mut o = Map::new();
307            o.insert("request_id".into(), vi(m.request_id.into_inner()));
308            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
309            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
310            o
311        }
312        ControlMessage::PublishNamespaceDone(m) => {
313            let mut o = Map::new();
314            o.insert("request_id".into(), vi(m.request_id.into_inner()));
315            o
316        }
317        ControlMessage::PublishNamespaceCancel(m) => {
318            let mut o = Map::new();
319            o.insert("request_id".into(), vi(m.request_id.into_inner()));
320            o.insert("error_code".into(), vi(m.error_code.into_inner()));
321            o.insert(
322                "reason_phrase".into(),
323                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
324            );
325            o
326        }
327        ControlMessage::Namespace(m) => {
328            let mut o = Map::new();
329            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
330            o
331        }
332        ControlMessage::NamespaceDone(m) => {
333            let mut o = Map::new();
334            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
335            o
336        }
337        ControlMessage::SubscribeNamespace(m) => {
338            let mut o = Map::new();
339            o.insert("request_id".into(), vi(m.request_id.into_inner()));
340            o.insert("namespace_prefix".into(), ns_to_json(&m.namespace_prefix));
341            o.insert("subscribe_options".into(), vi(m.subscribe_options.into_inner()));
342            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
343            o
344        }
345        ControlMessage::TrackStatus(m) => {
346            let mut o = Map::new();
347            o.insert("request_id".into(), vi(m.request_id.into_inner()));
348            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
349            o.insert(
350                "track_name".into(),
351                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
352            );
353            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
354            o
355        }
356        ControlMessage::Fetch(m) => {
357            let mut o = Map::new();
358            o.insert("request_id".into(), vi(m.request_id.into_inner()));
359            o.insert("fetch_type".into(), vi(m.fetch_type as u64));
360            match &m.fetch_payload {
361                crate::draft16::message::FetchPayload::Standalone {
362                    track_namespace,
363                    track_name,
364                    start_group,
365                    start_object,
366                    end_group,
367                    end_object,
368                } => {
369                    o.insert("track_namespace".into(), ns_to_json(track_namespace));
370                    o.insert(
371                        "track_name".into(),
372                        Value::Text(String::from_utf8_lossy(track_name).into_owned()),
373                    );
374                    o.insert("start_group".into(), vi(start_group.into_inner()));
375                    o.insert("start_object".into(), vi(start_object.into_inner()));
376                    o.insert("end_group".into(), vi(end_group.into_inner()));
377                    o.insert("end_object".into(), vi(end_object.into_inner()));
378                }
379                crate::draft16::message::FetchPayload::Joining {
380                    joining_request_id,
381                    joining_start,
382                } => {
383                    o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
384                    o.insert("joining_start".into(), vi(joining_start.into_inner()));
385                }
386            }
387            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
388            o
389        }
390        ControlMessage::FetchOk(m) => {
391            let mut o = Map::new();
392            o.insert("request_id".into(), vi(m.request_id.into_inner()));
393            o.insert("end_of_track".into(), vi(m.end_of_track as u64));
394            o.insert("end_group".into(), vi(m.end_group.into_inner()));
395            o.insert("end_object".into(), vi(m.end_object.into_inner()));
396            o.insert("parameters".into(), kvp_to_json_d16(&m.parameters));
397            if !m.track_extensions.is_empty() {
398                o.insert("track_extensions".into(), kvp_to_json_d16_track_ext(&m.track_extensions));
399            }
400            o
401        }
402        ControlMessage::FetchCancel(m) => {
403            let mut o = Map::new();
404            o.insert("request_id".into(), vi(m.request_id.into_inner()));
405            o
406        }
407    };
408    obj
409}