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