Skip to main content

moqtap_codec/draft14/
fields.rs

1use crate::draft14::message::{ControlMessage, FetchPayload};
2use crate::fields::{FieldMap as Map, FieldValue as Value};
3use crate::types::*;
4
5use crate::kvp::{KeyValuePair, KvpValue};
6use crate::varint::VarInt;
7
8fn vi(v: u64) -> Value {
9    Value::Uint(v)
10}
11
12fn ns_to_json(ns: &TrackNamespace) -> Value {
13    Value::Array(
14        ns.0.iter().map(|e| Value::Text(String::from_utf8_lossy(e).into_owned())).collect(),
15    )
16}
17
18fn loc_to_json(loc: &Location) -> Value {
19    let mut o = Map::new();
20    o.insert("group".into(), vi(loc.group.into_inner()));
21    o.insert("object".into(), vi(loc.object.into_inner()));
22    Value::Map(o)
23}
24
25/// Parse draft-14+ authorization_token bytes into structured JSON.
26/// Structure: alias_type (varint), [token_alias (varint)?], [token_type (varint), token_value (bytes)?]
27/// depending on alias_type (0=DELETE, 1=REGISTER, 2=USE_ALIAS, 3=USE_VALUE).
28fn auth_token_to_json_d14(bytes: &[u8]) -> Value {
29    let mut buf = bytes;
30    let alias_type = match VarInt::decode(&mut buf) {
31        Ok(v) => v,
32        Err(_) => return Value::Bytes(bytes.to_vec()),
33    };
34    let at = alias_type.into_inner();
35    let mut o = Map::new();
36    o.insert("alias_type".to_string(), Value::Uint(at));
37    match at {
38        0 | 2 => {
39            if let Ok(ta) = VarInt::decode(&mut buf) {
40                o.insert("token_alias".to_string(), Value::Uint(ta.into_inner()));
41            }
42        }
43        1 => {
44            if let Ok(ta) = VarInt::decode(&mut buf) {
45                o.insert("token_alias".to_string(), Value::Uint(ta.into_inner()));
46            }
47            if let Ok(tt) = VarInt::decode(&mut buf) {
48                o.insert("token_type".to_string(), Value::Uint(tt.into_inner()));
49            }
50            o.insert("token_value".to_string(), Value::Bytes(buf.to_vec()));
51        }
52        _ => {
53            if let Ok(tt) = VarInt::decode(&mut buf) {
54                o.insert("token_type".to_string(), Value::Uint(tt.into_inner()));
55            }
56            o.insert("token_value".to_string(), Value::Bytes(buf.to_vec()));
57        }
58    }
59    Value::Map(o)
60}
61
62/// Known parameter names for draft-14+ SETUP messages.
63fn d14_setup_param_name(key: u64) -> Option<&'static str> {
64    match key {
65        0x01 => Some("path"),
66        0x02 => Some("max_request_id"),
67        0x03 => Some("authorization_token"),
68        0x04 => Some("max_auth_token_cache_size"),
69        0x05 => Some("authority"),
70        _ => None,
71    }
72}
73
74/// Known parameter names for draft-14+ non-SETUP messages.
75fn d14_msg_param_name(key: u64) -> Option<&'static str> {
76    match key {
77        0x02 => Some("delivery_timeout"),
78        0x03 => Some("authorization_token"),
79        0x04 => Some("max_cache_duration"),
80        _ => None,
81    }
82}
83
84/// Convert KVP list to JSON Value matching test vector format.
85fn kvp_to_json(params: &[KeyValuePair], name_fn: fn(u64) -> Option<&'static str>) -> Value {
86    let mut obj = Map::new();
87    let mut unknown = Vec::new();
88
89    for p in params {
90        let key = p.key.into_inner();
91        if let Some(name) = name_fn(key) {
92            match &p.value {
93                KvpValue::Varint(v) => {
94                    obj.insert(name.to_string(), Value::Uint(v.into_inner()));
95                }
96                KvpValue::Bytes(b) => {
97                    if name == "authorization_token" {
98                        obj.insert(name.to_string(), auth_token_to_json_d14(b));
99                    } else {
100                        obj.insert(
101                            name.to_string(),
102                            Value::Text(String::from_utf8_lossy(b).into_owned()),
103                        );
104                    }
105                }
106            }
107        } else {
108            let mut entry = Map::new();
109            entry.insert("id".to_string(), Value::Text(format!("0x{:x}", key)));
110            match &p.value {
111                KvpValue::Varint(v) => {
112                    entry.insert("length".to_string(), Value::Uint(v.into_inner()));
113                }
114                KvpValue::Bytes(b) => {
115                    entry.insert("length".to_string(), Value::Uint(b.len() as u64));
116                    entry.insert("raw_hex".to_string(), Value::Bytes(b.to_vec()));
117                }
118            }
119            unknown.push(Value::Map(entry));
120        }
121    }
122
123    if !unknown.is_empty() {
124        obj.insert("unknown".to_string(), Value::Array(unknown));
125    }
126
127    Value::Map(obj)
128}
129
130fn kvp_to_json_d14(params: &[KeyValuePair]) -> Value {
131    kvp_to_json(params, d14_msg_param_name)
132}
133
134fn kvp_to_json_d14_setup(params: &[KeyValuePair]) -> Value {
135    kvp_to_json(params, d14_setup_param_name)
136}
137
138/// This draft's field names for a decoded control message.
139///
140/// Keys are the names this draft gives its fields, in the order it defines
141/// them. An optional field the message did not carry is absent rather than
142/// zero.
143pub fn message_fields(msg: &ControlMessage) -> Map {
144    let obj = match msg {
145        ControlMessage::ClientSetup(m) => {
146            let mut o = Map::new();
147            o.insert(
148                "supported_versions".into(),
149                Value::Array(m.supported_versions.iter().map(|v| vi(v.into_inner())).collect()),
150            );
151            o.insert("parameters".into(), kvp_to_json_d14_setup(&m.parameters));
152            o
153        }
154        ControlMessage::ServerSetup(m) => {
155            let mut o = Map::new();
156            o.insert("selected_version".into(), vi(m.selected_version.into_inner()));
157            o.insert("parameters".into(), kvp_to_json_d14_setup(&m.parameters));
158            o
159        }
160        ControlMessage::GoAway(m) => {
161            let mut o = Map::new();
162            o.insert(
163                "new_session_uri".into(),
164                Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
165            );
166            o
167        }
168        ControlMessage::MaxRequestId(m) => {
169            let mut o = Map::new();
170            o.insert("request_id".into(), vi(m.request_id.into_inner()));
171            o
172        }
173        ControlMessage::RequestsBlocked(m) => {
174            let mut o = Map::new();
175            o.insert("request_id".into(), vi(m.maximum_request_id.into_inner()));
176            o
177        }
178        ControlMessage::Subscribe(m) => {
179            let mut o = Map::new();
180            o.insert("request_id".into(), vi(m.request_id.into_inner()));
181            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
182            o.insert(
183                "track_name".into(),
184                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
185            );
186            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
187            o.insert("group_order".into(), vi(m.group_order as u64));
188            o.insert("forward".into(), vi(m.forward as u64));
189            o.insert("filter_type".into(), vi(m.filter_type as u64));
190            if let Some(loc) = &m.start_location {
191                o.insert("start_group".into(), vi(loc.group.into_inner()));
192                o.insert("start_object".into(), vi(loc.object.into_inner()));
193            }
194            if let Some(eg) = &m.end_group {
195                o.insert("end_group".into(), vi(eg.into_inner()));
196            }
197            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
198            o
199        }
200        ControlMessage::SubscribeOk(m) => {
201            let mut o = Map::new();
202            o.insert("request_id".into(), vi(m.request_id.into_inner()));
203            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
204            o.insert("expires".into(), vi(m.expires.into_inner()));
205            o.insert("group_order".into(), vi(m.group_order as u64));
206            o.insert("content_exists".into(), vi(m.content_exists as u64));
207            if let Some(loc) = &m.largest_location {
208                o.insert("largest_location".into(), loc_to_json(loc));
209            }
210            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
211            o
212        }
213        ControlMessage::SubscribeError(m) => {
214            let mut o = Map::new();
215            o.insert("request_id".into(), vi(m.request_id.into_inner()));
216            o.insert("error_code".into(), vi(m.error_code.into_inner()));
217            o.insert(
218                "reason_phrase".into(),
219                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
220            );
221            o
222        }
223        ControlMessage::SubscribeUpdate(m) => {
224            let mut o = Map::new();
225            o.insert("request_id".into(), vi(m.request_id.into_inner()));
226            o.insert("subscription_request_id".into(), vi(m.subscription_request_id.into_inner()));
227            o.insert("start_group".into(), vi(m.start_location.group.into_inner()));
228            o.insert("start_object".into(), vi(m.start_location.object.into_inner()));
229            o.insert("end_group".into(), vi(m.end_group.into_inner()));
230            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
231            o.insert("forward".into(), vi(m.forward as u64));
232            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
233            o
234        }
235        ControlMessage::Unsubscribe(m) => {
236            let mut o = Map::new();
237            o.insert("request_id".into(), vi(m.request_id.into_inner()));
238            o
239        }
240        ControlMessage::Publish(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("track_alias".into(), vi(m.track_alias.into_inner()));
249            o.insert("group_order".into(), vi(m.group_order as u64));
250            o.insert("content_exists".into(), vi(m.content_exists as u64));
251            if let Some(loc) = &m.largest_location {
252                o.insert("largest_location".into(), loc_to_json(loc));
253            }
254            o.insert("forward".into(), vi(m.forward as u64));
255            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
256            o
257        }
258        ControlMessage::PublishOk(m) => {
259            let mut o = Map::new();
260            o.insert("request_id".into(), vi(m.request_id.into_inner()));
261            o.insert("forward".into(), vi(m.forward as u64));
262            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
263            o.insert("group_order".into(), vi(m.group_order as u64));
264            o.insert("filter_type".into(), vi(m.filter_type as u64));
265            if let Some(loc) = &m.start_location {
266                o.insert("start_group".into(), vi(loc.group.into_inner()));
267                o.insert("start_object".into(), vi(loc.object.into_inner()));
268            }
269            if let Some(eg) = &m.end_group {
270                o.insert("end_group".into(), vi(eg.into_inner()));
271            }
272            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
273            o
274        }
275        ControlMessage::PublishError(m) => {
276            let mut o = Map::new();
277            o.insert("request_id".into(), vi(m.request_id.into_inner()));
278            o.insert("error_code".into(), vi(m.error_code.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::PublishDone(m) => {
286            let mut o = Map::new();
287            o.insert("request_id".into(), vi(m.request_id.into_inner()));
288            o.insert("status_code".into(), vi(m.status_code.into_inner()));
289            o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
290            o.insert(
291                "reason_phrase".into(),
292                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
293            );
294            o
295        }
296        ControlMessage::PublishNamespace(m) => {
297            let mut o = Map::new();
298            o.insert("request_id".into(), vi(m.request_id.into_inner()));
299            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
300            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
301            o
302        }
303        ControlMessage::PublishNamespaceOk(m) => {
304            let mut o = Map::new();
305            o.insert("request_id".into(), vi(m.request_id.into_inner()));
306            o
307        }
308        ControlMessage::PublishNamespaceError(m) => {
309            let mut o = Map::new();
310            o.insert("request_id".into(), vi(m.request_id.into_inner()));
311            o.insert("error_code".into(), vi(m.error_code.into_inner()));
312            o.insert(
313                "reason_phrase".into(),
314                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
315            );
316            o
317        }
318        ControlMessage::PublishNamespaceDone(m) => {
319            let mut o = Map::new();
320            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
321            o
322        }
323        ControlMessage::PublishNamespaceCancel(m) => {
324            let mut o = Map::new();
325            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
326            o.insert("error_code".into(), vi(m.error_code.into_inner()));
327            o.insert(
328                "reason_phrase".into(),
329                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
330            );
331            o
332        }
333        ControlMessage::SubscribeNamespace(m) => {
334            let mut o = Map::new();
335            o.insert("request_id".into(), vi(m.request_id.into_inner()));
336            o.insert("namespace_prefix".into(), ns_to_json(&m.track_namespace));
337            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
338            o
339        }
340        ControlMessage::SubscribeNamespaceOk(m) => {
341            let mut o = Map::new();
342            o.insert("request_id".into(), vi(m.request_id.into_inner()));
343            o
344        }
345        ControlMessage::SubscribeNamespaceError(m) => {
346            let mut o = Map::new();
347            o.insert("request_id".into(), vi(m.request_id.into_inner()));
348            o.insert("error_code".into(), vi(m.error_code.into_inner()));
349            o.insert(
350                "reason_phrase".into(),
351                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
352            );
353            o
354        }
355        ControlMessage::UnsubscribeNamespace(m) => {
356            let mut o = Map::new();
357            o.insert("track_namespace_prefix".into(), ns_to_json(&m.track_namespace_prefix));
358            o
359        }
360        ControlMessage::Fetch(m) => {
361            let mut o = Map::new();
362            o.insert("request_id".into(), vi(m.request_id.into_inner()));
363            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
364            o.insert("group_order".into(), vi(m.group_order as u64));
365            o.insert("fetch_type".into(), vi(m.fetch_type as u64));
366            match &m.fetch_payload {
367                FetchPayload::Standalone {
368                    track_namespace,
369                    track_name,
370                    start_group,
371                    start_object,
372                    end_group,
373                    end_object,
374                } => {
375                    o.insert("track_namespace".into(), ns_to_json(track_namespace));
376                    o.insert(
377                        "track_name".into(),
378                        Value::Text(String::from_utf8_lossy(track_name).into_owned()),
379                    );
380                    o.insert("start_group".into(), vi(start_group.into_inner()));
381                    o.insert("start_object".into(), vi(start_object.into_inner()));
382                    o.insert("end_group".into(), vi(end_group.into_inner()));
383                    o.insert("end_object".into(), vi(end_object.into_inner()));
384                }
385                FetchPayload::Joining { joining_request_id, joining_start } => {
386                    o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
387                    o.insert("joining_start".into(), vi(joining_start.into_inner()));
388                }
389            }
390            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
391            o
392        }
393        ControlMessage::FetchOk(m) => {
394            let mut o = Map::new();
395            o.insert("request_id".into(), vi(m.request_id.into_inner()));
396            o.insert("group_order".into(), vi(m.group_order as u64));
397            o.insert("end_of_track".into(), vi(m.end_of_track as u64));
398            o.insert("end_location".into(), loc_to_json(&m.end_location));
399            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
400            o
401        }
402        ControlMessage::FetchError(m) => {
403            let mut o = Map::new();
404            o.insert("request_id".into(), vi(m.request_id.into_inner()));
405            o.insert("error_code".into(), vi(m.error_code.into_inner()));
406            o.insert(
407                "reason_phrase".into(),
408                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
409            );
410            o
411        }
412        ControlMessage::FetchCancel(m) => {
413            let mut o = Map::new();
414            o.insert("request_id".into(), vi(m.request_id.into_inner()));
415            o
416        }
417        ControlMessage::TrackStatus(m) => {
418            let mut o = Map::new();
419            o.insert("request_id".into(), vi(m.request_id.into_inner()));
420            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
421            o.insert(
422                "track_name".into(),
423                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
424            );
425            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
426            o.insert("group_order".into(), vi(m.group_order as u64));
427            o.insert("forward".into(), vi(m.forward as u64));
428            o.insert("filter_type".into(), vi(m.filter_type as u64));
429            if let Some(loc) = &m.start_location {
430                o.insert("start_group".into(), vi(loc.group.into_inner()));
431                o.insert("start_object".into(), vi(loc.object.into_inner()));
432            }
433            if let Some(eg) = &m.end_group {
434                o.insert("end_group".into(), vi(eg.into_inner()));
435            }
436            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
437            o
438        }
439        ControlMessage::TrackStatusOk(m) => {
440            let mut o = Map::new();
441            o.insert("request_id".into(), vi(m.request_id.into_inner()));
442            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
443            o.insert("expires".into(), vi(m.expires.into_inner()));
444            o.insert("group_order".into(), vi(m.group_order as u64));
445            o.insert("content_exists".into(), vi(m.content_exists as u64));
446            if let Some(loc) = &m.largest_location {
447                o.insert("largest_location".into(), loc_to_json(loc));
448            }
449            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
450            o
451        }
452        ControlMessage::TrackStatusError(m) => {
453            let mut o = Map::new();
454            o.insert("request_id".into(), vi(m.request_id.into_inner()));
455            o.insert("error_code".into(), vi(m.error_code.into_inner()));
456            o.insert(
457                "reason_phrase".into(),
458                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
459            );
460            o
461        }
462    };
463    obj
464}