1use crate::draft12::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 token_type = VarInt::decode(&mut buf).unwrap();
28 let token_value = buf;
29 let mut o = Map::new();
30 o.insert("alias_type".into(), vi(alias_type.into_inner()));
31 o.insert("token_type".into(), vi(token_type.into_inner()));
32 o.insert("token_value".into(), Value::Bytes(token_value.to_vec()));
33 Value::Map(o)
34}
35
36fn kvp_to_json_setup(params: &[KeyValuePair]) -> Value {
37 crate::fields::kvp_entries(params, |key, value| match (key, value) {
38 (0x01, KvpValue::Bytes(b)) => {
39 (Some("path"), Some(Value::Text(String::from_utf8_lossy(b).into_owned())))
40 }
41 (0x02, KvpValue::Varint(v)) => (Some("max_request_id"), Some(vi(v.into_inner()))),
42 _ => (None, None),
43 })
44}
45
46fn kvp_to_json_msg(params: &[KeyValuePair]) -> Value {
47 crate::fields::kvp_entries(params, |key, value| match (key, value) {
48 (0x03, KvpValue::Bytes(b)) => (Some("authorization_token"), Some(auth_token_to_json(b))),
49 (0x02, KvpValue::Varint(v)) => (Some("delivery_timeout"), Some(vi(v.into_inner()))),
50 (0x04, KvpValue::Varint(v)) => (Some("max_cache_duration"), Some(vi(v.into_inner()))),
51 _ => (None, None),
52 })
53}
54
55pub fn message_fields(msg: &ControlMessage) -> Map {
61 let obj = match msg {
62 ControlMessage::ClientSetup(m) => {
63 let mut o = Map::new();
64 o.insert(
65 "supported_versions".into(),
66 Value::Array(m.supported_versions.iter().map(|v| vi(v.into_inner())).collect()),
67 );
68 o.insert("parameters".into(), kvp_to_json_setup(&m.parameters));
69 o
70 }
71 ControlMessage::ServerSetup(m) => {
72 let mut o = Map::new();
73 o.insert("selected_version".into(), vi(m.selected_version.into_inner()));
74 o.insert("parameters".into(), kvp_to_json_setup(&m.parameters));
75 o
76 }
77 ControlMessage::GoAway(m) => {
78 let mut o = Map::new();
79 o.insert(
80 "new_session_uri".into(),
81 Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
82 );
83 o
84 }
85 ControlMessage::MaxRequestId(m) => {
86 let mut o = Map::new();
87 o.insert("request_id".into(), vi(m.request_id.into_inner()));
88 o
89 }
90 ControlMessage::RequestsBlocked(m) => {
91 let mut o = Map::new();
92 o.insert("maximum_request_id".into(), vi(m.maximum_request_id.into_inner()));
93 o
94 }
95 ControlMessage::Subscribe(m) => {
96 let mut o = Map::new();
97 o.insert("request_id".into(), vi(m.request_id.into_inner()));
98 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
99 o.insert(
100 "track_name".into(),
101 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
102 );
103 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
104 o.insert("group_order".into(), vi(m.group_order as u64));
105 o.insert("forward".into(), vi(m.forward as u64));
106 o.insert("filter_type".into(), vi(m.filter_type.into_inner()));
107 if let Some(sg) = &m.start_group {
108 o.insert("start_group".into(), vi(sg.into_inner()));
109 }
110 if let Some(so) = &m.start_object {
111 o.insert("start_object".into(), vi(so.into_inner()));
112 }
113 if let Some(eg) = &m.end_group {
114 o.insert("end_group".into(), vi(eg.into_inner()));
115 }
116 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
117 o
118 }
119 ControlMessage::SubscribeOk(m) => {
120 let mut o = Map::new();
121 o.insert("request_id".into(), vi(m.request_id.into_inner()));
122 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
123 o.insert("expires".into(), vi(m.expires.into_inner()));
124 o.insert("group_order".into(), vi(m.group_order as u64));
125 o.insert("content_exists".into(), vi(m.content_exists as u64));
126 if let Some(loc) = &m.largest_location {
127 o.insert("largest_location".into(), loc_to_json(loc));
128 }
129 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
130 o
131 }
132 ControlMessage::SubscribeError(m) => {
133 let mut o = Map::new();
134 o.insert("request_id".into(), vi(m.request_id.into_inner()));
135 o.insert("error_code".into(), vi(m.error_code.into_inner()));
136 o.insert(
137 "reason_phrase".into(),
138 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
139 );
140 o
141 }
142 ControlMessage::SubscribeUpdate(m) => {
143 let mut o = Map::new();
144 o.insert("request_id".into(), vi(m.request_id.into_inner()));
145 o.insert("start_group".into(), vi(m.start_group.into_inner()));
146 o.insert("start_object".into(), vi(m.start_object.into_inner()));
147 o.insert("end_group".into(), vi(m.end_group.into_inner()));
148 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
149 o.insert("forward".into(), vi(m.forward as u64));
150 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
151 o
152 }
153 ControlMessage::SubscribeDone(m) => {
154 let mut o = Map::new();
155 o.insert("request_id".into(), vi(m.request_id.into_inner()));
156 o.insert("status_code".into(), vi(m.status_code.into_inner()));
157 o.insert("stream_count".into(), vi(m.stream_count.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::Unsubscribe(m) => {
165 let mut o = Map::new();
166 o.insert("request_id".into(), vi(m.request_id.into_inner()));
167 o
168 }
169 ControlMessage::Announce(m) => {
170 let mut o = Map::new();
171 o.insert("request_id".into(), vi(m.request_id.into_inner()));
172 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
173 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
174 o
175 }
176 ControlMessage::AnnounceOk(m) => {
177 let mut o = Map::new();
178 o.insert("request_id".into(), vi(m.request_id.into_inner()));
179 o
180 }
181 ControlMessage::AnnounceError(m) => {
182 let mut o = Map::new();
183 o.insert("request_id".into(), vi(m.request_id.into_inner()));
184 o.insert("error_code".into(), vi(m.error_code.into_inner()));
185 o.insert(
186 "reason_phrase".into(),
187 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
188 );
189 o
190 }
191 ControlMessage::AnnounceCancel(m) => {
192 let mut o = Map::new();
193 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
194 o.insert("error_code".into(), vi(m.error_code.into_inner()));
195 o.insert(
196 "reason_phrase".into(),
197 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
198 );
199 o
200 }
201 ControlMessage::Unannounce(m) => {
202 let mut o = Map::new();
203 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
204 o
205 }
206 ControlMessage::SubscribeAnnounces(m) => {
207 let mut o = Map::new();
208 o.insert("request_id".into(), vi(m.request_id.into_inner()));
209 o.insert("track_namespace_prefix".into(), ns_to_json(&m.track_namespace_prefix));
210 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
211 o
212 }
213 ControlMessage::SubscribeAnnouncesOk(m) => {
214 let mut o = Map::new();
215 o.insert("request_id".into(), vi(m.request_id.into_inner()));
216 o
217 }
218 ControlMessage::SubscribeAnnouncesError(m) => {
219 let mut o = Map::new();
220 o.insert("request_id".into(), vi(m.request_id.into_inner()));
221 o.insert("error_code".into(), vi(m.error_code.into_inner()));
222 o.insert(
223 "reason_phrase".into(),
224 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
225 );
226 o
227 }
228 ControlMessage::UnsubscribeAnnounces(m) => {
229 let mut o = Map::new();
230 o.insert("track_namespace_prefix".into(), ns_to_json(&m.track_namespace_prefix));
231 o
232 }
233 ControlMessage::TrackStatusRequest(m) => {
234 let mut o = Map::new();
235 o.insert("request_id".into(), vi(m.request_id.into_inner()));
236 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
237 o.insert(
238 "track_name".into(),
239 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
240 );
241 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
242 o
243 }
244 ControlMessage::TrackStatus(m) => {
245 let mut o = Map::new();
246 o.insert("request_id".into(), vi(m.request_id.into_inner()));
247 o.insert("status_code".into(), vi(m.status_code.into_inner()));
248 o.insert("largest_location".into(), loc_to_json(&m.largest_location));
249 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
250 o
251 }
252 ControlMessage::Fetch(m) => {
253 let mut o = Map::new();
254 o.insert("request_id".into(), vi(m.request_id.into_inner()));
255 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
256 o.insert("group_order".into(), vi(m.group_order as u64));
257 o.insert("fetch_type".into(), vi(m.fetch_type as u64));
258 match &m.fetch_payload {
259 FetchPayload::Standalone {
260 track_namespace,
261 track_name,
262 start_group,
263 start_object,
264 end_group,
265 end_object,
266 } => {
267 o.insert("track_namespace".into(), ns_to_json(track_namespace));
268 o.insert(
269 "track_name".into(),
270 Value::Text(String::from_utf8_lossy(track_name).into_owned()),
271 );
272 o.insert("start_group".into(), vi(start_group.into_inner()));
273 o.insert("start_object".into(), vi(start_object.into_inner()));
274 o.insert("end_group".into(), vi(end_group.into_inner()));
275 o.insert("end_object".into(), vi(end_object.into_inner()));
276 }
277 FetchPayload::Joining { joining_request_id, joining_start } => {
278 o.insert("joining_subscribe_id".into(), vi(joining_request_id.into_inner()));
284 o.insert("joining_start".into(), vi(joining_start.into_inner()));
285 }
286 }
287 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
288 o
289 }
290 ControlMessage::FetchOk(m) => {
291 let mut o = Map::new();
292 o.insert("request_id".into(), vi(m.request_id.into_inner()));
293 o.insert("group_order".into(), vi(m.group_order as u64));
294 o.insert("end_of_track".into(), vi(m.end_of_track as u64));
295 o.insert("end_location".into(), loc_to_json(&m.end_location));
296 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
297 o
298 }
299 ControlMessage::FetchError(m) => {
300 let mut o = Map::new();
301 o.insert("request_id".into(), vi(m.request_id.into_inner()));
302 o.insert("error_code".into(), vi(m.error_code.into_inner()));
303 o.insert(
304 "reason_phrase".into(),
305 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
306 );
307 o
308 }
309 ControlMessage::FetchCancel(m) => {
310 let mut o = Map::new();
311 o.insert("request_id".into(), vi(m.request_id.into_inner()));
312 o
313 }
314 ControlMessage::Publish(m) => {
315 let mut o = Map::new();
316 o.insert("request_id".into(), vi(m.request_id.into_inner()));
317 o.insert("track_namespace".into(), ns_to_json(&m.track_namespace));
318 o.insert(
319 "track_name".into(),
320 Value::Text(String::from_utf8_lossy(&m.track_name).into_owned()),
321 );
322 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
323 o.insert("group_order".into(), vi(m.group_order as u64));
324 o.insert("content_exists".into(), vi(m.content_exists as u64));
325 if let Some(loc) = &m.largest_location {
326 o.insert("largest_location".into(), loc_to_json(loc));
327 }
328 o.insert("forward".into(), vi(m.forward as u64));
329 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
330 o
331 }
332 ControlMessage::PublishOk(m) => {
333 let mut o = Map::new();
334 o.insert("request_id".into(), vi(m.request_id.into_inner()));
335 o.insert("forward".into(), vi(m.forward as u64));
336 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
337 o.insert("group_order".into(), vi(m.group_order as u64));
338 o.insert("filter_type".into(), vi(m.filter_type.into_inner()));
339 if let Some(sg) = &m.start_group {
340 o.insert("start_group".into(), vi(sg.into_inner()));
341 }
342 if let Some(so) = &m.start_object {
343 o.insert("start_object".into(), vi(so.into_inner()));
344 }
345 if let Some(eg) = &m.end_group {
346 o.insert("end_group".into(), vi(eg.into_inner()));
347 }
348 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
349 o
350 }
351 ControlMessage::PublishError(m) => {
352 let mut o = Map::new();
353 o.insert("request_id".into(), vi(m.request_id.into_inner()));
354 o.insert("error_code".into(), vi(m.error_code.into_inner()));
355 o.insert(
356 "reason_phrase".into(),
357 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
358 );
359 o
360 }
361 };
362 obj
363}