1use crate::draft17::message::ControlMessage;
2use crate::fields::{FieldMap as Map, FieldValue as Value};
3use crate::kvp::{KeyValuePair, KvpValue};
4use crate::types::*;
5use crate::varint::{Moqt17 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 d17_param_name(key: u64) -> Option<&'static str> {
19 match key {
20 0x02 => Some("delivery_timeout"),
21 0x03 => Some("authorization_token"),
22 0x04 => Some("rendezvous_timeout"),
23 0x08 => Some("expires"),
24 0x09 => Some("largest_object"),
25 0x10 => Some("forward"),
26 0x20 => Some("subscriber_priority"),
27 0x21 => Some("subscription_filter"),
28 0x22 => Some("group_order"),
29 0x32 => Some("new_group_request"),
30 _ => None,
31 }
32}
33
34fn d17_option_name(key: u64) -> Option<&'static str> {
36 match key {
37 0x01 => Some("path"),
38 0x03 => Some("authorization_token"),
39 0x04 => Some("max_auth_token_cache_size"),
40 0x05 => Some("authority"),
41 0x07 => Some("moqt_implementation"),
42 _ => None,
43 }
44}
45
46fn decode_subscription_filter(bytes: &[u8]) -> Value {
47 let mut buf = bytes;
48 let filter_type = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
49 let mut obj = Map::new();
50 obj.insert("filter_type".into(), vi(filter_type));
51 match filter_type {
52 3 => {
53 let start_group = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
54 let start_object = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
55 obj.insert("start_group".into(), vi(start_group));
56 obj.insert("start_object".into(), vi(start_object));
57 }
58 4 => {
59 let start_group = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
60 let start_object = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
61 let end_group = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
62 obj.insert("start_group".into(), vi(start_group));
63 obj.insert("start_object".into(), vi(start_object));
64 obj.insert("end_group".into(), vi(end_group));
65 }
66 _ => {}
67 }
68 Value::Map(obj)
69}
70
71fn auth_token_to_json_d17(bytes: &[u8]) -> Value {
72 let mut buf = bytes;
73 let alias_type = match VarInt::decode_moqt::<Wire>(&mut buf) {
74 Ok(v) => v,
75 Err(_) => return Value::Bytes(bytes.to_vec()),
76 };
77 let at = alias_type.into_inner();
78 let mut o = Map::new();
79 o.insert("alias_type".into(), vi(at));
80 match at {
81 0 | 2 => {
82 if let Ok(ta) = VarInt::decode_moqt::<Wire>(&mut buf) {
83 o.insert("token_alias".into(), vi(ta.into_inner()));
84 }
85 }
86 1 => {
87 if let Ok(ta) = VarInt::decode_moqt::<Wire>(&mut buf) {
88 o.insert("token_alias".into(), vi(ta.into_inner()));
89 }
90 if let Ok(tt) = VarInt::decode_moqt::<Wire>(&mut buf) {
91 o.insert("token_type".into(), vi(tt.into_inner()));
92 }
93 let tv = match VarInt::decode_moqt::<Wire>(&mut buf) {
95 Ok(len) => {
96 let n = len.into_inner() as usize;
97 if buf.len() >= n {
98 &buf[..n]
99 } else {
100 buf
101 }
102 }
103 Err(_) => buf,
104 };
105 o.insert("token_value".into(), Value::Bytes(tv.to_vec()));
106 }
107 _ => {
108 if let Ok(tt) = VarInt::decode_moqt::<Wire>(&mut buf) {
109 o.insert("token_type".into(), vi(tt.into_inner()));
110 }
111 let tv = match VarInt::decode_moqt::<Wire>(&mut buf) {
112 Ok(len) => {
113 let n = len.into_inner() as usize;
114 if buf.len() >= n {
115 &buf[..n]
116 } else {
117 buf
118 }
119 }
120 Err(_) => buf,
121 };
122 o.insert("token_value".into(), Value::Bytes(tv.to_vec()));
123 }
124 }
125 Value::Map(o)
126}
127
128fn decode_largest_object(bytes: &[u8]) -> Value {
129 let mut buf = bytes;
130 let group = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
131 let object = VarInt::decode_moqt::<Wire>(&mut buf).unwrap().into_inner();
132 let mut obj = Map::new();
133 obj.insert("group".into(), vi(group));
134 obj.insert("object".into(), vi(object));
135 Value::Map(obj)
136}
137
138fn params_to_json(params: &[KeyValuePair]) -> Value {
139 let mut obj = Map::new();
140 let mut unknown = Vec::new();
141
142 for p in params {
143 let key = p.key.into_inner();
144 if let Some(name) = d17_param_name(key) {
145 match (&p.value, key) {
146 (KvpValue::Bytes(b), 0x21) => {
147 obj.insert(name.to_string(), decode_subscription_filter(b));
148 }
149 (KvpValue::Bytes(b), 0x09) => {
150 obj.insert(name.to_string(), decode_largest_object(b));
151 }
152 (KvpValue::Bytes(b), _) if name == "authorization_token" => {
153 obj.insert(name.to_string(), auth_token_to_json_d17(b));
154 }
155 (KvpValue::Varint(v), _) => {
156 obj.insert(name.to_string(), vi(v.into_inner()));
157 }
158 (KvpValue::Bytes(b), _) => {
159 obj.insert(
160 name.to_string(),
161 Value::Text(String::from_utf8_lossy(b).into_owned()),
162 );
163 }
164 }
165 } else {
166 let mut entry = Map::new();
167 entry.insert("id".to_string(), Value::Text(format!("0x{:x}", key)));
168 match &p.value {
169 KvpValue::Varint(v) => {
170 entry.insert("length".to_string(), vi(v.into_inner()));
171 }
172 KvpValue::Bytes(b) => {
173 entry.insert("length".to_string(), vi(b.len() as u64));
174 entry.insert("raw_hex".to_string(), Value::Bytes(b.to_vec()));
175 }
176 }
177 unknown.push(Value::Map(entry));
178 }
179 }
180
181 if !unknown.is_empty() {
182 obj.insert("unknown".to_string(), Value::Array(unknown));
183 }
184
185 Value::Map(obj)
186}
187
188fn options_to_json(options: &[KeyValuePair]) -> Value {
189 let mut obj = Map::new();
190 for p in options {
191 let key = p.key.into_inner();
192 if let Some(name) = d17_option_name(key) {
193 match &p.value {
194 KvpValue::Varint(v) => {
195 obj.insert(name.to_string(), vi(v.into_inner()));
196 }
197 KvpValue::Bytes(b) if name == "authorization_token" => {
198 obj.insert(name.to_string(), auth_token_to_json_d17(b));
199 }
200 KvpValue::Bytes(b) => {
201 obj.insert(
202 name.to_string(),
203 Value::Text(String::from_utf8_lossy(b).into_owned()),
204 );
205 }
206 }
207 }
208 }
209 Value::Map(obj)
210}
211
212fn d17_track_prop_name(key: u64) -> Option<&'static str> {
213 match key {
214 0x02 => Some("delivery_timeout"),
215 0x04 => Some("max_cache_duration"),
216 0x0b => Some("immutable_properties"),
217 0x0e => Some("default_publisher_priority"),
218 0x22 => Some("default_publisher_group_order"),
219 0x30 => Some("dynamic_groups"),
220 _ => None,
221 }
222}
223
224fn track_props_to_json(props: &[KeyValuePair]) -> Value {
225 let mut obj = Map::new();
226 for p in props {
227 let key = p.key.into_inner();
228 let name = d17_track_prop_name(key)
229 .map(|s| s.to_string())
230 .unwrap_or_else(|| format!("0x{:x}", key));
231 match &p.value {
232 KvpValue::Varint(v) => {
233 obj.insert(name, vi(v.into_inner()));
234 }
235 KvpValue::Bytes(b) => {
236 obj.insert(name, Value::Bytes(b.to_vec()));
237 }
238 }
239 }
240 Value::Map(obj)
241}
242
243pub fn message_fields(msg: &ControlMessage) -> Map {
249 let obj = match msg {
250 ControlMessage::Setup(m) => {
251 let mut o = Map::new();
252 o.insert("options".into(), options_to_json(&m.options));
253 o
254 }
255 ControlMessage::GoAway(m) => {
256 let mut o = Map::new();
257 o.insert(
258 "new_session_uri".into(),
259 Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
260 );
261 o.insert("timeout".into(), vi(m.timeout.into_inner()));
262 o
263 }
264 ControlMessage::RequestOk(m) => {
265 let mut o = Map::new();
266 o.insert("parameters".into(), params_to_json(&m.parameters));
267 o
268 }
269 ControlMessage::RequestError(m) => {
270 let mut o = Map::new();
271 o.insert("error_code".into(), vi(m.error_code.into_inner()));
272 o.insert("retry_interval".into(), vi(m.retry_interval.into_inner()));
273 o.insert(
274 "reason_phrase".into(),
275 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
276 );
277 o
278 }
279 ControlMessage::Subscribe(m) => {
280 let mut o = Map::new();
281 o.insert("request_id".into(), vi(m.request_id.into_inner()));
282 o.insert(
283 "required_request_id_delta".into(),
284 vi(m.required_request_id_delta.into_inner()),
285 );
286 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
287 o.insert(
288 "track_name".into(),
289 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
290 );
291 o.insert("parameters".into(), params_to_json(&m.parameters));
292 o
293 }
294 ControlMessage::SubscribeOk(m) => {
295 let mut o = Map::new();
296 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
297 o.insert("parameters".into(), params_to_json(&m.parameters));
298 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
299 o
300 }
301 ControlMessage::RequestUpdate(m) => {
302 let mut o = Map::new();
303 o.insert("request_id".into(), vi(m.request_id.into_inner()));
304 o.insert(
305 "required_request_id_delta".into(),
306 vi(m.required_request_id_delta.into_inner()),
307 );
308 o.insert("parameters".into(), params_to_json(&m.parameters));
309 o
310 }
311 ControlMessage::Publish(m) => {
312 let mut o = Map::new();
313 o.insert("request_id".into(), vi(m.request_id.into_inner()));
314 o.insert(
315 "required_request_id_delta".into(),
316 vi(m.required_request_id_delta.into_inner()),
317 );
318 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
319 o.insert(
320 "track_name".into(),
321 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
322 );
323 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
324 o.insert("parameters".into(), params_to_json(&m.parameters));
325 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
326 o
327 }
328 ControlMessage::PublishOk(m) => {
329 let mut o = Map::new();
330 o.insert("parameters".into(), params_to_json(&m.parameters));
331 o
332 }
333 ControlMessage::PublishDone(m) => {
334 let mut o = Map::new();
335 o.insert("status_code".into(), vi(m.status_code.into_inner()));
336 o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
337 o.insert(
338 "reason_phrase".into(),
339 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
340 );
341 o
342 }
343 ControlMessage::PublishNamespace(m) => {
344 let mut o = Map::new();
345 o.insert("request_id".into(), vi(m.request_id.into_inner()));
346 o.insert(
347 "required_request_id_delta".into(),
348 vi(m.required_request_id_delta.into_inner()),
349 );
350 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
351 o.insert("parameters".into(), params_to_json(&m.parameters));
352 o
353 }
354 ControlMessage::Namespace(m) => {
355 let mut o = Map::new();
356 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
357 o
358 }
359 ControlMessage::NamespaceDone(m) => {
360 let mut o = Map::new();
361 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
362 o
363 }
364 ControlMessage::SubscribeNamespace(m) => {
365 let mut o = Map::new();
366 o.insert("request_id".into(), vi(m.request_id.into_inner()));
367 o.insert(
368 "required_request_id_delta".into(),
369 vi(m.required_request_id_delta.into_inner()),
370 );
371 o.insert("namespace_prefix".into(), ns_to_json(&m.namespace_prefix));
372 o.insert("subscribe_options".into(), vi(m.subscribe_options.into_inner()));
373 o.insert("parameters".into(), params_to_json(&m.parameters));
374 o
375 }
376 ControlMessage::TrackStatus(m) => {
377 let mut o = Map::new();
378 o.insert("request_id".into(), vi(m.request_id.into_inner()));
379 o.insert(
380 "required_request_id_delta".into(),
381 vi(m.required_request_id_delta.into_inner()),
382 );
383 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
384 o.insert(
385 "track_name".into(),
386 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
387 );
388 o.insert("parameters".into(), params_to_json(&m.parameters));
389 o
390 }
391 ControlMessage::Fetch(m) => {
392 let mut o = Map::new();
393 o.insert("request_id".into(), vi(m.request_id.into_inner()));
394 o.insert(
395 "required_request_id_delta".into(),
396 vi(m.required_request_id_delta.into_inner()),
397 );
398 o.insert("fetch_type".into(), vi(m.fetch_type as u64));
399 match &m.fetch_payload {
400 crate::draft17::message::FetchPayload::Standalone {
401 track_namespace,
402 track_name,
403 start_group,
404 start_object,
405 end_group,
406 end_object,
407 } => {
408 o.insert("track_namespace".into(), ns_to_json(track_namespace));
409 o.insert(
410 "track_name".into(),
411 Value::Text(String::from_utf8_lossy(track_name).into_owned()),
412 );
413 o.insert("start_group".into(), vi(start_group.into_inner()));
414 o.insert("start_object".into(), vi(start_object.into_inner()));
415 o.insert("end_group".into(), vi(end_group.into_inner()));
416 o.insert("end_object".into(), vi(end_object.into_inner()));
417 }
418 crate::draft17::message::FetchPayload::Joining {
419 joining_request_id,
420 joining_start,
421 } => {
422 o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
423 o.insert("joining_start".into(), vi(joining_start.into_inner()));
424 }
425 }
426 o.insert("parameters".into(), params_to_json(&m.parameters));
427 o
428 }
429 ControlMessage::FetchOk(m) => {
430 let mut o = Map::new();
431 o.insert("end_of_track".into(), vi(m.end_of_track as u64));
432 o.insert("end_group".into(), vi(m.end_group.into_inner()));
433 o.insert("end_object".into(), vi(m.end_object.into_inner()));
434 o.insert("parameters".into(), params_to_json(&m.parameters));
435 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
436 o
437 }
438 ControlMessage::PublishBlocked(m) => {
439 let mut o = Map::new();
440 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
441 o.insert(
442 "track_name".into(),
443 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
444 );
445 o
446 }
447 };
448 obj
449}