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    crate::fields::kvp_entries(params, |key, value| {
87        let Some(name) = name_fn(key) else {
88            return (None, None);
89        };
90        let rendered = match value {
91            KvpValue::Varint(v) => Value::Uint(v.into_inner()),
92            KvpValue::Bytes(b) if name == "authorization_token" => auth_token_to_json_d14(b),
93            KvpValue::Bytes(b) => Value::Text(String::from_utf8_lossy(b).into_owned()),
94        };
95        (Some(name), Some(rendered))
96    })
97}
98
99fn kvp_to_json_d14(params: &[KeyValuePair]) -> Value {
100    kvp_to_json(params, d14_msg_param_name)
101}
102
103fn kvp_to_json_d14_setup(params: &[KeyValuePair]) -> Value {
104    kvp_to_json(params, d14_setup_param_name)
105}
106
107/// This draft's field names for a decoded control message.
108///
109/// Keys are the names this draft gives its fields, in the order it defines
110/// them. An optional field the message did not carry is absent rather than
111/// zero.
112pub fn message_fields(msg: &ControlMessage) -> Map {
113    let obj = match msg {
114        ControlMessage::ClientSetup(m) => {
115            let mut o = Map::new();
116            o.insert(
117                "supported_versions".into(),
118                Value::Array(m.supported_versions.iter().map(|v| vi(v.into_inner())).collect()),
119            );
120            o.insert("parameters".into(), kvp_to_json_d14_setup(&m.parameters));
121            o
122        }
123        ControlMessage::ServerSetup(m) => {
124            let mut o = Map::new();
125            o.insert("selected_version".into(), vi(m.selected_version.into_inner()));
126            o.insert("parameters".into(), kvp_to_json_d14_setup(&m.parameters));
127            o
128        }
129        ControlMessage::GoAway(m) => {
130            let mut o = Map::new();
131            o.insert(
132                "new_session_uri".into(),
133                Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
134            );
135            o
136        }
137        ControlMessage::MaxRequestId(m) => {
138            let mut o = Map::new();
139            o.insert("request_id".into(), vi(m.request_id.into_inner()));
140            o
141        }
142        ControlMessage::RequestsBlocked(m) => {
143            let mut o = Map::new();
144            o.insert("request_id".into(), vi(m.maximum_request_id.into_inner()));
145            o
146        }
147        ControlMessage::Subscribe(m) => {
148            let mut o = Map::new();
149            o.insert("request_id".into(), vi(m.request_id.into_inner()));
150            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
151            o.insert(
152                "track_name".into(),
153                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
154            );
155            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
156            o.insert("group_order".into(), vi(m.group_order as u64));
157            o.insert("forward".into(), vi(m.forward as u64));
158            o.insert("filter_type".into(), vi(m.filter_type as u64));
159            if let Some(loc) = &m.start_location {
160                o.insert("start_group".into(), vi(loc.group.into_inner()));
161                o.insert("start_object".into(), vi(loc.object.into_inner()));
162            }
163            if let Some(eg) = &m.end_group {
164                o.insert("end_group".into(), vi(eg.into_inner()));
165            }
166            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
167            o
168        }
169        ControlMessage::SubscribeOk(m) => {
170            let mut o = Map::new();
171            o.insert("request_id".into(), vi(m.request_id.into_inner()));
172            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
173            o.insert("expires".into(), vi(m.expires.into_inner()));
174            o.insert("group_order".into(), vi(m.group_order as u64));
175            o.insert("content_exists".into(), vi(m.content_exists as u64));
176            if let Some(loc) = &m.largest_location {
177                o.insert("largest_location".into(), loc_to_json(loc));
178            }
179            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
180            o
181        }
182        ControlMessage::SubscribeError(m) => {
183            let mut o = Map::new();
184            o.insert("request_id".into(), vi(m.request_id.into_inner()));
185            o.insert("error_code".into(), vi(m.error_code.into_inner()));
186            o.insert(
187                "reason_phrase".into(),
188                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
189            );
190            o
191        }
192        ControlMessage::SubscribeUpdate(m) => {
193            let mut o = Map::new();
194            o.insert("request_id".into(), vi(m.request_id.into_inner()));
195            o.insert("subscription_request_id".into(), vi(m.subscription_request_id.into_inner()));
196            o.insert("start_group".into(), vi(m.start_location.group.into_inner()));
197            o.insert("start_object".into(), vi(m.start_location.object.into_inner()));
198            o.insert("end_group".into(), vi(m.end_group.into_inner()));
199            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
200            o.insert("forward".into(), vi(m.forward as u64));
201            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
202            o
203        }
204        ControlMessage::Unsubscribe(m) => {
205            let mut o = Map::new();
206            o.insert("request_id".into(), vi(m.request_id.into_inner()));
207            o
208        }
209        ControlMessage::Publish(m) => {
210            let mut o = Map::new();
211            o.insert("request_id".into(), vi(m.request_id.into_inner()));
212            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
213            o.insert(
214                "track_name".into(),
215                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
216            );
217            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
218            o.insert("group_order".into(), vi(m.group_order as u64));
219            o.insert("content_exists".into(), vi(m.content_exists as u64));
220            if let Some(loc) = &m.largest_location {
221                o.insert("largest_location".into(), loc_to_json(loc));
222            }
223            o.insert("forward".into(), vi(m.forward as u64));
224            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
225            o
226        }
227        ControlMessage::PublishOk(m) => {
228            let mut o = Map::new();
229            o.insert("request_id".into(), vi(m.request_id.into_inner()));
230            o.insert("forward".into(), vi(m.forward as u64));
231            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
232            o.insert("group_order".into(), vi(m.group_order as u64));
233            o.insert("filter_type".into(), vi(m.filter_type as u64));
234            if let Some(loc) = &m.start_location {
235                o.insert("start_group".into(), vi(loc.group.into_inner()));
236                o.insert("start_object".into(), vi(loc.object.into_inner()));
237            }
238            if let Some(eg) = &m.end_group {
239                o.insert("end_group".into(), vi(eg.into_inner()));
240            }
241            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
242            o
243        }
244        ControlMessage::PublishError(m) => {
245            let mut o = Map::new();
246            o.insert("request_id".into(), vi(m.request_id.into_inner()));
247            o.insert("error_code".into(), vi(m.error_code.into_inner()));
248            o.insert(
249                "reason_phrase".into(),
250                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
251            );
252            o
253        }
254        ControlMessage::PublishDone(m) => {
255            let mut o = Map::new();
256            o.insert("request_id".into(), vi(m.request_id.into_inner()));
257            o.insert("status_code".into(), vi(m.status_code.into_inner()));
258            o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
259            o.insert(
260                "reason_phrase".into(),
261                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
262            );
263            o
264        }
265        ControlMessage::PublishNamespace(m) => {
266            let mut o = Map::new();
267            o.insert("request_id".into(), vi(m.request_id.into_inner()));
268            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
269            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
270            o
271        }
272        ControlMessage::PublishNamespaceOk(m) => {
273            let mut o = Map::new();
274            o.insert("request_id".into(), vi(m.request_id.into_inner()));
275            o
276        }
277        ControlMessage::PublishNamespaceError(m) => {
278            let mut o = Map::new();
279            o.insert("request_id".into(), vi(m.request_id.into_inner()));
280            o.insert("error_code".into(), vi(m.error_code.into_inner()));
281            o.insert(
282                "reason_phrase".into(),
283                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
284            );
285            o
286        }
287        ControlMessage::PublishNamespaceDone(m) => {
288            let mut o = Map::new();
289            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
290            o
291        }
292        ControlMessage::PublishNamespaceCancel(m) => {
293            let mut o = Map::new();
294            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
295            o.insert("error_code".into(), vi(m.error_code.into_inner()));
296            o.insert(
297                "reason_phrase".into(),
298                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
299            );
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.track_namespace));
306            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
307            o
308        }
309        ControlMessage::SubscribeNamespaceOk(m) => {
310            let mut o = Map::new();
311            o.insert("request_id".into(), vi(m.request_id.into_inner()));
312            o
313        }
314        ControlMessage::SubscribeNamespaceError(m) => {
315            let mut o = Map::new();
316            o.insert("request_id".into(), vi(m.request_id.into_inner()));
317            o.insert("error_code".into(), vi(m.error_code.into_inner()));
318            o.insert(
319                "reason_phrase".into(),
320                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
321            );
322            o
323        }
324        ControlMessage::UnsubscribeNamespace(m) => {
325            let mut o = Map::new();
326            o.insert("track_namespace_prefix".into(), ns_to_json(&m.track_namespace_prefix));
327            o
328        }
329        ControlMessage::Fetch(m) => {
330            let mut o = Map::new();
331            o.insert("request_id".into(), vi(m.request_id.into_inner()));
332            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
333            o.insert("group_order".into(), vi(m.group_order as u64));
334            o.insert("fetch_type".into(), vi(m.fetch_type as u64));
335            match &m.fetch_payload {
336                FetchPayload::Standalone {
337                    track_namespace,
338                    track_name,
339                    start_group,
340                    start_object,
341                    end_group,
342                    end_object,
343                } => {
344                    o.insert("track_namespace".into(), ns_to_json(track_namespace));
345                    o.insert(
346                        "track_name".into(),
347                        Value::Text(String::from_utf8_lossy(track_name).into_owned()),
348                    );
349                    o.insert("start_group".into(), vi(start_group.into_inner()));
350                    o.insert("start_object".into(), vi(start_object.into_inner()));
351                    o.insert("end_group".into(), vi(end_group.into_inner()));
352                    o.insert("end_object".into(), vi(end_object.into_inner()));
353                }
354                FetchPayload::Joining { joining_request_id, joining_start } => {
355                    o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
356                    o.insert("joining_start".into(), vi(joining_start.into_inner()));
357                }
358            }
359            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
360            o
361        }
362        ControlMessage::FetchOk(m) => {
363            let mut o = Map::new();
364            o.insert("request_id".into(), vi(m.request_id.into_inner()));
365            o.insert("group_order".into(), vi(m.group_order as u64));
366            o.insert("end_of_track".into(), vi(m.end_of_track as u64));
367            o.insert("end_location".into(), loc_to_json(&m.end_location));
368            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
369            o
370        }
371        ControlMessage::FetchError(m) => {
372            let mut o = Map::new();
373            o.insert("request_id".into(), vi(m.request_id.into_inner()));
374            o.insert("error_code".into(), vi(m.error_code.into_inner()));
375            o.insert(
376                "reason_phrase".into(),
377                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
378            );
379            o
380        }
381        ControlMessage::FetchCancel(m) => {
382            let mut o = Map::new();
383            o.insert("request_id".into(), vi(m.request_id.into_inner()));
384            o
385        }
386        ControlMessage::TrackStatus(m) => {
387            let mut o = Map::new();
388            o.insert("request_id".into(), vi(m.request_id.into_inner()));
389            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
390            o.insert(
391                "track_name".into(),
392                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
393            );
394            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
395            o.insert("group_order".into(), vi(m.group_order as u64));
396            o.insert("forward".into(), vi(m.forward as u64));
397            o.insert("filter_type".into(), vi(m.filter_type as u64));
398            if let Some(loc) = &m.start_location {
399                o.insert("start_group".into(), vi(loc.group.into_inner()));
400                o.insert("start_object".into(), vi(loc.object.into_inner()));
401            }
402            if let Some(eg) = &m.end_group {
403                o.insert("end_group".into(), vi(eg.into_inner()));
404            }
405            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
406            o
407        }
408        ControlMessage::TrackStatusOk(m) => {
409            let mut o = Map::new();
410            o.insert("request_id".into(), vi(m.request_id.into_inner()));
411            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
412            o.insert("expires".into(), vi(m.expires.into_inner()));
413            o.insert("group_order".into(), vi(m.group_order as u64));
414            o.insert("content_exists".into(), vi(m.content_exists as u64));
415            if let Some(loc) = &m.largest_location {
416                o.insert("largest_location".into(), loc_to_json(loc));
417            }
418            o.insert("parameters".into(), kvp_to_json_d14(&m.parameters));
419            o
420        }
421        ControlMessage::TrackStatusError(m) => {
422            let mut o = Map::new();
423            o.insert("request_id".into(), vi(m.request_id.into_inner()));
424            o.insert("error_code".into(), vi(m.error_code.into_inner()));
425            o.insert(
426                "reason_phrase".into(),
427                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
428            );
429            o
430        }
431    };
432    obj
433}