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    crate::fields::kvp_entries(params, |key, value| {
129        let Some(name) = d18_param_name(key) else {
130            return (None, None);
131        };
132        let rendered = match (value, key) {
133            (KvpValue::Bytes(b), 0x21) => decode_subscription_filter(b),
134            (KvpValue::Bytes(b), 0x09) => decode_largest_object(b),
135            (KvpValue::Bytes(b), 0x34) => decode_track_namespace_prefix(b),
136            (KvpValue::Bytes(b), _) if name == "authorization_token" => auth_token_to_json_d18(b),
137            (KvpValue::Varint(v), _) => vi(v.into_inner()),
138            (KvpValue::Bytes(b), _) => Value::Text(String::from_utf8_lossy(b).into_owned()),
139        };
140        (Some(name), Some(rendered))
141    })
142}
143
144fn options_to_json(options: &[KeyValuePair]) -> Value {
145    crate::fields::kvp_entries(options, |key, value| {
146        let Some(name) = d18_option_name(key) else {
147            return (None, None);
148        };
149        let rendered = match value {
150            KvpValue::Varint(v) => vi(v.into_inner()),
151            KvpValue::Bytes(b) if name == "authorization_token" => auth_token_to_json_d18(b),
152            KvpValue::Bytes(b) => Value::Text(String::from_utf8_lossy(b).into_owned()),
153        };
154        (Some(name), Some(rendered))
155    })
156}
157
158fn d18_track_prop_name(key: u64) -> Option<&'static str> {
159    match key {
160        0x02 => Some("object_delivery_timeout"),
161        0x04 => Some("max_cache_duration"),
162        0x06 => Some("subgroup_delivery_timeout"),
163        0x0b => Some("immutable_properties"),
164        0x0e => Some("default_publisher_priority"),
165        0x22 => Some("default_publisher_group_order"),
166        0x30 => Some("dynamic_groups"),
167        _ => None,
168    }
169}
170
171fn track_props_to_json(props: &[KeyValuePair]) -> Value {
172    crate::fields::kvp_entries(props, |key, value| {
173        let name = d18_track_prop_name(key);
174        let rendered = match value {
175            KvpValue::Varint(v) => Some(vi(v.into_inner())),
176            // A property this draft does not name keeps its bytes rather than
177            // a name invented from its type, which is what an entry's absent
178            // `name` already says.
179            KvpValue::Bytes(_) => None,
180        };
181        (name, rendered)
182    })
183}
184
185/// This draft's field names for a decoded control message.
186///
187/// Keys are the names this draft gives its fields, in the order it defines
188/// them. An optional field the message did not carry is absent rather than
189/// zero.
190pub fn message_fields(msg: &ControlMessage) -> Map {
191    let obj = match msg {
192        ControlMessage::Setup(m) => {
193            let mut o = Map::new();
194            o.insert("options".into(), options_to_json(&m.options));
195            o
196        }
197        ControlMessage::GoAway(m) => {
198            let mut o = Map::new();
199            o.insert(
200                "new_session_uri".into(),
201                Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
202            );
203            o.insert("timeout".into(), vi(m.timeout.into_inner()));
204            if let Some(rid) = &m.request_id {
205                o.insert("request_id".into(), vi(rid.into_inner()));
206            }
207            o
208        }
209        ControlMessage::RequestOk(m) => {
210            let mut o = Map::new();
211            o.insert("parameters".into(), params_to_json(&m.parameters));
212            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
213            o
214        }
215        ControlMessage::RequestError(m) => {
216            let mut o = Map::new();
217            o.insert("error_code".into(), vi(m.error_code.into_inner()));
218            o.insert("retry_interval".into(), vi(m.retry_interval.into_inner()));
219            o.insert(
220                "reason_phrase".into(),
221                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
222            );
223            if let Some(r) = &m.redirect {
224                let mut r_obj = Map::new();
225                r_obj.insert(
226                    "connect_uri".into(),
227                    Value::Text(String::from_utf8_lossy(&r.connect_uri).into_owned()),
228                );
229                r_obj.insert("track_namespace".into(), ns_to_json(&r.track_namespace));
230                r_obj.insert(
231                    "track_name".into(),
232                    Value::Text(String::from_utf8_lossy(&r.track_name).into_owned()),
233                );
234                o.insert("redirect".into(), Value::Map(r_obj));
235            }
236            o
237        }
238        ControlMessage::Subscribe(m) => {
239            let mut o = Map::new();
240            o.insert("request_id".into(), vi(m.request_id.into_inner()));
241            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
242            o.insert(
243                "track_name".into(),
244                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
245            );
246            o.insert("parameters".into(), params_to_json(&m.parameters));
247            o
248        }
249        ControlMessage::SubscribeOk(m) => {
250            let mut o = Map::new();
251            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
252            o.insert("parameters".into(), params_to_json(&m.parameters));
253            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
254            o
255        }
256        ControlMessage::RequestUpdate(m) => {
257            let mut o = Map::new();
258            o.insert("request_id".into(), vi(m.request_id.into_inner()));
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("track_namespace".into(), ns_to_json(&m.track_namespace));
266            o.insert(
267                "track_name".into(),
268                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
269            );
270            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
271            o.insert("parameters".into(), params_to_json(&m.parameters));
272            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
273            o
274        }
275        ControlMessage::PublishDone(m) => {
276            let mut o = Map::new();
277            o.insert("status_code".into(), vi(m.status_code.into_inner()));
278            o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
279            o.insert(
280                "reason_phrase".into(),
281                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
282            );
283            o
284        }
285        ControlMessage::PublishNamespace(m) => {
286            let mut o = Map::new();
287            o.insert("request_id".into(), vi(m.request_id.into_inner()));
288            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
289            o.insert("parameters".into(), params_to_json(&m.parameters));
290            o
291        }
292        ControlMessage::Namespace(m) => {
293            let mut o = Map::new();
294            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
295            o
296        }
297        ControlMessage::NamespaceDone(m) => {
298            let mut o = Map::new();
299            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
300            o
301        }
302        ControlMessage::SubscribeNamespace(m) => {
303            let mut o = Map::new();
304            o.insert("request_id".into(), vi(m.request_id.into_inner()));
305            o.insert("namespace_prefix".into(), ns_to_json(&m.namespace_prefix));
306            o.insert("parameters".into(), params_to_json(&m.parameters));
307            o
308        }
309        ControlMessage::SubscribeTracks(m) => {
310            let mut o = Map::new();
311            o.insert("request_id".into(), vi(m.request_id.into_inner()));
312            o.insert("namespace_prefix".into(), ns_to_json(&m.namespace_prefix));
313            o.insert("parameters".into(), params_to_json(&m.parameters));
314            o
315        }
316        ControlMessage::TrackStatus(m) => {
317            let mut o = Map::new();
318            o.insert("request_id".into(), vi(m.request_id.into_inner()));
319            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
320            o.insert(
321                "track_name".into(),
322                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
323            );
324            o.insert("parameters".into(), params_to_json(&m.parameters));
325            o
326        }
327        ControlMessage::Fetch(m) => {
328            let mut o = Map::new();
329            o.insert("request_id".into(), vi(m.request_id.into_inner()));
330            o.insert("fetch_type".into(), vi(m.fetch_type as u64));
331            match &m.fetch_payload {
332                crate::draft18::message::FetchPayload::Standalone {
333                    track_namespace,
334                    track_name,
335                    start_group,
336                    start_object,
337                    end_group,
338                    end_object,
339                } => {
340                    o.insert("track_namespace".into(), ns_to_json(track_namespace));
341                    o.insert(
342                        "track_name".into(),
343                        Value::Text(String::from_utf8_lossy(track_name).into_owned()),
344                    );
345                    o.insert("start_group".into(), vi(start_group.into_inner()));
346                    o.insert("start_object".into(), vi(start_object.into_inner()));
347                    o.insert("end_group".into(), vi(end_group.into_inner()));
348                    o.insert("end_object".into(), vi(end_object.into_inner()));
349                }
350                crate::draft18::message::FetchPayload::Joining {
351                    joining_request_id,
352                    joining_start,
353                } => {
354                    o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
355                    o.insert("joining_start".into(), vi(joining_start.into_inner()));
356                }
357            }
358            o.insert("parameters".into(), params_to_json(&m.parameters));
359            o
360        }
361        ControlMessage::FetchOk(m) => {
362            let mut o = Map::new();
363            o.insert("end_of_track".into(), vi(m.end_of_track as u64));
364            o.insert("end_group".into(), vi(m.end_group.into_inner()));
365            o.insert("end_object".into(), vi(m.end_object.into_inner()));
366            o.insert("parameters".into(), params_to_json(&m.parameters));
367            o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
368            o
369        }
370        ControlMessage::PublishBlocked(m) => {
371            let mut o = Map::new();
372            o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
373            o.insert(
374                "track_name".into(),
375                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
376            );
377            o
378        }
379    };
380    obj
381}