1use serde_json::{Map, Value};
6
7use crate::model::{Frame, Message, MsgKind, OpResult};
8use crate::segment::{ColumnValue, RowsSegment};
9use crate::util::base64;
10
11fn obj(entries: Vec<(&str, Value)>) -> Value {
12 let mut map = Map::new();
13 for (k, v) in entries {
14 map.insert(k.to_string(), v);
15 }
16 Value::Object(map)
17}
18
19fn scope_map_value(scopes: &[(String, Vec<String>)]) -> Value {
20 let mut map = Map::new();
21 for (k, vs) in scopes {
22 map.insert(
23 k.clone(),
24 Value::Array(vs.iter().map(|v| Value::String(v.clone())).collect()),
25 );
26 }
27 Value::Object(map)
28}
29
30pub fn render_message(msg: &Message) -> Value {
33 obj(vec![
34 ("magic", Value::from("SSP2")),
35 ("wireVersion", Value::from(1)),
36 (
37 "msgKind",
38 Value::from(match msg.msg_kind {
39 MsgKind::Request => "request",
40 MsgKind::Response => "response",
41 }),
42 ),
43 (
44 "frames",
45 Value::Array(msg.frames.iter().map(render_frame).collect()),
46 ),
47 ])
48}
49
50fn render_frame(frame: &Frame) -> Value {
51 let mut entries: Vec<(&str, Value)> = Vec::new();
55 match frame {
56 Frame::ReqHeader {
57 client_id,
58 schema_version,
59 } => {
60 entries.push(("type", Value::from("REQ_HEADER")));
61 entries.push(("clientId", Value::from(client_id.clone())));
62 entries.push(("schemaVersion", Value::from(*schema_version)));
63 }
64 Frame::PushCommit {
65 client_commit_id,
66 operations,
67 } => {
68 entries.push(("type", Value::from("PUSH_COMMIT")));
69 entries.push(("clientCommitId", Value::from(client_commit_id.clone())));
70 let ops = operations
71 .iter()
72 .map(|op| {
73 let mut e: Vec<(&str, Value)> = vec![
74 ("table", Value::from(op.table.clone())),
75 ("rowId", Value::from(op.row_id.clone())),
76 ("op", Value::from(op.op.name())),
77 ];
78 if let Some(v) = op.base_version {
79 e.push(("baseVersion", Value::from(v)));
80 }
81 if let Some(p) = &op.payload {
82 e.push(("payload", Value::from(base64(p))));
83 }
84 obj(e)
85 })
86 .collect();
87 entries.push(("operations", Value::Array(ops)));
88 }
89 Frame::PullHeader {
90 limit_commits,
91 limit_snapshot_rows,
92 max_snapshot_pages,
93 accept,
94 } => {
95 entries.push(("type", Value::from("PULL_HEADER")));
96 entries.push(("limitCommits", Value::from(*limit_commits)));
97 entries.push(("limitSnapshotRows", Value::from(*limit_snapshot_rows)));
98 entries.push(("maxSnapshotPages", Value::from(*max_snapshot_pages)));
99 entries.push(("accept", Value::from(*accept)));
101 }
102 Frame::Subscription {
103 id,
104 table,
105 scopes,
106 params,
107 cursor,
108 bootstrap_state,
109 } => {
110 entries.push(("type", Value::from("SUBSCRIPTION")));
111 entries.push(("id", Value::from(id.clone())));
112 entries.push(("table", Value::from(table.clone())));
113 entries.push(("scopes", scope_map_value(scopes)));
114 if let Some(p) = params {
115 entries.push(("params", p.parse()));
117 }
118 entries.push(("cursor", Value::from(*cursor)));
119 if let Some(b) = bootstrap_state {
120 entries.push(("bootstrapState", b.parse()));
121 }
122 }
123 Frame::RespHeader {
124 required_schema_version,
125 latest_schema_version,
126 } => {
127 entries.push(("type", Value::from("RESP_HEADER")));
128 if let Some(v) = required_schema_version {
129 entries.push(("requiredSchemaVersion", Value::from(*v)));
130 }
131 if let Some(v) = latest_schema_version {
132 entries.push(("latestSchemaVersion", Value::from(*v)));
133 }
134 }
135 Frame::Lease {
136 lease_id,
137 expires_at_ms,
138 } => {
139 entries.push(("type", Value::from("LEASE")));
140 entries.push(("leaseId", Value::from(lease_id.clone())));
141 entries.push(("expiresAtMs", Value::from(*expires_at_ms)));
142 }
143 Frame::PushResult {
144 client_commit_id,
145 status,
146 commit_seq,
147 results,
148 } => {
149 entries.push(("type", Value::from("PUSH_RESULT")));
150 entries.push(("clientCommitId", Value::from(client_commit_id.clone())));
151 entries.push(("status", Value::from(status.name())));
152 if let Some(v) = commit_seq {
153 entries.push(("commitSeq", Value::from(*v)));
154 }
155 entries.push((
156 "results",
157 Value::Array(results.iter().map(render_result).collect()),
158 ));
159 }
160 Frame::PushResultDetails {
161 client_commit_id,
162 entries: details_entries,
163 } => {
164 entries.push(("type", Value::from("PUSH_RESULT_DETAILS")));
165 entries.push(("clientCommitId", Value::from(client_commit_id.clone())));
166 entries.push((
167 "entries",
168 Value::Array(
169 details_entries
170 .iter()
171 .map(|entry| {
172 obj(vec![
173 ("opIndex", Value::from(entry.op_index)),
174 ("details", entry.details.parse()),
175 ])
176 })
177 .collect(),
178 ),
179 ));
180 }
181 Frame::SubStart {
182 id,
183 status,
184 reason_code,
185 effective_scopes,
186 bootstrap,
187 } => {
188 entries.push(("type", Value::from("SUB_START")));
189 entries.push(("id", Value::from(id.clone())));
190 entries.push(("status", Value::from(status.name())));
191 entries.push(("reasonCode", Value::from(reason_code.clone())));
192 entries.push(("effectiveScopes", scope_map_value(effective_scopes)));
193 entries.push(("bootstrap", Value::from(*bootstrap)));
194 }
195 Frame::Commit {
196 commit_seq,
197 created_at_ms,
198 actor_id,
199 tables,
200 changes,
201 } => {
202 entries.push(("type", Value::from("COMMIT")));
203 entries.push(("commitSeq", Value::from(*commit_seq)));
204 entries.push(("createdAtMs", Value::from(*created_at_ms)));
205 entries.push(("actorId", Value::from(actor_id.clone())));
206 entries.push((
207 "tables",
208 Value::Array(tables.iter().map(|t| Value::from(t.clone())).collect()),
209 ));
210 let rendered_changes = changes
211 .iter()
212 .map(|c| {
213 let mut e: Vec<(&str, Value)> = vec![
214 ("tableIndex", Value::from(c.table_index)),
215 ("rowId", Value::from(c.row_id.clone())),
216 ("op", Value::from(c.op.name())),
217 ];
218 if let Some(v) = c.row_version {
219 e.push(("rowVersion", Value::from(v)));
220 }
221 let mut scopes = Map::new();
222 for (k, v) in &c.scopes {
223 scopes.insert(k.clone(), Value::from(v.clone()));
224 }
225 e.push(("scopes", Value::Object(scopes)));
226 if let Some(row) = &c.row {
227 e.push(("row", Value::from(base64(row))));
229 }
230 obj(e)
231 })
232 .collect();
233 entries.push(("changes", Value::Array(rendered_changes)));
234 }
235 Frame::SegmentRef {
236 segment_id,
237 media_type,
238 table,
239 byte_length,
240 row_count,
241 as_of_commit_seq,
242 scope_digest,
243 row_cursor,
244 next_row_cursor,
245 url,
246 url_expires_at_ms,
247 } => {
248 entries.push(("type", Value::from("SEGMENT_REF")));
249 entries.push(("segmentId", Value::from(segment_id.clone())));
250 entries.push(("mediaType", Value::from(media_type.name())));
251 entries.push(("table", Value::from(table.clone())));
252 entries.push(("byteLength", Value::from(*byte_length)));
253 entries.push(("rowCount", Value::from(*row_count)));
254 entries.push(("asOfCommitSeq", Value::from(*as_of_commit_seq)));
255 entries.push(("scopeDigest", Value::from(scope_digest.clone())));
256 if let Some(v) = row_cursor {
257 entries.push(("rowCursor", Value::from(v.clone())));
258 }
259 if let Some(v) = next_row_cursor {
260 entries.push(("nextRowCursor", Value::from(v.clone())));
261 }
262 if let Some(v) = url {
263 entries.push(("url", Value::from(v.clone())));
264 }
265 if let Some(v) = url_expires_at_ms {
266 entries.push(("urlExpiresAtMs", Value::from(*v)));
267 }
268 }
269 Frame::SegmentInline { payload } => {
270 entries.push(("type", Value::from("SEGMENT_INLINE")));
273 entries.push(("payload", Value::from(base64(payload))));
274 }
275 Frame::SubEnd {
276 next_cursor,
277 bootstrap_state,
278 } => {
279 entries.push(("type", Value::from("SUB_END")));
280 entries.push(("nextCursor", Value::from(*next_cursor)));
281 if let Some(b) = bootstrap_state {
282 entries.push(("bootstrapState", b.parse()));
283 }
284 }
285 Frame::Error {
286 code,
287 message,
288 category,
289 retryable,
290 recommended_action,
291 details,
292 } => {
293 entries.push(("type", Value::from("ERROR")));
294 entries.push(("code", Value::from(code.clone())));
295 entries.push(("message", Value::from(message.clone())));
296 entries.push(("category", Value::from(category.clone())));
297 entries.push(("retryable", Value::from(*retryable)));
298 entries.push(("recommendedAction", Value::from(recommended_action.clone())));
299 if let Some(d) = details {
300 entries.push(("details", d.parse()));
301 }
302 }
303 Frame::Unknown {
304 frame_type,
305 payload,
306 } => {
307 entries.push(("type", Value::from("UNKNOWN")));
309 entries.push(("frameType", Value::from(*frame_type)));
310 entries.push(("payload", Value::from(base64(payload))));
311 }
312 }
313 obj(entries)
314}
315
316fn render_result(result: &OpResult) -> Value {
317 match result {
318 OpResult::Applied { op_index } => obj(vec![
319 ("opIndex", Value::from(*op_index)),
320 ("status", Value::from("applied")),
321 ]),
322 OpResult::Conflict {
323 op_index,
324 code,
325 message,
326 server_version,
327 server_row,
328 } => obj(vec![
329 ("opIndex", Value::from(*op_index)),
330 ("status", Value::from("conflict")),
331 ("code", Value::from(code.clone())),
332 ("message", Value::from(message.clone())),
333 ("serverVersion", Value::from(*server_version)),
334 ("serverRow", Value::from(base64(server_row))),
335 ]),
336 OpResult::Error {
337 op_index,
338 code,
339 message,
340 retryable,
341 } => obj(vec![
342 ("opIndex", Value::from(*op_index)),
343 ("status", Value::from("error")),
344 ("code", Value::from(code.clone())),
345 ("message", Value::from(message.clone())),
346 ("retryable", Value::from(*retryable)),
347 ]),
348 }
349}
350
351pub fn render_rows_segment(seg: &RowsSegment) -> Value {
355 let columns = seg
356 .columns
357 .iter()
358 .map(|c| {
359 obj(vec![
360 ("name", Value::from(c.name.clone())),
361 ("type", Value::from(c.ty.name())),
362 ("nullable", Value::from(c.nullable)),
363 ])
364 })
365 .collect();
366 let blocks = seg
367 .blocks
368 .iter()
369 .map(|block| {
370 Value::Array(
371 block
372 .iter()
373 .map(|row| {
374 let mut values = Map::new();
375 for (col, value) in seg.columns.iter().zip(row.values.iter()) {
376 values.insert(col.name.clone(), render_column_value(value));
377 }
378 obj(vec![
379 ("serverVersion", Value::from(row.server_version)),
380 ("values", Value::Object(values)),
381 ])
382 })
383 .collect(),
384 )
385 })
386 .collect();
387 obj(vec![
388 ("magic", Value::from("SSG2")),
389 ("formatVersion", Value::from(1)),
390 ("table", Value::from(seg.table.clone())),
391 ("schemaVersion", Value::from(seg.schema_version)),
392 ("columns", Value::Array(columns)),
393 ("blocks", Value::Array(blocks)),
394 ])
395}
396
397fn render_column_value(value: &Option<ColumnValue>) -> Value {
398 match value {
399 None => Value::Null,
400 Some(ColumnValue::String(s)) => Value::from(s.clone()),
401 Some(ColumnValue::Integer(v)) => Value::from(*v),
402 Some(ColumnValue::Float(v)) => {
403 serde_json::Number::from_f64(*v).map_or(Value::Null, Value::Number)
404 }
405 Some(ColumnValue::Boolean(v)) => Value::from(*v),
406 Some(ColumnValue::Json(j)) => j.parse(),
407 Some(ColumnValue::Bytes(b)) => Value::from(base64(b)),
408 Some(ColumnValue::BlobRef(j)) => j.parse(),
410 Some(ColumnValue::Crdt(b)) => Value::from(base64(b)),
412 }
413}