Skip to main content

moqtap_codec/draft13/
fields.rs

1use crate::draft13::message::*;
2use crate::fields::{FieldMap as Map, FieldValue as Value};
3use crate::kvp::{KeyValuePair, KvpValue};
4use crate::types::*;
5use crate::varint::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
17fn loc_to_json(loc: &Location) -> Value {
18    let mut o = Map::new();
19    o.insert("group".into(), vi(loc.group.into_inner()));
20    o.insert("object".into(), vi(loc.object.into_inner()));
21    Value::Map(o)
22}
23
24fn auth_token_to_json(bytes: &[u8]) -> Value {
25    let mut buf = bytes;
26    let alias_type = VarInt::decode(&mut buf).unwrap();
27    let at = alias_type.into_inner();
28    let mut o = Map::new();
29    o.insert("alias_type".into(), vi(at));
30    match at {
31        // DELETE, USE_ALIAS: only token_alias
32        0 | 2 => {
33            let token_alias = VarInt::decode(&mut buf).unwrap();
34            o.insert("token_alias".into(), vi(token_alias.into_inner()));
35        }
36        // REGISTER: token_alias + token_type + token_value
37        1 => {
38            let token_alias = VarInt::decode(&mut buf).unwrap();
39            let token_type = VarInt::decode(&mut buf).unwrap();
40            o.insert("token_alias".into(), vi(token_alias.into_inner()));
41            o.insert("token_type".into(), vi(token_type.into_inner()));
42            o.insert("token_value".into(), Value::Bytes(buf.to_vec()));
43        }
44        // USE_VALUE (3) or other: token_type + token_value
45        _ => {
46            let token_type = VarInt::decode(&mut buf).unwrap();
47            o.insert("token_type".into(), vi(token_type.into_inner()));
48            o.insert("token_value".into(), Value::Bytes(buf.to_vec()));
49        }
50    }
51    Value::Map(o)
52}
53
54fn kvp_to_json_setup(params: &[KeyValuePair]) -> Value {
55    crate::fields::kvp_entries(params, |key, value| match (key, value) {
56        (0x01, KvpValue::Bytes(b)) => {
57            (Some("path"), Some(Value::Text(String::from_utf8_lossy(b).into_owned())))
58        }
59        (0x02, KvpValue::Varint(v)) => (Some("max_request_id"), Some(vi(v.into_inner()))),
60        (0x03, KvpValue::Bytes(b)) => (Some("authorization_token"), Some(auth_token_to_json(b))),
61        (0x04, KvpValue::Varint(v)) => {
62            (Some("max_auth_token_cache_size"), Some(vi(v.into_inner())))
63        }
64        _ => (None, None),
65    })
66}
67
68fn kvp_to_json_msg(params: &[KeyValuePair]) -> Value {
69    crate::fields::kvp_entries(params, |key, value| match (key, value) {
70        (0x03, KvpValue::Bytes(b)) => (Some("authorization_token"), Some(auth_token_to_json(b))),
71        (0x02, KvpValue::Varint(v)) => (Some("delivery_timeout"), Some(vi(v.into_inner()))),
72        (0x04, KvpValue::Varint(v)) => (Some("max_cache_duration"), Some(vi(v.into_inner()))),
73        _ => (None, None),
74    })
75}
76
77/// This draft's field names for a decoded control message.
78///
79/// Keys are the names this draft gives its fields, in the order it defines
80/// them. An optional field the message did not carry is absent rather than
81/// zero.
82pub fn message_fields(msg: &ControlMessage) -> Map {
83    let obj = match msg {
84        ControlMessage::ClientSetup(m) => {
85            let mut o = Map::new();
86            o.insert(
87                "supported_versions".into(),
88                Value::Array(m.supported_versions.iter().map(|v| vi(v.into_inner())).collect()),
89            );
90            o.insert("parameters".into(), kvp_to_json_setup(&m.parameters));
91            o
92        }
93        ControlMessage::ServerSetup(m) => {
94            let mut o = Map::new();
95            o.insert("selected_version".into(), vi(m.selected_version.into_inner()));
96            o.insert("parameters".into(), kvp_to_json_setup(&m.parameters));
97            o
98        }
99        ControlMessage::GoAway(m) => {
100            let mut o = Map::new();
101            o.insert(
102                "new_session_uri".into(),
103                Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
104            );
105            o
106        }
107        ControlMessage::MaxRequestId(m) => {
108            let mut o = Map::new();
109            o.insert("request_id".into(), vi(m.request_id.into_inner()));
110            o
111        }
112        ControlMessage::RequestsBlocked(m) => {
113            let mut o = Map::new();
114            o.insert("maximum_request_id".into(), vi(m.maximum_request_id.into_inner()));
115            o
116        }
117        ControlMessage::Subscribe(m) => {
118            let mut o = Map::new();
119            o.insert("request_id".into(), vi(m.request_id.into_inner()));
120            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
121            o.insert(
122                "track_name".into(),
123                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
124            );
125            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
126            o.insert("group_order".into(), vi(m.group_order as u64));
127            o.insert("forward".into(), vi(m.forward as u64));
128            o.insert("filter_type".into(), vi(m.filter_type as u64));
129            if let Some(sg) = &m.start_group {
130                o.insert("start_group".into(), vi(sg.into_inner()));
131            }
132            if let Some(so) = &m.start_object {
133                o.insert("start_object".into(), vi(so.into_inner()));
134            }
135            if let Some(eg) = &m.end_group {
136                o.insert("end_group".into(), vi(eg.into_inner()));
137            }
138            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
139            o
140        }
141        ControlMessage::SubscribeOk(m) => {
142            let mut o = Map::new();
143            o.insert("request_id".into(), vi(m.request_id.into_inner()));
144            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
145            o.insert("expires".into(), vi(m.expires.into_inner()));
146            o.insert("group_order".into(), vi(m.group_order as u64));
147            o.insert("content_exists".into(), vi(m.content_exists as u64));
148            if let Some(loc) = &m.largest_location {
149                o.insert("largest_location".into(), loc_to_json(loc));
150            }
151            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
152            o
153        }
154        ControlMessage::SubscribeError(m) => {
155            let mut o = Map::new();
156            o.insert("request_id".into(), vi(m.request_id.into_inner()));
157            o.insert("error_code".into(), vi(m.error_code.into_inner()));
158            o.insert(
159                "reason_phrase".into(),
160                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
161            );
162            o
163        }
164        ControlMessage::SubscribeUpdate(m) => {
165            let mut o = Map::new();
166            o.insert("request_id".into(), vi(m.request_id.into_inner()));
167            o.insert("start_group".into(), vi(m.start_group.into_inner()));
168            o.insert("start_object".into(), vi(m.start_object.into_inner()));
169            o.insert("end_group".into(), vi(m.end_group.into_inner()));
170            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
171            o.insert("forward".into(), vi(m.forward as u64));
172            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
173            o
174        }
175        ControlMessage::SubscribeDone(m) => {
176            let mut o = Map::new();
177            o.insert("request_id".into(), vi(m.request_id.into_inner()));
178            o.insert("status_code".into(), vi(m.status_code.into_inner()));
179            o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
180            o.insert(
181                "reason_phrase".into(),
182                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
183            );
184            o
185        }
186        ControlMessage::Unsubscribe(m) => {
187            let mut o = Map::new();
188            o.insert("request_id".into(), vi(m.request_id.into_inner()));
189            o
190        }
191        ControlMessage::Announce(m) => {
192            let mut o = Map::new();
193            o.insert("request_id".into(), vi(m.request_id.into_inner()));
194            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
195            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
196            o
197        }
198        ControlMessage::AnnounceOk(m) => {
199            let mut o = Map::new();
200            o.insert("request_id".into(), vi(m.request_id.into_inner()));
201            o
202        }
203        ControlMessage::AnnounceError(m) => {
204            let mut o = Map::new();
205            o.insert("request_id".into(), vi(m.request_id.into_inner()));
206            o.insert("error_code".into(), vi(m.error_code.into_inner()));
207            o.insert(
208                "reason_phrase".into(),
209                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
210            );
211            o
212        }
213        ControlMessage::AnnounceCancel(m) => {
214            let mut o = Map::new();
215            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
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::Unannounce(m) => {
224            let mut o = Map::new();
225            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
226            o
227        }
228        ControlMessage::SubscribeNamespace(m) => {
229            let mut o = Map::new();
230            o.insert("request_id".into(), vi(m.request_id.into_inner()));
231            o.insert("track_namespace_prefix".into(), ns_to_json(&m.track_namespace_prefix));
232            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
233            o
234        }
235        ControlMessage::SubscribeNamespaceOk(m) => {
236            let mut o = Map::new();
237            o.insert("request_id".into(), vi(m.request_id.into_inner()));
238            o
239        }
240        ControlMessage::SubscribeNamespaceError(m) => {
241            let mut o = Map::new();
242            o.insert("request_id".into(), vi(m.request_id.into_inner()));
243            o.insert("error_code".into(), vi(m.error_code.into_inner()));
244            o.insert(
245                "reason_phrase".into(),
246                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
247            );
248            o
249        }
250        ControlMessage::UnsubscribeNamespace(m) => {
251            let mut o = Map::new();
252            o.insert("track_namespace_prefix".into(), ns_to_json(&m.track_namespace_prefix));
253            o
254        }
255        ControlMessage::TrackStatus(m) => {
256            let mut o = Map::new();
257            o.insert("request_id".into(), vi(m.request_id.into_inner()));
258            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
259            o.insert(
260                "track_name".into(),
261                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
262            );
263            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
264            o.insert("group_order".into(), vi(m.group_order as u64));
265            o.insert("forward".into(), vi(m.forward as u64));
266            o.insert("filter_type".into(), vi(m.filter_type as u64));
267            if let Some(sg) = &m.start_group {
268                o.insert("start_group".into(), vi(sg.into_inner()));
269            }
270            if let Some(so) = &m.start_object {
271                o.insert("start_object".into(), vi(so.into_inner()));
272            }
273            if let Some(eg) = &m.end_group {
274                o.insert("end_group".into(), vi(eg.into_inner()));
275            }
276            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
277            o
278        }
279        ControlMessage::TrackStatusOk(m) => {
280            let mut o = Map::new();
281            o.insert("request_id".into(), vi(m.request_id.into_inner()));
282            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
283            o.insert("expires".into(), vi(m.expires.into_inner()));
284            o.insert("group_order".into(), vi(m.group_order as u64));
285            o.insert("content_exists".into(), vi(m.content_exists as u64));
286            if let Some(loc) = &m.largest_location {
287                o.insert("largest_location".into(), loc_to_json(loc));
288            }
289            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
290            o
291        }
292        ControlMessage::TrackStatusError(m) => {
293            let mut o = Map::new();
294            o.insert("request_id".into(), vi(m.request_id.into_inner()));
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::Fetch(m) => {
303            let mut o = Map::new();
304            o.insert("request_id".into(), vi(m.request_id.into_inner()));
305            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
306            o.insert("group_order".into(), vi(m.group_order as u64));
307            o.insert("fetch_type".into(), vi(m.fetch_type as u64));
308            match &m.fetch_payload {
309                FetchPayload::Standalone {
310                    track_namespace,
311                    track_name,
312                    start_group,
313                    start_object,
314                    end_group,
315                    end_object,
316                } => {
317                    o.insert("track_namespace".into(), ns_to_json(track_namespace));
318                    o.insert(
319                        "track_name".into(),
320                        Value::Text(String::from_utf8_lossy(track_name).into_owned()),
321                    );
322                    o.insert("start_group".into(), vi(start_group.into_inner()));
323                    o.insert("start_object".into(), vi(start_object.into_inner()));
324                    o.insert("end_group".into(), vi(end_group.into_inner()));
325                    o.insert("end_object".into(), vi(end_object.into_inner()));
326                }
327                FetchPayload::Joining { joining_request_id, joining_start } => {
328                    // The key is the corpus's spelling, not this draft's. The
329                    // draft renamed the field to Joining Request ID; the shared
330                    // vector files still say joining_subscribe_id, and they are
331                    // maintained elsewhere, so the name is translated here
332                    // rather than changed there.
333                    o.insert("joining_subscribe_id".into(), vi(joining_request_id.into_inner()));
334                    o.insert("joining_start".into(), vi(joining_start.into_inner()));
335                }
336            }
337            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
338            o
339        }
340        ControlMessage::FetchOk(m) => {
341            let mut o = Map::new();
342            o.insert("request_id".into(), vi(m.request_id.into_inner()));
343            o.insert("group_order".into(), vi(m.group_order as u64));
344            o.insert("end_of_track".into(), vi(m.end_of_track as u64));
345            o.insert("end_location".into(), loc_to_json(&m.end_location));
346            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
347            o
348        }
349        ControlMessage::FetchError(m) => {
350            let mut o = Map::new();
351            o.insert("request_id".into(), vi(m.request_id.into_inner()));
352            o.insert("error_code".into(), vi(m.error_code.into_inner()));
353            o.insert(
354                "reason_phrase".into(),
355                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
356            );
357            o
358        }
359        ControlMessage::FetchCancel(m) => {
360            let mut o = Map::new();
361            o.insert("request_id".into(), vi(m.request_id.into_inner()));
362            o
363        }
364        ControlMessage::Publish(m) => {
365            let mut o = Map::new();
366            o.insert("request_id".into(), vi(m.request_id.into_inner()));
367            o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
368            o.insert(
369                "track_name".into(),
370                Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
371            );
372            o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
373            o.insert("group_order".into(), vi(m.group_order as u64));
374            o.insert("content_exists".into(), vi(m.content_exists as u64));
375            if let Some(loc) = &m.largest_location {
376                o.insert("largest_location".into(), loc_to_json(loc));
377            }
378            o.insert("forward".into(), vi(m.forward as u64));
379            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
380            o
381        }
382        ControlMessage::PublishOk(m) => {
383            let mut o = Map::new();
384            o.insert("request_id".into(), vi(m.request_id.into_inner()));
385            o.insert("forward".into(), vi(m.forward as u64));
386            o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
387            o.insert("group_order".into(), vi(m.group_order as u64));
388            o.insert("filter_type".into(), vi(m.filter_type as u64));
389            if let Some(sg) = &m.start_group {
390                o.insert("start_group".into(), vi(sg.into_inner()));
391            }
392            if let Some(so) = &m.start_object {
393                o.insert("start_object".into(), vi(so.into_inner()));
394            }
395            if let Some(eg) = &m.end_group {
396                o.insert("end_group".into(), vi(eg.into_inner()));
397            }
398            o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
399            o
400        }
401        ControlMessage::PublishError(m) => {
402            let mut o = Map::new();
403            o.insert("request_id".into(), vi(m.request_id.into_inner()));
404            o.insert("error_code".into(), vi(m.error_code.into_inner()));
405            o.insert(
406                "reason_phrase".into(),
407                Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
408            );
409            o
410        }
411    };
412    obj
413}