Skip to main content

moqtap_codec/draft18/
fields.rs

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