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 {
73 let mut buf = bytes;
74 let Ok(filter_type) = VarInt::decode_moqt::<Wire>(&mut buf) else {
75 return Value::Bytes(bytes.to_vec());
76 };
77 let filter_type = filter_type.into_inner();
78 let mut obj = Map::new();
79 obj.insert("filter_type".into(), vi(filter_type));
80 match filter_type {
81 3 => {
82 let Ok(start_group) = VarInt::decode_moqt::<Wire>(&mut buf) else {
83 return Value::Bytes(bytes.to_vec());
84 };
85 let start_group = start_group.into_inner();
86 let Ok(start_object) = VarInt::decode_moqt::<Wire>(&mut buf) else {
87 return Value::Bytes(bytes.to_vec());
88 };
89 let start_object = start_object.into_inner();
90 obj.insert("start_group".into(), vi(start_group));
91 obj.insert("start_object".into(), vi(start_object));
92 }
93 4 => {
94 let Ok(start_group) = VarInt::decode_moqt::<Wire>(&mut buf) else {
95 return Value::Bytes(bytes.to_vec());
96 };
97 let start_group = start_group.into_inner();
98 let Ok(start_object) = VarInt::decode_moqt::<Wire>(&mut buf) else {
99 return Value::Bytes(bytes.to_vec());
100 };
101 let start_object = start_object.into_inner();
102 let Ok(end_group) = VarInt::decode_moqt::<Wire>(&mut buf) else {
103 return Value::Bytes(bytes.to_vec());
104 };
105 let end_group = end_group.into_inner();
106 obj.insert("start_group".into(), vi(start_group));
107 obj.insert("start_object".into(), vi(start_object));
108 obj.insert("end_group".into(), vi(end_group));
109 }
110 _ => {}
111 }
112 Value::Map(obj)
113}
114
115fn auth_token_to_json_d18(bytes: &[u8]) -> Value {
116 let mut buf = bytes;
117 let alias_type = match VarInt::decode_moqt::<Wire>(&mut buf) {
118 Ok(v) => v,
119 Err(_) => return Value::Bytes(bytes.to_vec()),
120 };
121 let at = alias_type.into_inner();
122 let mut o = Map::new();
123 o.insert("alias_type".into(), vi(at));
124 match at {
125 0 | 2 => {
126 if let Ok(ta) = VarInt::decode_moqt::<Wire>(&mut buf) {
127 o.insert("token_alias".into(), vi(ta.into_inner()));
128 }
129 }
130 1 => {
131 if let Ok(ta) = VarInt::decode_moqt::<Wire>(&mut buf) {
132 o.insert("token_alias".into(), vi(ta.into_inner()));
133 }
134 if let Ok(tt) = VarInt::decode_moqt::<Wire>(&mut buf) {
135 o.insert("token_type".into(), vi(tt.into_inner()));
136 }
137 o.insert("token_value".into(), Value::Bytes(buf.to_vec()));
139 }
140 _ => {
141 if let Ok(tt) = VarInt::decode_moqt::<Wire>(&mut buf) {
142 o.insert("token_type".into(), vi(tt.into_inner()));
143 }
144 o.insert("token_value".into(), Value::Bytes(buf.to_vec()));
145 }
146 }
147 Value::Map(o)
148}
149
150fn decode_largest_object(bytes: &[u8]) -> Value {
163 let mut buf = bytes;
164 let Ok(group) = VarInt::decode_moqt::<Wire>(&mut buf) else {
165 return Value::Bytes(bytes.to_vec());
166 };
167 let Ok(object) = VarInt::decode_moqt::<Wire>(&mut buf) else {
168 return Value::Bytes(bytes.to_vec());
169 };
170 let (group, object) = (group.into_inner(), object.into_inner());
171 let mut obj = Map::new();
172 obj.insert("group".into(), vi(group));
173 obj.insert("object".into(), vi(object));
174 Value::Map(obj)
175}
176
177fn decode_track_namespace_prefix(bytes: &[u8]) -> Value {
178 let mut buf = bytes;
179 match TrackNamespace::decode_allow_empty_moqt::<Wire>(&mut buf) {
180 Ok(ns) => ns_to_json(&ns),
181 Err(_) => Value::Bytes(bytes.to_vec()),
182 }
183}
184
185fn params_to_json(params: &[KeyValuePair]) -> Value {
186 crate::fields::kvp_entries(params, |key, value| {
187 let Some(name) = d18_param_name(key) else {
188 return (None, None);
189 };
190 let rendered = match (value, key) {
191 (KvpValue::Bytes(b), 0x21) => decode_subscription_filter(b),
192 (KvpValue::Bytes(b), 0x09) => decode_largest_object(b),
193 (KvpValue::Bytes(b), 0x34) => decode_track_namespace_prefix(b),
194 (KvpValue::Bytes(b), _) if name == "authorization_token" => auth_token_to_json_d18(b),
195 (KvpValue::Varint(v), _) => vi(v.into_inner()),
196 (KvpValue::Bytes(b), _) => Value::Text(String::from_utf8_lossy(b).into_owned()),
197 };
198 (Some(name), Some(rendered))
199 })
200}
201
202pub(crate) fn options_to_json(options: &[KeyValuePair]) -> Value {
203 crate::fields::kvp_entries(options, |key, value| {
204 let Some(name) = d18_option_name(key) else {
205 return (None, None);
206 };
207 let rendered = match value {
208 KvpValue::Varint(v) => vi(v.into_inner()),
209 KvpValue::Bytes(b) if name == "authorization_token" => auth_token_to_json_d18(b),
210 KvpValue::Bytes(b) => Value::Text(String::from_utf8_lossy(b).into_owned()),
211 };
212 (Some(name), Some(rendered))
213 })
214}
215
216fn d18_track_prop_name(key: u64) -> Option<&'static str> {
217 match key {
218 0x02 => Some("object_delivery_timeout"),
219 0x04 => Some("max_cache_duration"),
220 0x06 => Some("subgroup_delivery_timeout"),
221 0x0b => Some("immutable_properties"),
222 0x0e => Some("default_publisher_priority"),
223 0x22 => Some("default_publisher_group_order"),
224 0x30 => Some("dynamic_groups"),
225 _ => None,
226 }
227}
228
229fn track_props_to_json(props: &[KeyValuePair]) -> Value {
230 crate::fields::kvp_entries(props, |key, value| {
231 let name = d18_track_prop_name(key);
232 let rendered = match value {
233 KvpValue::Varint(v) => Some(vi(v.into_inner())),
234 KvpValue::Bytes(_) => None,
238 };
239 (name, rendered)
240 })
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 if let Some(rid) = &m.request_id {
263 o.insert("request_id".into(), vi(rid.into_inner()));
264 }
265 o
266 }
267 ControlMessage::RequestOk(m) => {
268 let mut o = Map::new();
269 o.insert("parameters".into(), params_to_json(&m.parameters));
270 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
271 o
272 }
273 ControlMessage::RequestError(m) => {
274 let mut o = Map::new();
275 o.insert("error_code".into(), vi(m.error_code.into_inner()));
276 o.insert("retry_interval".into(), vi(m.retry_interval.into_inner()));
277 o.insert(
278 "reason_phrase".into(),
279 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
280 );
281 if let Some(r) = &m.redirect {
282 let mut r_obj = Map::new();
283 r_obj.insert(
284 "connect_uri".into(),
285 Value::Text(String::from_utf8_lossy(&r.connect_uri).into_owned()),
286 );
287 r_obj.insert("track_namespace".into(), ns_to_json(&r.track_namespace));
288 r_obj.insert(
289 "track_name".into(),
290 Value::Text(String::from_utf8_lossy(&r.track_name).into_owned()),
291 );
292 o.insert("redirect".into(), Value::Map(r_obj));
293 }
294 o
295 }
296 ControlMessage::Subscribe(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(
301 "track_name".into(),
302 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
303 );
304 o.insert("parameters".into(), params_to_json(&m.parameters));
305 o
306 }
307 ControlMessage::SubscribeOk(m) => {
308 let mut o = Map::new();
309 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
310 o.insert("parameters".into(), params_to_json(&m.parameters));
311 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
312 o
313 }
314 ControlMessage::RequestUpdate(m) => {
315 let mut o = Map::new();
316 o.insert("request_id".into(), vi(m.request_id.into_inner()));
317 o.insert("parameters".into(), params_to_json(&m.parameters));
318 o
319 }
320 ControlMessage::Publish(m) => {
321 let mut o = Map::new();
322 o.insert("request_id".into(), vi(m.request_id.into_inner()));
323 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
324 o.insert(
325 "track_name".into(),
326 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
327 );
328 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
329 o.insert("parameters".into(), params_to_json(&m.parameters));
330 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
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("track_namespace".into(), ns_to_json(&m.track_namespace));
347 o.insert("parameters".into(), params_to_json(&m.parameters));
348 o
349 }
350 ControlMessage::Namespace(m) => {
351 let mut o = Map::new();
352 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
353 o
354 }
355 ControlMessage::NamespaceDone(m) => {
356 let mut o = Map::new();
357 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
358 o
359 }
360 ControlMessage::SubscribeNamespace(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::SubscribeTracks(m) => {
368 let mut o = Map::new();
369 o.insert("request_id".into(), vi(m.request_id.into_inner()));
370 o.insert("namespace_prefix".into(), ns_to_json(&m.namespace_prefix));
371 o.insert("parameters".into(), params_to_json(&m.parameters));
372 o
373 }
374 ControlMessage::TrackStatus(m) => {
375 let mut o = Map::new();
376 o.insert("request_id".into(), vi(m.request_id.into_inner()));
377 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
378 o.insert(
379 "track_name".into(),
380 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
381 );
382 o.insert("parameters".into(), params_to_json(&m.parameters));
383 o
384 }
385 ControlMessage::Fetch(m) => {
386 let mut o = Map::new();
387 o.insert("request_id".into(), vi(m.request_id.into_inner()));
388 o.insert("fetch_type".into(), vi(m.fetch_type as u64));
389 match &m.fetch_payload {
390 crate::draft18::message::FetchPayload::Standalone {
391 track_namespace,
392 track_name,
393 start_group,
394 start_object,
395 end_group,
396 end_object,
397 } => {
398 o.insert("track_namespace".into(), ns_to_json(track_namespace));
399 o.insert(
400 "track_name".into(),
401 Value::Text(String::from_utf8_lossy(track_name).into_owned()),
402 );
403 o.insert("start_group".into(), vi(start_group.into_inner()));
404 o.insert("start_object".into(), vi(start_object.into_inner()));
405 o.insert("end_group".into(), vi(end_group.into_inner()));
406 o.insert("end_object".into(), vi(end_object.into_inner()));
407 }
408 crate::draft18::message::FetchPayload::Joining {
409 joining_request_id,
410 joining_start,
411 } => {
412 o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
413 o.insert("joining_start".into(), vi(joining_start.into_inner()));
414 }
415 }
416 o.insert("parameters".into(), params_to_json(&m.parameters));
417 o
418 }
419 ControlMessage::FetchOk(m) => {
420 let mut o = Map::new();
421 o.insert("end_of_track".into(), vi(m.end_of_track as u64));
422 o.insert("end_group".into(), vi(m.end_group.into_inner()));
423 o.insert("end_object".into(), vi(m.end_object.into_inner()));
424 o.insert("parameters".into(), params_to_json(&m.parameters));
425 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
426 o
427 }
428 ControlMessage::PublishBlocked(m) => {
429 let mut o = Map::new();
430 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
431 o.insert(
432 "track_name".into(),
433 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
434 );
435 o
436 }
437 };
438 obj
439}