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