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 crate::fields::kvp_entries(params, |key, value| {
129 let Some(name) = d18_param_name(key) else {
130 return (None, None);
131 };
132 let rendered = match (value, key) {
133 (KvpValue::Bytes(b), 0x21) => decode_subscription_filter(b),
134 (KvpValue::Bytes(b), 0x09) => decode_largest_object(b),
135 (KvpValue::Bytes(b), 0x34) => decode_track_namespace_prefix(b),
136 (KvpValue::Bytes(b), _) if name == "authorization_token" => auth_token_to_json_d18(b),
137 (KvpValue::Varint(v), _) => vi(v.into_inner()),
138 (KvpValue::Bytes(b), _) => Value::Text(String::from_utf8_lossy(b).into_owned()),
139 };
140 (Some(name), Some(rendered))
141 })
142}
143
144fn options_to_json(options: &[KeyValuePair]) -> Value {
145 crate::fields::kvp_entries(options, |key, value| {
146 let Some(name) = d18_option_name(key) else {
147 return (None, None);
148 };
149 let rendered = match value {
150 KvpValue::Varint(v) => vi(v.into_inner()),
151 KvpValue::Bytes(b) if name == "authorization_token" => auth_token_to_json_d18(b),
152 KvpValue::Bytes(b) => Value::Text(String::from_utf8_lossy(b).into_owned()),
153 };
154 (Some(name), Some(rendered))
155 })
156}
157
158fn d18_track_prop_name(key: u64) -> Option<&'static str> {
159 match key {
160 0x02 => Some("object_delivery_timeout"),
161 0x04 => Some("max_cache_duration"),
162 0x06 => Some("subgroup_delivery_timeout"),
163 0x0b => Some("immutable_properties"),
164 0x0e => Some("default_publisher_priority"),
165 0x22 => Some("default_publisher_group_order"),
166 0x30 => Some("dynamic_groups"),
167 _ => None,
168 }
169}
170
171fn track_props_to_json(props: &[KeyValuePair]) -> Value {
172 crate::fields::kvp_entries(props, |key, value| {
173 let name = d18_track_prop_name(key);
174 let rendered = match value {
175 KvpValue::Varint(v) => Some(vi(v.into_inner())),
176 KvpValue::Bytes(_) => None,
180 };
181 (name, rendered)
182 })
183}
184
185pub fn message_fields(msg: &ControlMessage) -> Map {
191 let obj = match msg {
192 ControlMessage::Setup(m) => {
193 let mut o = Map::new();
194 o.insert("options".into(), options_to_json(&m.options));
195 o
196 }
197 ControlMessage::GoAway(m) => {
198 let mut o = Map::new();
199 o.insert(
200 "new_session_uri".into(),
201 Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
202 );
203 o.insert("timeout".into(), vi(m.timeout.into_inner()));
204 if let Some(rid) = &m.request_id {
205 o.insert("request_id".into(), vi(rid.into_inner()));
206 }
207 o
208 }
209 ControlMessage::RequestOk(m) => {
210 let mut o = Map::new();
211 o.insert("parameters".into(), params_to_json(&m.parameters));
212 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
213 o
214 }
215 ControlMessage::RequestError(m) => {
216 let mut o = Map::new();
217 o.insert("error_code".into(), vi(m.error_code.into_inner()));
218 o.insert("retry_interval".into(), vi(m.retry_interval.into_inner()));
219 o.insert(
220 "reason_phrase".into(),
221 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
222 );
223 if let Some(r) = &m.redirect {
224 let mut r_obj = Map::new();
225 r_obj.insert(
226 "connect_uri".into(),
227 Value::Text(String::from_utf8_lossy(&r.connect_uri).into_owned()),
228 );
229 r_obj.insert("track_namespace".into(), ns_to_json(&r.track_namespace));
230 r_obj.insert(
231 "track_name".into(),
232 Value::Text(String::from_utf8_lossy(&r.track_name).into_owned()),
233 );
234 o.insert("redirect".into(), Value::Map(r_obj));
235 }
236 o
237 }
238 ControlMessage::Subscribe(m) => {
239 let mut o = Map::new();
240 o.insert("request_id".into(), vi(m.request_id.into_inner()));
241 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
242 o.insert(
243 "track_name".into(),
244 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
245 );
246 o.insert("parameters".into(), params_to_json(&m.parameters));
247 o
248 }
249 ControlMessage::SubscribeOk(m) => {
250 let mut o = Map::new();
251 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
252 o.insert("parameters".into(), params_to_json(&m.parameters));
253 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
254 o
255 }
256 ControlMessage::RequestUpdate(m) => {
257 let mut o = Map::new();
258 o.insert("request_id".into(), vi(m.request_id.into_inner()));
259 o.insert("parameters".into(), params_to_json(&m.parameters));
260 o
261 }
262 ControlMessage::Publish(m) => {
263 let mut o = Map::new();
264 o.insert("request_id".into(), vi(m.request_id.into_inner()));
265 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
266 o.insert(
267 "track_name".into(),
268 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
269 );
270 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
271 o.insert("parameters".into(), params_to_json(&m.parameters));
272 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
273 o
274 }
275 ControlMessage::PublishDone(m) => {
276 let mut o = Map::new();
277 o.insert("status_code".into(), vi(m.status_code.into_inner()));
278 o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
279 o.insert(
280 "reason_phrase".into(),
281 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
282 );
283 o
284 }
285 ControlMessage::PublishNamespace(m) => {
286 let mut o = Map::new();
287 o.insert("request_id".into(), vi(m.request_id.into_inner()));
288 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
289 o.insert("parameters".into(), params_to_json(&m.parameters));
290 o
291 }
292 ControlMessage::Namespace(m) => {
293 let mut o = Map::new();
294 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
295 o
296 }
297 ControlMessage::NamespaceDone(m) => {
298 let mut o = Map::new();
299 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
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.namespace_prefix));
306 o.insert("parameters".into(), params_to_json(&m.parameters));
307 o
308 }
309 ControlMessage::SubscribeTracks(m) => {
310 let mut o = Map::new();
311 o.insert("request_id".into(), vi(m.request_id.into_inner()));
312 o.insert("namespace_prefix".into(), ns_to_json(&m.namespace_prefix));
313 o.insert("parameters".into(), params_to_json(&m.parameters));
314 o
315 }
316 ControlMessage::TrackStatus(m) => {
317 let mut o = Map::new();
318 o.insert("request_id".into(), vi(m.request_id.into_inner()));
319 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
320 o.insert(
321 "track_name".into(),
322 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
323 );
324 o.insert("parameters".into(), params_to_json(&m.parameters));
325 o
326 }
327 ControlMessage::Fetch(m) => {
328 let mut o = Map::new();
329 o.insert("request_id".into(), vi(m.request_id.into_inner()));
330 o.insert("fetch_type".into(), vi(m.fetch_type as u64));
331 match &m.fetch_payload {
332 crate::draft18::message::FetchPayload::Standalone {
333 track_namespace,
334 track_name,
335 start_group,
336 start_object,
337 end_group,
338 end_object,
339 } => {
340 o.insert("track_namespace".into(), ns_to_json(track_namespace));
341 o.insert(
342 "track_name".into(),
343 Value::Text(String::from_utf8_lossy(track_name).into_owned()),
344 );
345 o.insert("start_group".into(), vi(start_group.into_inner()));
346 o.insert("start_object".into(), vi(start_object.into_inner()));
347 o.insert("end_group".into(), vi(end_group.into_inner()));
348 o.insert("end_object".into(), vi(end_object.into_inner()));
349 }
350 crate::draft18::message::FetchPayload::Joining {
351 joining_request_id,
352 joining_start,
353 } => {
354 o.insert("joining_request_id".into(), vi(joining_request_id.into_inner()));
355 o.insert("joining_start".into(), vi(joining_start.into_inner()));
356 }
357 }
358 o.insert("parameters".into(), params_to_json(&m.parameters));
359 o
360 }
361 ControlMessage::FetchOk(m) => {
362 let mut o = Map::new();
363 o.insert("end_of_track".into(), vi(m.end_of_track as u64));
364 o.insert("end_group".into(), vi(m.end_group.into_inner()));
365 o.insert("end_object".into(), vi(m.end_object.into_inner()));
366 o.insert("parameters".into(), params_to_json(&m.parameters));
367 o.insert("track_properties".into(), track_props_to_json(&m.track_properties));
368 o
369 }
370 ControlMessage::PublishBlocked(m) => {
371 let mut o = Map::new();
372 o.insert("namespace_suffix".into(), ns_to_json(&m.namespace_suffix));
373 o.insert(
374 "track_name".into(),
375 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
376 );
377 o
378 }
379 };
380 obj
381}