1use crate::draft11::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 {
26 let mut buf = bytes;
27 let alias_type = VarInt::decode(&mut buf).unwrap();
28 let token_type = VarInt::decode(&mut buf).unwrap();
29 let token_value = buf; let mut o = Map::new();
31 o.insert("alias_type".into(), vi(alias_type.into_inner()));
32 o.insert("token_type".into(), vi(token_type.into_inner()));
33 o.insert("token_value".into(), Value::Bytes(token_value.to_vec()));
34 Value::Map(o)
35}
36
37fn kvp_to_json_setup(params: &[KeyValuePair]) -> Value {
39 let mut obj = Map::new();
40 for p in params {
41 let key = p.key.into_inner();
42 match (key, &p.value) {
43 (0x01, KvpValue::Bytes(b)) => {
44 obj.insert("path".into(), Value::Text(String::from_utf8_lossy(b).into_owned()));
45 }
46 (0x02, KvpValue::Varint(v)) => {
47 obj.insert("max_request_id".into(), vi(v.into_inner()));
48 }
49 _ => {}
50 }
51 }
52 Value::Map(obj)
53}
54
55fn kvp_to_json_msg(params: &[KeyValuePair]) -> Value {
57 let mut obj = Map::new();
58 for p in params {
59 let key = p.key.into_inner();
60 match (key, &p.value) {
61 (0x01, KvpValue::Bytes(b)) => {
62 obj.insert("authorization_token".into(), auth_token_to_json(b));
63 }
64 (0x02, KvpValue::Varint(v)) => {
65 obj.insert("delivery_timeout".into(), vi(v.into_inner()));
66 }
67 (0x04, KvpValue::Varint(v)) => {
68 obj.insert("max_cache_duration".into(), vi(v.into_inner()));
69 }
70 _ => {}
71 }
72 }
73 Value::Map(obj)
74}
75
76pub fn message_fields(msg: &ControlMessage) -> Map {
82 let obj = match msg {
83 ControlMessage::ClientSetup(m) => {
84 let mut o = Map::new();
85 o.insert(
86 "supported_versions".into(),
87 Value::Array(m.supported_versions.iter().map(|v| vi(v.into_inner())).collect()),
88 );
89 o.insert("parameters".into(), kvp_to_json_setup(&m.parameters));
90 o
91 }
92 ControlMessage::ServerSetup(m) => {
93 let mut o = Map::new();
94 o.insert("selected_version".into(), vi(m.selected_version.into_inner()));
95 o.insert("parameters".into(), kvp_to_json_setup(&m.parameters));
96 o
97 }
98 ControlMessage::GoAway(m) => {
99 let mut o = Map::new();
100 o.insert(
101 "new_session_uri".into(),
102 Value::Text(String::from_utf8_lossy(&m.new_session_uri).into_owned()),
103 );
104 o
105 }
106 ControlMessage::MaxRequestId(m) => {
107 let mut o = Map::new();
108 o.insert("request_id".into(), vi(m.request_id.into_inner()));
109 o
110 }
111 ControlMessage::RequestsBlocked(m) => {
112 let mut o = Map::new();
113 o.insert("maximum_request_id".into(), vi(m.maximum_request_id.into_inner()));
114 o
115 }
116 ControlMessage::Subscribe(m) => {
117 let mut o = Map::new();
118 o.insert("request_id".into(), vi(m.request_id.into_inner()));
119 o.insert("track_alias".into(), vi(m.track_alias.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.into_inner()));
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("expires".into(), vi(m.expires.into_inner()));
145 o.insert("group_order".into(), vi(m.group_order as u64));
146 o.insert("content_exists".into(), vi(m.content_exists as u64));
147 if let Some(loc) = &m.largest_location {
148 o.insert("largest_location".into(), loc_to_json(loc));
149 }
150 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
151 o
152 }
153 ControlMessage::SubscribeError(m) => {
154 let mut o = Map::new();
155 o.insert("request_id".into(), vi(m.request_id.into_inner()));
156 o.insert("error_code".into(), vi(m.error_code.into_inner()));
157 o.insert(
158 "reason_phrase".into(),
159 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
160 );
161 o.insert("track_alias".into(), vi(m.track_alias.into_inner()));
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::SubscribeAnnounces(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::SubscribeAnnouncesOk(m) => {
236 let mut o = Map::new();
237 o.insert("request_id".into(), vi(m.request_id.into_inner()));
238 o
239 }
240 ControlMessage::SubscribeAnnouncesError(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::UnsubscribeAnnounces(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::TrackStatusRequest(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("parameters".into(), kvp_to_json_msg(&m.parameters));
264 o
265 }
266 ControlMessage::TrackStatus(m) => {
267 let mut o = Map::new();
268 o.insert("request_id".into(), vi(m.request_id.into_inner()));
269 o.insert("status_code".into(), vi(m.status_code.into_inner()));
270 o.insert("largest_location".into(), loc_to_json(&m.largest_location));
271 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
272 o
273 }
274 ControlMessage::Fetch(m) => {
275 let mut o = Map::new();
276 o.insert("request_id".into(), vi(m.request_id.into_inner()));
277 o.insert("subscriber_priority".into(), vi(m.subscriber_priority as u64));
278 o.insert("group_order".into(), vi(m.group_order as u64));
279 o.insert("fetch_type".into(), vi(m.fetch_type as u64));
280 match &m.fetch_payload {
281 FetchPayload::Standalone {
282 track_namespace,
283 track_name,
284 start_group,
285 start_object,
286 end_group,
287 end_object,
288 } => {
289 o.insert("track_namespace".into(), ns_to_json(track_namespace));
290 o.insert(
291 "track_name".into(),
292 Value::Text(String::from_utf8_lossy(track_name).into_owned()),
293 );
294 o.insert("start_group".into(), vi(start_group.into_inner()));
295 o.insert("start_object".into(), vi(start_object.into_inner()));
296 o.insert("end_group".into(), vi(end_group.into_inner()));
297 o.insert("end_object".into(), vi(end_object.into_inner()));
298 }
299 FetchPayload::Joining { joining_subscribe_id, joining_start } => {
300 o.insert("joining_subscribe_id".into(), vi(joining_subscribe_id.into_inner()));
301 o.insert("joining_start".into(), vi(joining_start.into_inner()));
302 }
303 }
304 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
305 o
306 }
307 ControlMessage::FetchOk(m) => {
308 let mut o = Map::new();
309 o.insert("request_id".into(), vi(m.request_id.into_inner()));
310 o.insert("group_order".into(), vi(m.group_order as u64));
311 o.insert("end_of_track".into(), vi(m.end_of_track as u64));
312 o.insert("end_location".into(), loc_to_json(&m.end_location));
313 o.insert("parameters".into(), kvp_to_json_msg(&m.parameters));
314 o
315 }
316 ControlMessage::FetchError(m) => {
317 let mut o = Map::new();
318 o.insert("request_id".into(), vi(m.request_id.into_inner()));
319 o.insert("error_code".into(), vi(m.error_code.into_inner()));
320 o.insert(
321 "reason_phrase".into(),
322 Value::Text(String::from_utf8_lossy(&m.reason_phrase).into_owned()),
323 );
324 o
325 }
326 ControlMessage::FetchCancel(m) => {
327 let mut o = Map::new();
328 o.insert("request_id".into(), vi(m.request_id.into_inner()));
329 o
330 }
331 };
332 obj
333}