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    let mut obj = Map::new();
140    let mut unknown = Vec::new();
141
142    for p in params {
143        let key = p.key.into_inner();
144        if let Some(name) = d17_param_name(key) {
145            match (&p.value, key) {
146                (KvpValue::Bytes(b), 0x21) => {
147                    obj.insert(name.to_string(), decode_subscription_filter(b));
148                }
149                (KvpValue::Bytes(b), 0x09) => {
150                    obj.insert(name.to_string(), decode_largest_object(b));
151                }
152                (KvpValue::Bytes(b), _) if name == "authorization_token" => {
153                    obj.insert(name.to_string(), auth_token_to_json_d17(b));
154                }
155                (KvpValue::Varint(v), _) => {
156                    obj.insert(name.to_string(), vi(v.into_inner()));
157                }
158                (KvpValue::Bytes(b), _) => {
159                    obj.insert(
160                        name.to_string(),
161                        Value::Text(String::from_utf8_lossy(b).into_owned()),
162                    );
163                }
164            }
165        } else {
166            let mut entry = Map::new();
167            entry.insert("id".to_string(), Value::Text(format!("0x{:x}", key)));
168            match &p.value {
169                KvpValue::Varint(v) => {
170                    entry.insert("length".to_string(), vi(v.into_inner()));
171                }
172                KvpValue::Bytes(b) => {
173                    entry.insert("length".to_string(), vi(b.len() as u64));
174                    entry.insert("raw_hex".to_string(), Value::Bytes(b.to_vec()));
175                }
176            }
177            unknown.push(Value::Map(entry));
178        }
179    }
180
181    if !unknown.is_empty() {
182        obj.insert("unknown".to_string(), Value::Array(unknown));
183    }
184
185    Value::Map(obj)
186}
187
188fn options_to_json(options: &[KeyValuePair]) -> Value {
189    let mut obj = Map::new();
190    for p in options {
191        let key = p.key.into_inner();
192        if let Some(name) = d17_option_name(key) {
193            match &p.value {
194                KvpValue::Varint(v) => {
195                    obj.insert(name.to_string(), vi(v.into_inner()));
196                }
197                KvpValue::Bytes(b) if name == "authorization_token" => {
198                    obj.insert(name.to_string(), auth_token_to_json_d17(b));
199                }
200                KvpValue::Bytes(b) => {
201                    obj.insert(
202                        name.to_string(),
203                        Value::Text(String::from_utf8_lossy(b).into_owned()),
204                    );
205                }
206            }
207        }
208    }
209    Value::Map(obj)
210}
211
212fn d17_track_prop_name(key: u64) -> Option<&'static str> {
213    match key {
214        0x02 => Some("delivery_timeout"),
215        0x04 => Some("max_cache_duration"),
216        0x0b => Some("immutable_properties"),
217        0x0e => Some("default_publisher_priority"),
218        0x22 => Some("default_publisher_group_order"),
219        0x30 => Some("dynamic_groups"),
220        _ => None,
221    }
222}
223
224fn track_props_to_json(props: &[KeyValuePair]) -> Value {
225    let mut obj = Map::new();
226    for p in props {
227        let key = p.key.into_inner();
228        let name = d17_track_prop_name(key)
229            .map(|s| s.to_string())
230            .unwrap_or_else(|| format!("0x{:x}", key));
231        match &p.value {
232            KvpValue::Varint(v) => {
233                obj.insert(name, vi(v.into_inner()));
234            }
235            KvpValue::Bytes(b) => {
236                obj.insert(name, Value::Bytes(b.to_vec()));
237            }
238        }
239    }
240    Value::Map(obj)
241}
242
243/// This draft's field names for a decoded control message.
244///
245/// Keys are the names this draft gives its fields, in the order it defines
246/// them. An optional field the message did not carry is absent rather than
247/// zero.
248pub fn message_fields(msg: &ControlMessage) -> Map {
249    let obj = match msg {
250        ControlMessage::Setup(m) => {
251            let mut o = Map::new();
252            o.insert("options".into(), options_to_json(&m.options));
253            o
254        }
255        ControlMessage::GoAway(m) => {
256            let mut o = Map::new();
257            o.insert(
258                "new_session_uri".into(),
259                Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
260            );
261            o.insert("timeout".into(), vi(m.timeout.into_inner()));
262            o
263        }
264        ControlMessage::RequestOk(m) => {
265            let mut o = Map::new();
266            o.insert("parameters".into(), params_to_json(&m.parameters));
267            o
268        }
269        ControlMessage::RequestError(m) => {
270            let mut o = Map::new();
271            o.insert("error_code".into(), vi(m.error_code.into_inner()));
272            o.insert("retry_interval".into(), vi(m.retry_interval.into_inner()));
273            o.insert(
274                "reason_phrase".into(),
275                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
276            );
277            o
278        }
279        ControlMessage::Subscribe(m) => {
280            let mut o = Map::new();
281            o.insert("request_id".into(), vi(m.request_id.into_inner()));
282            o.insert(
283                "required_request_id_delta".into(),
284                vi(m.required_request_id_delta.into_inner()),
285            );
286            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
287            o.insert(
288                "track_name".into(),
289                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
290            );
291            o.insert("parameters".into(), params_to_json(&m.parameters));
292            o
293        }
294        ControlMessage::SubscribeOk(m) => {
295            let mut o = Map::new();
296            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
297            o.insert("parameters".into(), params_to_json(&m.parameters));
298            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
299            o
300        }
301        ControlMessage::RequestUpdate(m) => {
302            let mut o = Map::new();
303            o.insert("request_id".into(), vi(m.request_id.into_inner()));
304            o.insert(
305                "required_request_id_delta".into(),
306                vi(m.required_request_id_delta.into_inner()),
307            );
308            o.insert("parameters".into(), params_to_json(&m.parameters));
309            o
310        }
311        ControlMessage::Publish(m) => {
312            let mut o = Map::new();
313            o.insert("request_id".into(), vi(m.request_id.into_inner()));
314            o.insert(
315                "required_request_id_delta".into(),
316                vi(m.required_request_id_delta.into_inner()),
317            );
318            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
319            o.insert(
320                "track_name".into(),
321                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
322            );
323            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
324            o.insert("parameters".into(), params_to_json(&m.parameters));
325            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
326            o
327        }
328        ControlMessage::PublishOk(m) => {
329            let mut o = Map::new();
330            o.insert("parameters".into(), params_to_json(&m.parameters));
331            o
332        }
333        ControlMessage::PublishDone(m) => {
334            let mut o = Map::new();
335            o.insert("status_code".into(), vi(m.status_code.into_inner()));
336            o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
337            o.insert(
338                "reason_phrase".into(),
339                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
340            );
341            o
342        }
343        ControlMessage::PublishNamespace(m) => {
344            let mut o = Map::new();
345            o.insert("request_id".into(), vi(m.request_id.into_inner()));
346            o.insert(
347                "required_request_id_delta".into(),
348                vi(m.required_request_id_delta.into_inner()),
349            );
350            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
351            o.insert("parameters".into(), params_to_json(&m.parameters));
352            o
353        }
354        ControlMessage::Namespace(m) => {
355            let mut o = Map::new();
356            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
357            o
358        }
359        ControlMessage::NamespaceDone(m) => {
360            let mut o = Map::new();
361            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
362            o
363        }
364        ControlMessage::SubscribeNamespace(m) => {
365            let mut o = Map::new();
366            o.insert("request_id".into(), vi(m.request_id.into_inner()));
367            o.insert(
368                "required_request_id_delta".into(),
369                vi(m.required_request_id_delta.into_inner()),
370            );
371            o.insert("namespace_prefix".into(), ns_to_json(&m.namespace_prefix));
372            o.insert("subscribe_options".into(), vi(m.subscribe_options.into_inner()));
373            o.insert("parameters".into(), params_to_json(&m.parameters));
374            o
375        }
376        ControlMessage::TrackStatus(m) => {
377            let mut o = Map::new();
378            o.insert("request_id".into(), vi(m.request_id.into_inner()));
379            o.insert(
380                "required_request_id_delta".into(),
381                vi(m.required_request_id_delta.into_inner()),
382            );
383            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
384            o.insert(
385                "track_name".into(),
386                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
387            );
388            o.insert("parameters".into(), params_to_json(&m.parameters));
389            o
390        }
391        ControlMessage::Fetch(m) => {
392            let mut o = Map::new();
393            o.insert("request_id".into(), vi(m.request_id.into_inner()));
394            o.insert(
395                "required_request_id_delta".into(),
396                vi(m.required_request_id_delta.into_inner()),
397            );
398            o.insert("fetch_type".into(), vi(m.fetch_type as u64));
399            match &m.fetch_payload {
400                crate::draft17::message::FetchPayload::Standalone {
401                    track_namespace,
402                    track_name,
403                    start_group,
404                    start_object,
405                    end_group,
406                    end_object,
407                } => {
408                    o.insert("track_namespace".into(), ns_to_json(track_namespace));
409                    o.insert(
410                        "track_name".into(),
411                        Value::Text(String::from_utf8_lossy(track_name).into_owned()),
412                    );
413                    o.insert("start_group".into(), vi(start_group.into_inner()));
414                    o.insert("start_object".into(), vi(start_object.into_inner()));
415                    o.insert("end_group".into(), vi(end_group.into_inner()));
416                    o.insert("end_object".into(), vi(end_object.into_inner()));
417                }
418                crate::draft17::message::FetchPayload::Joining {
419                    joining_request_id,
420                    joining_start,
421                } => {
422                    o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
423                    o.insert("joining_start".into(), vi(joining_start.into_inner()));
424                }
425            }
426            o.insert("parameters".into(), params_to_json(&m.parameters));
427            o
428        }
429        ControlMessage::FetchOk(m) => {
430            let mut o = Map::new();
431            o.insert("end_of_track".into(), vi(m.end_of_track as u64));
432            o.insert("end_group".into(), vi(m.end_group.into_inner()));
433            o.insert("end_object".into(), vi(m.end_object.into_inner()));
434            o.insert("parameters".into(), params_to_json(&m.parameters));
435            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
436            o
437        }
438        ControlMessage::PublishBlocked(m) => {
439            let mut o = Map::new();
440            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
441            o.insert(
442                "track_name".into(),
443                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
444            );
445            o
446        }
447    };
448    obj
449}