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