1use crate::draft13::message::*;
2use crate::fields::{FieldMap as Map, FieldValue as Value};
3use crate::kvp::{KeyValuePair, KvpValue};
4use crate::types::*;
5use crate::varint::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 loc_to_json(loc: &Location) -> Value {
18 let mut o = Map::new();
19 o.insert("group".into(), vi(loc.group.into_inner()));
20 o.insert("object".into(), vi(loc.object.into_inner()));
21 Value::Map(o)
22}
23
24fn auth_token_to_json(bytes: &[u8]) -> Value {
25 let mut buf = bytes;
26 let alias_type = VarInt::decode(&mut buf).unwrap();
27 let at = alias_type.into_inner();
28 let mut o = Map::new();
29 o.insert("alias_type".into(), vi(at));
30 match at {
31 0 | 2 => {
33 let token_alias = VarInt::decode(&mut buf).unwrap();
34 o.insert("token_alias".into(), vi(token_alias.into_inner()));
35 }
36 1 => {
38 let token_alias = VarInt::decode(&mut buf).unwrap();
39 let token_type = VarInt::decode(&mut buf).unwrap();
40 o.insert("token_alias".into(), vi(token_alias.into_inner()));
41 o.insert("token_type".into(), vi(token_type.into_inner()));
42 o.insert("token_value".into(), Value::Bytes(buf.to_vec()));
43 }
44 _ => {
46 let token_type = VarInt::decode(&mut buf).unwrap();
47 o.insert("token_type".into(), vi(token_type.into_inner()));
48 o.insert("token_value".into(), Value::Bytes(buf.to_vec()));
49 }
50 }
51 Value::Map(o)
52}
53
54fn kvp_to_json_setup(params: &[KeyValuePair]) -> Value {
55 crate::fields::kvp_entries(params, |key, value| match (key, value) {
56 (0x01, KvpValue::Bytes(b)) => {
57 (Some("path"), Some(Value::Text(String::from_utf8_lossy(b).into_owned())))
58 }
59 (0x02, KvpValue::Varint(v)) => (Some("max_request_id"), Some(vi(v.into_inner()))),
60 (0x03, KvpValue::Bytes(b)) => (Some("authorization_token"), Some(auth_token_to_json(b))),
61 (0x04, KvpValue::Varint(v)) => {
62 (Some("max_auth_token_cache_size"), Some(vi(v.into_inner())))
63 }
64 _ => (None, None),
65 })
66}
67
68fn kvp_to_json_msg(params: &[KeyValuePair]) -> Value {
69 crate::fields::kvp_entries(params, |key, value| match (key, value) {
70 (0x03, KvpValue::Bytes(b)) => (Some("authorization_token"), Some(auth_token_to_json(b))),
71 (0x02, KvpValue::Varint(v)) => (Some("delivery_timeout"), Some(vi(v.into_inner()))),
72 (0x04, KvpValue::Varint(v)) => (Some("max_cache_duration"), Some(vi(v.into_inner()))),
73 _ => (None, None),
74 })
75}
76
77pub fn message_fields(msg: &ControlMessage) -> Map {
83 let obj = match msg {
84 ControlMessage::ClientSetup(m) => {
85 let mut o = Map::new();
86 o.insert(
87 "supported_versions".into(),
88 Value::Array(m.supported_versions.iter().map(|v| vi(v.into_inner())).collect()),
89 );
90 o.insert("parameters".into(), kvp_to_json_setup(&m.parameters));
91 o
92 }
93 ControlMessage::ServerSetup(m) => {
94 let mut o = Map::new();
95 o.insert("selected_version".into(), vi(m.selected_version.into_inner()));
96 o.insert("parameters".into(), kvp_to_json_setup(&m.parameters));
97 o
98 }
99 ControlMessage::GoAway(m) => {
100 let mut o = Map::new();
101 o.insert(
102 "new_session_uri".into(),
103 Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
104 );
105 o
106 }
107 ControlMessage::MaxRequestId(m) => {
108 let mut o = Map::new();
109 o.insert("request_id".into(), vi(m.request_id.into_inner()));
110 o
111 }
112 ControlMessage::RequestsBlocked(m) => {
113 let mut o = Map::new();
114 o.insert("maximum_request_id".into(), vi(m.maximum_request_id.into_inner()));
115 o
116 }
117 ControlMessage::Subscribe(m) => {
118 let mut o = Map::new();
119 o.insert("request_id".into(), vi(m.request_id.into_inner()));
120 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
121 o.insert(
122 "track_name".into(),
123 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
124 );
125 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
126 o.insert("group_order".into(), vi(m.group_order as u64));
127 o.insert("forward".into(), vi(m.forward as u64));
128 o.insert("filter_type".into(), vi(m.filter_type as u64));
129 if let Some(sg) = &m.start_group {
130 o.insert("start_group".into(), vi(sg.into_inner()));
131 }
132 if let Some(so) = &m.start_object {
133 o.insert("start_object".into(), vi(so.into_inner()));
134 }
135 if let Some(eg) = &m.end_group {
136 o.insert("end_group".into(), vi(eg.into_inner()));
137 }
138 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
139 o
140 }
141 ControlMessage::SubscribeOk(m) => {
142 let mut o = Map::new();
143 o.insert("request_id".into(), vi(m.request_id.into_inner()));
144 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
145 o.insert("expires".into(), vi(m.expires.into_inner()));
146 o.insert("group_order".into(), vi(m.group_order as u64));
147 o.insert("content_exists".into(), vi(m.content_exists as u64));
148 if let Some(loc) = &m.largest_location {
149 o.insert("largest_location".into(), loc_to_json(loc));
150 }
151 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
152 o
153 }
154 ControlMessage::SubscribeError(m) => {
155 let mut o = Map::new();
156 o.insert("request_id".into(), vi(m.request_id.into_inner()));
157 o.insert("error_code".into(), vi(m.error_code.into_inner()));
158 o.insert(
159 "reason_phrase".into(),
160 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
161 );
162 o
163 }
164 ControlMessage::SubscribeUpdate(m) => {
165 let mut o = Map::new();
166 o.insert("request_id".into(), vi(m.request_id.into_inner()));
167 o.insert("start_group".into(), vi(m.start_group.into_inner()));
168 o.insert("start_object".into(), vi(m.start_object.into_inner()));
169 o.insert("end_group".into(), vi(m.end_group.into_inner()));
170 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
171 o.insert("forward".into(), vi(m.forward as u64));
172 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
173 o
174 }
175 ControlMessage::SubscribeDone(m) => {
176 let mut o = Map::new();
177 o.insert("request_id".into(), vi(m.request_id.into_inner()));
178 o.insert("status_code".into(), vi(m.status_code.into_inner()));
179 o.insert("stream_count".into(), vi(m.stream_count.into_inner()));
180 o.insert(
181 "reason_phrase".into(),
182 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
183 );
184 o
185 }
186 ControlMessage::Unsubscribe(m) => {
187 let mut o = Map::new();
188 o.insert("request_id".into(), vi(m.request_id.into_inner()));
189 o
190 }
191 ControlMessage::Announce(m) => {
192 let mut o = Map::new();
193 o.insert("request_id".into(), vi(m.request_id.into_inner()));
194 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
195 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
196 o
197 }
198 ControlMessage::AnnounceOk(m) => {
199 let mut o = Map::new();
200 o.insert("request_id".into(), vi(m.request_id.into_inner()));
201 o
202 }
203 ControlMessage::AnnounceError(m) => {
204 let mut o = Map::new();
205 o.insert("request_id".into(), vi(m.request_id.into_inner()));
206 o.insert("error_code".into(), vi(m.error_code.into_inner()));
207 o.insert(
208 "reason_phrase".into(),
209 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
210 );
211 o
212 }
213 ControlMessage::AnnounceCancel(m) => {
214 let mut o = Map::new();
215 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
216 o.insert("error_code".into(), vi(m.error_code.into_inner()));
217 o.insert(
218 "reason_phrase".into(),
219 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
220 );
221 o
222 }
223 ControlMessage::Unannounce(m) => {
224 let mut o = Map::new();
225 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
226 o
227 }
228 ControlMessage::SubscribeNamespace(m) => {
229 let mut o = Map::new();
230 o.insert("request_id".into(), vi(m.request_id.into_inner()));
231 o.insert("track_namespace_prefix".into(), ns_to_json(&m.track_namespace_prefix));
232 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
233 o
234 }
235 ControlMessage::SubscribeNamespaceOk(m) => {
236 let mut o = Map::new();
237 o.insert("request_id".into(), vi(m.request_id.into_inner()));
238 o
239 }
240 ControlMessage::SubscribeNamespaceError(m) => {
241 let mut o = Map::new();
242 o.insert("request_id".into(), vi(m.request_id.into_inner()));
243 o.insert("error_code".into(), vi(m.error_code.into_inner()));
244 o.insert(
245 "reason_phrase".into(),
246 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
247 );
248 o
249 }
250 ControlMessage::UnsubscribeNamespace(m) => {
251 let mut o = Map::new();
252 o.insert("track_namespace_prefix".into(), ns_to_json(&m.track_namespace_prefix));
253 o
254 }
255 ControlMessage::TrackStatus(m) => {
256 let mut o = Map::new();
257 o.insert("request_id".into(), vi(m.request_id.into_inner()));
258 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
259 o.insert(
260 "track_name".into(),
261 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
262 );
263 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
264 o.insert("group_order".into(), vi(m.group_order as u64));
265 o.insert("forward".into(), vi(m.forward as u64));
266 o.insert("filter_type".into(), vi(m.filter_type as u64));
267 if let Some(sg) = &m.start_group {
268 o.insert("start_group".into(), vi(sg.into_inner()));
269 }
270 if let Some(so) = &m.start_object {
271 o.insert("start_object".into(), vi(so.into_inner()));
272 }
273 if let Some(eg) = &m.end_group {
274 o.insert("end_group".into(), vi(eg.into_inner()));
275 }
276 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
277 o
278 }
279 ControlMessage::TrackStatusOk(m) => {
280 let mut o = Map::new();
281 o.insert("request_id".into(), vi(m.request_id.into_inner()));
282 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
283 o.insert("expires".into(), vi(m.expires.into_inner()));
284 o.insert("group_order".into(), vi(m.group_order as u64));
285 o.insert("content_exists".into(), vi(m.content_exists as u64));
286 if let Some(loc) = &m.largest_location {
287 o.insert("largest_location".into(), loc_to_json(loc));
288 }
289 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
290 o
291 }
292 ControlMessage::TrackStatusError(m) => {
293 let mut o = Map::new();
294 o.insert("request_id".into(), vi(m.request_id.into_inner()));
295 o.insert("error_code".into(), vi(m.error_code.into_inner()));
296 o.insert(
297 "reason_phrase".into(),
298 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
299 );
300 o
301 }
302 ControlMessage::Fetch(m) => {
303 let mut o = Map::new();
304 o.insert("request_id".into(), vi(m.request_id.into_inner()));
305 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
306 o.insert("group_order".into(), vi(m.group_order as u64));
307 o.insert("fetch_type".into(), vi(m.fetch_type as u64));
308 match &m.fetch_payload {
309 FetchPayload::Standalone {
310 track_namespace,
311 track_name,
312 start_group,
313 start_object,
314 end_group,
315 end_object,
316 } => {
317 o.insert("track_namespace".into(), ns_to_json(track_namespace));
318 o.insert(
319 "track_name".into(),
320 Value::Text(String::from_utf8_lossy(track_name).into_owned()),
321 );
322 o.insert("start_group".into(), vi(start_group.into_inner()));
323 o.insert("start_object".into(), vi(start_object.into_inner()));
324 o.insert("end_group".into(), vi(end_group.into_inner()));
325 o.insert("end_object".into(), vi(end_object.into_inner()));
326 }
327 FetchPayload::Joining { joining_request_id, joining_start } => {
328 o.insert("joining_subscribe_id".into(), vi(joining_request_id.into_inner()));
334 o.insert("joining_start".into(), vi(joining_start.into_inner()));
335 }
336 }
337 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
338 o
339 }
340 ControlMessage::FetchOk(m) => {
341 let mut o = Map::new();
342 o.insert("request_id".into(), vi(m.request_id.into_inner()));
343 o.insert("group_order".into(), vi(m.group_order as u64));
344 o.insert("end_of_track".into(), vi(m.end_of_track as u64));
345 o.insert("end_location".into(), loc_to_json(&m.end_location));
346 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
347 o
348 }
349 ControlMessage::FetchError(m) => {
350 let mut o = Map::new();
351 o.insert("request_id".into(), vi(m.request_id.into_inner()));
352 o.insert("error_code".into(), vi(m.error_code.into_inner()));
353 o.insert(
354 "reason_phrase".into(),
355 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
356 );
357 o
358 }
359 ControlMessage::FetchCancel(m) => {
360 let mut o = Map::new();
361 o.insert("request_id".into(), vi(m.request_id.into_inner()));
362 o
363 }
364 ControlMessage::Publish(m) => {
365 let mut o = Map::new();
366 o.insert("request_id".into(), vi(m.request_id.into_inner()));
367 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
368 o.insert(
369 "track_name".into(),
370 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
371 );
372 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
373 o.insert("group_order".into(), vi(m.group_order as u64));
374 o.insert("content_exists".into(), vi(m.content_exists as u64));
375 if let Some(loc) = &m.largest_location {
376 o.insert("largest_location".into(), loc_to_json(loc));
377 }
378 o.insert("forward".into(), vi(m.forward as u64));
379 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
380 o
381 }
382 ControlMessage::PublishOk(m) => {
383 let mut o = Map::new();
384 o.insert("request_id".into(), vi(m.request_id.into_inner()));
385 o.insert("forward".into(), vi(m.forward as u64));
386 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
387 o.insert("group_order".into(), vi(m.group_order as u64));
388 o.insert("filter_type".into(), vi(m.filter_type as u64));
389 if let Some(sg) = &m.start_group {
390 o.insert("start_group".into(), vi(sg.into_inner()));
391 }
392 if let Some(so) = &m.start_object {
393 o.insert("start_object".into(), vi(so.into_inner()));
394 }
395 if let Some(eg) = &m.end_group {
396 o.insert("end_group".into(), vi(eg.into_inner()));
397 }
398 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
399 o
400 }
401 ControlMessage::PublishError(m) => {
402 let mut o = Map::new();
403 o.insert("request_id".into(), vi(m.request_id.into_inner()));
404 o.insert("error_code".into(), vi(m.error_code.into_inner()));
405 o.insert(
406 "reason_phrase".into(),
407 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
408 );
409 o
410 }
411 };
412 obj
413}