Skip to main content

moqtap_codec/draft17/
fields.rs

1use crate::draft17::message::ControlMessage;
2use crate::fields::{FieldMap as Map, FieldValue as Value};
3use crate::kvp::{KeyValuePair, KvpValue};
4use crate::types::*;
5use crate::varint::{Moqt17 as Wire, 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
17// Draft-17 known parameter types and their encodings
18fn d17_param_name(key: u64) -> Option<&'static str> {
19    match key {
20        0x02 => Some("delivery_timeout"),
21        0x03 => Some("authorization_token"),
22        0x04 => Some("rendezvous_timeout"),
23        0x08 => Some("expires"),
24        0x09 => Some("largest_object"),
25        0x10 => Some("forward"),
26        0x20 => Some("subscriber_priority"),
27        0x21 => Some("subscription_filter"),
28        0x22 => Some("group_order"),
29        0x32 => Some("new_group_request"),
30        _ => None,
31    }
32}
33
34// Draft-17 setup option names
35fn d17_option_name(key: u64) -> Option<&'static str> {
36    match key {
37        0x01 => Some("path"),
38        0x03 => Some("authorization_token"),
39        0x04 => Some("max_auth_token_cache_size"),
40        0x05 => Some("authority"),
41        0x07 => Some("moqt_implementation"),
42        _ => None,
43    }
44}
45
46fn decode_subscription_filter(bytes: &[u8]) -> Value {
47    let mut buf = bytes;
48    let filter_type = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
49    let mut obj = Map::new();
50    obj.insert("filter_type".into(), vi(filter_type));
51    match filter_type {
52        3 => {
53            let start_group = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
54            let start_object = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
55            obj.insert("start_group".into(), vi(start_group));
56            obj.insert("start_object".into(), vi(start_object));
57        }
58        4 => {
59            let start_group = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
60            let start_object = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
61            let end_group = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
62            obj.insert("start_group".into(), vi(start_group));
63            obj.insert("start_object".into(), vi(start_object));
64            obj.insert("end_group".into(), vi(end_group));
65        }
66        _ => {}
67    }
68    Value::Map(obj)
69}
70
71fn auth_token_to_json_d17(bytes: &[u8]) -> Value {
72    let mut buf = bytes;
73    let alias_type = match VarInt::decode_moqt::<Wire>(&mut buf) {
74        Ok(v) => v,
75        Err(_) => return Value::Bytes(bytes.to_vec()),
76    };
77    let at = alias_type.into_inner();
78    let mut o = Map::new();
79    o.insert("alias_type".into(), vi(at));
80    match at {
81        0 | 2 => {
82            if let Ok(ta) = VarInt::decode_moqt::<Wire>(&mut buf) {
83                o.insert("token_alias".into(), vi(ta.into_inner()));
84            }
85        }
86        1 => {
87            if let Ok(ta) = VarInt::decode_moqt::<Wire>(&mut buf) {
88                o.insert("token_alias".into(), vi(ta.into_inner()));
89            }
90            if let Ok(tt) = VarInt::decode_moqt::<Wire>(&mut buf) {
91                o.insert("token_type".into(), vi(tt.into_inner()));
92            }
93            // Draft-17: token_value is length-prefixed.
94            let tv = match VarInt::decode_moqt::<Wire>(&mut buf) {
95                Ok(len) => {
96                    let n = len.into_inner() as usize;
97                    if buf.len() >= n {
98                        &buf[..n]
99                    } else {
100                        buf
101                    }
102                }
103                Err(_) => buf,
104            };
105            o.insert("token_value".into(), Value::Bytes(tv.to_vec()));
106        }
107        _ => {
108            if let Ok(tt) = VarInt::decode_moqt::<Wire>(&mut buf) {
109                o.insert("token_type".into(), vi(tt.into_inner()));
110            }
111            let tv = match VarInt::decode_moqt::<Wire>(&mut buf) {
112                Ok(len) => {
113                    let n = len.into_inner() as usize;
114                    if buf.len() >= n {
115                        &buf[..n]
116                    } else {
117                        buf
118                    }
119                }
120                Err(_) => buf,
121            };
122            o.insert("token_value".into(), Value::Bytes(tv.to_vec()));
123        }
124    }
125    Value::Map(o)
126}
127
128fn decode_largest_object(bytes: &[u8]) -> Value {
129    let mut buf = bytes;
130    let group = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
131    let object = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
132    let mut obj = Map::new();
133    obj.insert("group".into(), vi(group));
134    obj.insert("object".into(), vi(object));
135    Value::Map(obj)
136}
137
138fn params_to_json(params: &[KeyValuePair]) -> Value {
139    crate::fields::kvp_entries(params, |key, value| {
140        let Some(name) = d17_param_name(key) else {
141            return (None, None);
142        };
143        let rendered = match (value, key) {
144            (KvpValue::Bytes(b), 0x21) => decode_subscription_filter(b),
145            (KvpValue::Bytes(b), 0x09) => decode_largest_object(b),
146            (KvpValue::Bytes(b), _) if name == "authorization_token" => auth_token_to_json_d17(b),
147            (KvpValue::Varint(v), _) => vi(v.into_inner()),
148            (KvpValue::Bytes(b), _) => Value::Text(String::from_utf8_lossy(b).into_owned()),
149        };
150        (Some(name), Some(rendered))
151    })
152}
153
154fn options_to_json(options: &[KeyValuePair]) -> Value {
155    crate::fields::kvp_entries(options, |key, value| {
156        let Some(name) = d17_option_name(key) else {
157            return (None, None);
158        };
159        let rendered = match value {
160            KvpValue::Varint(v) => vi(v.into_inner()),
161            KvpValue::Bytes(b) if name == "authorization_token" => auth_token_to_json_d17(b),
162            KvpValue::Bytes(b) => Value::Text(String::from_utf8_lossy(b).into_owned()),
163        };
164        (Some(name), Some(rendered))
165    })
166}
167
168fn d17_track_prop_name(key: u64) -> Option<&'static str> {
169    match key {
170        0x02 => Some("delivery_timeout"),
171        0x04 => Some("max_cache_duration"),
172        0x0b => Some("immutable_properties"),
173        0x0e => Some("default_publisher_priority"),
174        0x22 => Some("default_publisher_group_order"),
175        0x30 => Some("dynamic_groups"),
176        _ => None,
177    }
178}
179
180fn track_props_to_json(props: &[KeyValuePair]) -> Value {
181    crate::fields::kvp_entries(props, |key, value| {
182        let name = d17_track_prop_name(key);
183        let rendered = match value {
184            KvpValue::Varint(v) => Some(vi(v.into_inner())),
185            // A property this draft does not name keeps its bytes rather than
186            // a name invented from its type, which is what an entry's absent
187            // `name` already says.
188            KvpValue::Bytes(_) => None,
189        };
190        (name, rendered)
191    })
192}
193
194/// This draft's field names for a decoded control message.
195///
196/// Keys are the names this draft gives its fields, in the order it defines
197/// them. An optional field the message did not carry is absent rather than
198/// zero.
199pub fn message_fields(msg: &ControlMessage) -> Map {
200    let obj = match msg {
201        ControlMessage::Setup(m) => {
202            let mut o = Map::new();
203            o.insert("options".into(), options_to_json(&m.options));
204            o
205        }
206        ControlMessage::GoAway(m) => {
207            let mut o = Map::new();
208            o.insert(
209                "new_session_uri".into(),
210                Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
211            );
212            o.insert("timeout".into(), vi(m.timeout.into_inner()));
213            o
214        }
215        ControlMessage::RequestOk(m) => {
216            let mut o = Map::new();
217            o.insert("parameters".into(), params_to_json(&m.parameters));
218            o
219        }
220        ControlMessage::RequestError(m) => {
221            let mut o = Map::new();
222            o.insert("error_code".into(), vi(m.error_code.into_inner()));
223            o.insert("retry_interval".into(), vi(m.retry_interval.into_inner()));
224            o.insert(
225                "reason_phrase".into(),
226                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
227            );
228            o
229        }
230        ControlMessage::Subscribe(m) => {
231            let mut o = Map::new();
232            o.insert("request_id".into(), vi(m.request_id.into_inner()));
233            o.insert(
234                "required_request_id_delta".into(),
235                vi(m.required_request_id_delta.into_inner()),
236            );
237            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
238            o.insert(
239                "track_name".into(),
240                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
241            );
242            o.insert("parameters".into(), params_to_json(&m.parameters));
243            o
244        }
245        ControlMessage::SubscribeOk(m) => {
246            let mut o = Map::new();
247            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
248            o.insert("parameters".into(), params_to_json(&m.parameters));
249            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
250            o
251        }
252        ControlMessage::RequestUpdate(m) => {
253            let mut o = Map::new();
254            o.insert("request_id".into(), vi(m.request_id.into_inner()));
255            o.insert(
256                "required_request_id_delta".into(),
257                vi(m.required_request_id_delta.into_inner()),
258            );
259            o.insert("parameters".into(), params_to_json(&m.parameters));
260            o
261        }
262        ControlMessage::Publish(m) => {
263            let mut o = Map::new();
264            o.insert("request_id".into(), vi(m.request_id.into_inner()));
265            o.insert(
266                "required_request_id_delta".into(),
267                vi(m.required_request_id_delta.into_inner()),
268            );
269            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
270            o.insert(
271                "track_name".into(),
272                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
273            );
274            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
275            o.insert("parameters".into(), params_to_json(&m.parameters));
276            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
277            o
278        }
279        ControlMessage::PublishOk(m) => {
280            let mut o = Map::new();
281            o.insert("parameters".into(), params_to_json(&m.parameters));
282            o
283        }
284        ControlMessage::PublishDone(m) => {
285            let mut o = Map::new();
286            o.insert("status_code".into(), vi(m.status_code.into_inner()));
287            o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
288            o.insert(
289                "reason_phrase".into(),
290                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
291            );
292            o
293        }
294        ControlMessage::PublishNamespace(m) => {
295            let mut o = Map::new();
296            o.insert("request_id".into(), vi(m.request_id.into_inner()));
297            o.insert(
298                "required_request_id_delta".into(),
299                vi(m.required_request_id_delta.into_inner()),
300            );
301            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
302            o.insert("parameters".into(), params_to_json(&m.parameters));
303            o
304        }
305        ControlMessage::Namespace(m) => {
306            let mut o = Map::new();
307            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
308            o
309        }
310        ControlMessage::NamespaceDone(m) => {
311            let mut o = Map::new();
312            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
313            o
314        }
315        ControlMessage::SubscribeNamespace(m) => {
316            let mut o = Map::new();
317            o.insert("request_id".into(), vi(m.request_id.into_inner()));
318            o.insert(
319                "required_request_id_delta".into(),
320                vi(m.required_request_id_delta.into_inner()),
321            );
322            o.insert("namespace_prefix".into(), ns_to_json(&m.namespace_prefix));
323            o.insert("subscribe_options".into(), vi(m.subscribe_options.into_inner()));
324            o.insert("parameters".into(), params_to_json(&m.parameters));
325            o
326        }
327        ControlMessage::TrackStatus(m) => {
328            let mut o = Map::new();
329            o.insert("request_id".into(), vi(m.request_id.into_inner()));
330            o.insert(
331                "required_request_id_delta".into(),
332                vi(m.required_request_id_delta.into_inner()),
333            );
334            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
335            o.insert(
336                "track_name".into(),
337                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
338            );
339            o.insert("parameters".into(), params_to_json(&m.parameters));
340            o
341        }
342        ControlMessage::Fetch(m) => {
343            let mut o = Map::new();
344            o.insert("request_id".into(), vi(m.request_id.into_inner()));
345            o.insert(
346                "required_request_id_delta".into(),
347                vi(m.required_request_id_delta.into_inner()),
348            );
349            o.insert("fetch_type".into(), vi(m.fetch_type as u64));
350            match &m.fetch_payload {
351                crate::draft17::message::FetchPayload::Standalone {
352                    track_namespace,
353                    track_name,
354                    start_group,
355                    start_object,
356                    end_group,
357                    end_object,
358                } => {
359                    o.insert("track_namespace".into(), ns_to_json(track_namespace));
360                    o.insert(
361                        "track_name".into(),
362                        Value::Text(String::from_utf8_lossy(track_name).into_owned()),
363                    );
364                    o.insert("start_group".into(), vi(start_group.into_inner()));
365                    o.insert("start_object".into(), vi(start_object.into_inner()));
366                    o.insert("end_group".into(), vi(end_group.into_inner()));
367                    o.insert("end_object".into(), vi(end_object.into_inner()));
368                }
369                crate::draft17::message::FetchPayload::Joining {
370                    joining_request_id,
371                    joining_start,
372                } => {
373                    o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
374                    o.insert("joining_start".into(), vi(joining_start.into_inner()));
375                }
376            }
377            o.insert("parameters".into(), params_to_json(&m.parameters));
378            o
379        }
380        ControlMessage::FetchOk(m) => {
381            let mut o = Map::new();
382            o.insert("end_of_track".into(), vi(m.end_of_track as u64));
383            o.insert("end_group".into(), vi(m.end_group.into_inner()));
384            o.insert("end_object".into(), vi(m.end_object.into_inner()));
385            o.insert("parameters".into(), params_to_json(&m.parameters));
386            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
387            o
388        }
389        ControlMessage::PublishBlocked(m) => {
390            let mut o = Map::new();
391            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
392            o.insert(
393                "track_name".into(),
394                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
395            );
396            o
397        }
398    };
399    obj
400}