1use tabnas_alchemy::shared::{
14 Cell, Code, Fail, Flow, JsonEvent, Number, PublicColumn, Sink, TableEvent, TableSink,
15};
16
17#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
19pub enum MissingRecord {
20 #[default]
23 Skip,
24 Null,
26 Error,
28}
29
30#[derive(Clone, Copy, Debug, PartialEq, Eq)]
31enum Phase {
32 BeforeSchema,
33 Rows,
34 Done,
35}
36
37pub struct RecordsToJson<S: Sink> {
52 sink: S,
53 missing: MissingRecord,
54 phase: Phase,
55 labels: Vec<Box<str>>,
56 next_same: Vec<Option<usize>>,
60 rows: u64,
61 forwarded: bool,
62}
63
64fn contributes(cell: &Cell, missing: MissingRecord) -> bool {
67 !(cell.is_missing() && missing == MissingRecord::Skip)
68}
69
70impl<S: Sink> RecordsToJson<S> {
71 pub fn new(sink: S) -> Self {
72 RecordsToJson {
73 sink,
74 missing: MissingRecord::Skip,
75 phase: Phase::BeforeSchema,
76 labels: Vec::new(),
77 next_same: Vec::new(),
78 rows: 0,
79 forwarded: false,
80 }
81 }
82
83 pub fn with_missing(mut self, missing: MissingRecord) -> Self {
84 self.missing = missing;
85 self
86 }
87
88 pub fn rows(&self) -> u64 {
90 self.rows
91 }
92
93 pub fn is_done(&self) -> bool {
95 self.phase == Phase::Done
96 }
97
98 pub fn into_inner(self) -> S {
99 self.sink
100 }
101
102 fn fail(&self, f: Fail) -> Fail {
103 if self.forwarded {
104 f.committed()
105 } else {
106 f
107 }
108 }
109
110 fn send(&mut self, ev: JsonEvent<'_>) -> Result<Flow, Fail> {
111 self.forwarded = true;
112 self.sink.event(ev)
113 }
114
115 fn schema(&mut self, columns: &[PublicColumn]) -> Result<Flow, Fail> {
116 match self.phase {
117 Phase::BeforeSchema => {}
118 Phase::Rows => return Err(self.fail(Fail::protocol("a second schema"))),
119 Phase::Done => return Err(self.fail(Fail::protocol("a schema after the end"))),
120 }
121 self.labels = columns.iter().map(|c| c.label.clone()).collect();
122 self.next_same = (0..self.labels.len())
123 .map(|i| (i + 1..self.labels.len()).find(|&j| self.labels[j] == self.labels[i]))
124 .collect();
125 self.phase = Phase::Rows;
126 self.send(JsonEvent::ArrayStart)
127 }
128
129 fn row(&mut self, cells: &[Cell]) -> Result<Flow, Fail> {
130 match self.phase {
131 Phase::Rows => {}
132 Phase::BeforeSchema => return Err(Fail::protocol("a row before the schema")),
133 Phase::Done => return Err(self.fail(Fail::protocol("a row after the end"))),
134 }
135 if cells.len() != self.labels.len() {
136 return Err(self.fail(Fail::protocol(format!(
137 "row {} has {} cells; the schema has {} columns",
138 self.rows + 1,
139 cells.len(),
140 self.labels.len()
141 ))));
142 }
143 if self.missing == MissingRecord::Error {
144 if let Some(i) = cells.iter().position(Cell::is_missing) {
147 return Err(self.fail(Fail::new(
148 Code::MissingValue,
149 format!(
150 "row {} has no value for column {:?}",
151 self.rows + 1,
152 self.labels[i]
153 ),
154 )));
155 }
156 }
157 let RecordsToJson {
162 sink,
163 missing,
164 labels,
165 next_same,
166 rows,
167 forwarded,
168 ..
169 } = self;
170 let mut send = |ev: JsonEvent<'_>| -> Result<Flow, Fail> {
171 *forwarded = true;
172 sink.event(ev)
173 };
174 macro_rules! send {
175 ($ev:expr) => {
176 if send($ev)? == Flow::Stop {
177 return Ok(Flow::Stop);
178 }
179 };
180 }
181 send!(JsonEvent::ObjectStart);
182 for (i, cell) in cells.iter().enumerate() {
183 if !contributes(cell, *missing) {
184 continue;
185 }
186 let mut later = next_same.get(i).copied().flatten();
190 while let Some(j) = later {
191 if cells.get(j).is_some_and(|c| contributes(c, *missing)) {
192 break;
193 }
194 later = next_same.get(j).copied().flatten();
195 }
196 if later.is_some() {
197 continue;
198 }
199 let value = match cell {
200 Cell::Null => JsonEvent::Null,
201 Cell::Bool(b) => JsonEvent::Bool(*b),
202 Cell::Number { value, lexeme } => JsonEvent::Number(Number {
203 value: *value,
204 lexeme: lexeme.as_deref(),
205 }),
206 Cell::String(s) => JsonEvent::String(s),
207 Cell::Missing => match *missing {
208 MissingRecord::Null => JsonEvent::Null,
209 MissingRecord::Skip | MissingRecord::Error => continue,
212 },
213 };
214 send!(JsonEvent::Key(&labels[i]));
215 send!(value);
216 }
217 send!(JsonEvent::ObjectEnd);
218 *rows += 1;
219 Ok(Flow::Continue)
220 }
221
222 fn end(&mut self) -> Result<Flow, Fail> {
223 match self.phase {
224 Phase::Rows => {}
225 Phase::BeforeSchema => return Err(Fail::protocol("the end before the schema")),
226 Phase::Done => return Err(self.fail(Fail::protocol("a second end"))),
227 }
228 if self.send(JsonEvent::ArrayEnd)? == Flow::Stop {
229 return Ok(Flow::Stop);
230 }
231 let flow = self.send(JsonEvent::End)?;
234 self.phase = Phase::Done;
235 Ok(flow)
236 }
237}
238
239impl<S: Sink> TableSink for RecordsToJson<S> {
240 fn table_event(&mut self, ev: TableEvent<'_>) -> Result<Flow, Fail> {
241 match ev {
242 TableEvent::Schema(columns) => self.schema(columns),
243 TableEvent::Row(cells) => self.row(cells),
244 TableEvent::End => self.end(),
245 }
246 }
247}
248
249#[cfg(test)]
250mod tests {
251 use super::*;
252 use crate::json::{JsonOptions, JsonRenderer};
253 use crate::text::StringOut;
254 use tabnas_alchemy::shared::{FnSink, OwnedJsonEvent};
255 use OwnedJsonEvent::*;
256
257 fn cols(labels: &[&str]) -> Vec<PublicColumn> {
258 labels.iter().map(|l| PublicColumn::new(*l)).collect()
259 }
260
261 fn s(text: &str) -> Cell {
262 Cell::String(text.into())
263 }
264
265 fn key(k: &str) -> OwnedJsonEvent {
266 Key(k.into())
267 }
268
269 fn str(v: &str) -> OwnedJsonEvent {
270 String(v.into())
271 }
272
273 fn run(
274 missing: MissingRecord,
275 labels: &[&str],
276 rows: &[Vec<Cell>],
277 ) -> Result<Vec<OwnedJsonEvent>, Fail> {
278 let mut r = RecordsToJson::new(Vec::new()).with_missing(missing);
279 let columns = cols(labels);
280 r.table_event(TableEvent::Schema(&columns))?;
281 for row in rows {
282 r.table_event(TableEvent::Row(row))?;
283 }
284 r.table_event(TableEvent::End)?;
285 assert!(r.is_done());
286 Ok(r.into_inner())
287 }
288
289 #[test]
290 fn a_table_becomes_an_array_of_objects_keyed_by_label() {
291 let rows = vec![
292 vec![
293 Cell::Number {
294 value: 1.0,
295 lexeme: Some("1.0".into()),
296 },
297 s("ada"),
298 Cell::Bool(true),
299 ],
300 vec![
301 Cell::Number {
302 value: 2.0,
303 lexeme: None,
304 },
305 Cell::Null,
306 Cell::Bool(false),
307 ],
308 ];
309 assert_eq!(
310 run(MissingRecord::Skip, &["id", "name", "ok"], &rows).unwrap(),
311 vec![
312 ArrayStart,
313 ObjectStart,
314 key("id"),
315 OwnedJsonEvent::Number {
316 value: 1.0,
317 lexeme: Some("1.0".into())
318 },
319 key("name"),
320 str("ada"),
321 key("ok"),
322 Bool(true),
323 ObjectEnd,
324 ObjectStart,
325 key("id"),
326 OwnedJsonEvent::Number {
327 value: 2.0,
328 lexeme: None
329 },
330 key("name"),
331 Null,
332 key("ok"),
333 Bool(false),
334 ObjectEnd,
335 ArrayEnd,
336 End,
337 ]
338 );
339 }
340
341 #[test]
342 fn no_rows_is_an_empty_array() {
343 assert_eq!(
344 run(MissingRecord::Skip, &["a"], &[]).unwrap(),
345 vec![ArrayStart, ArrayEnd, End]
346 );
347 }
348
349 #[test]
350 fn zero_columns_gives_empty_objects() {
351 assert_eq!(
352 run(MissingRecord::Skip, &[], &[vec![], vec![]]).unwrap(),
353 vec![
354 ArrayStart,
355 ObjectStart,
356 ObjectEnd,
357 ObjectStart,
358 ObjectEnd,
359 ArrayEnd,
360 End
361 ]
362 );
363 }
364
365 #[test]
366 fn a_missing_cell_is_skipped_by_default() {
367 assert_eq!(
368 run(
369 MissingRecord::Skip,
370 &["a", "b"],
371 &[vec![Cell::Missing, s("x")]]
372 )
373 .unwrap(),
374 vec![
375 ArrayStart,
376 ObjectStart,
377 key("b"),
378 str("x"),
379 ObjectEnd,
380 ArrayEnd,
381 End
382 ]
383 );
384 }
385
386 #[test]
387 fn a_missing_cell_can_be_null_or_an_error() {
388 assert_eq!(
389 run(
390 MissingRecord::Null,
391 &["a", "b"],
392 &[vec![Cell::Missing, s("x")]]
393 )
394 .unwrap(),
395 vec![
396 ArrayStart,
397 ObjectStart,
398 key("a"),
399 Null,
400 key("b"),
401 str("x"),
402 ObjectEnd,
403 ArrayEnd,
404 End
405 ]
406 );
407 let mut r = RecordsToJson::new(Vec::new()).with_missing(MissingRecord::Error);
408 let columns = cols(&["a", "b"]);
409 r.table_event(TableEvent::Schema(&columns)).unwrap();
410 let err = r
411 .table_event(TableEvent::Row(&[s("x"), Cell::Missing]))
412 .unwrap_err();
413 assert_eq!(err.code, Code::MissingValue);
414 assert!(err.message.contains("\"b\""));
415 assert!(err.committed_output, "the array start was forwarded");
416 assert_eq!(r.into_inner(), vec![ArrayStart], "nothing of the row was");
417 }
418
419 #[test]
420 fn a_repeated_label_keeps_the_last_value_that_is_present() {
421 let rows = vec![
425 vec![s("1"), s("2"), s("3")],
426 vec![s("4"), s("5"), Cell::Missing],
427 vec![Cell::Missing, s("6"), Cell::Missing],
428 ];
429 assert_eq!(
430 run(MissingRecord::Skip, &["a", "b", "a"], &rows).unwrap(),
431 vec![
432 ArrayStart,
433 ObjectStart,
434 key("b"),
435 str("2"),
436 key("a"),
437 str("3"),
438 ObjectEnd,
439 ObjectStart,
440 key("a"),
441 str("4"),
442 key("b"),
443 str("5"),
444 ObjectEnd,
445 ObjectStart,
446 key("b"),
447 str("6"),
448 ObjectEnd,
449 ArrayEnd,
450 End
451 ]
452 );
453 assert_eq!(
454 run(MissingRecord::Null, &["a", "b", "a"], &rows[1..2]).unwrap(),
455 vec![
456 ArrayStart,
457 ObjectStart,
458 key("b"),
459 str("5"),
460 key("a"),
461 Null,
462 ObjectEnd,
463 ArrayEnd,
464 End
465 ]
466 );
467 assert_eq!(
470 run(
471 MissingRecord::Skip,
472 &["a", "a", "a"],
473 &[vec![s("1"), s("2"), Cell::Missing]]
474 )
475 .unwrap(),
476 vec![
477 ArrayStart,
478 ObjectStart,
479 key("a"),
480 str("2"),
481 ObjectEnd,
482 ArrayEnd,
483 End
484 ]
485 );
486 }
487
488 #[test]
489 fn protocol_errors_match_the_csv_renderers() {
490 let columns = cols(&["a"]);
491
492 let mut r = RecordsToJson::new(Vec::new());
493 let err = r.table_event(TableEvent::Row(&[s("x")])).unwrap_err();
494 assert_eq!(err.code, Code::ProtocolOrderError);
495 assert!(!err.committed_output);
496 let err = r.table_event(TableEvent::End).unwrap_err();
497 assert_eq!(err.code, Code::ProtocolOrderError);
498
499 r.table_event(TableEvent::Schema(&columns)).unwrap();
500 let err = r.table_event(TableEvent::Schema(&columns)).unwrap_err();
501 assert_eq!(err.code, Code::ProtocolOrderError);
502 assert!(err.committed_output);
503 let err = r
504 .table_event(TableEvent::Row(&[s("x"), s("y")]))
505 .unwrap_err();
506 assert_eq!(err.code, Code::ProtocolOrderError);
507 assert!(err.message.contains("row 1 has 2 cells"));
508 let err = r.table_event(TableEvent::Row(&[])).unwrap_err();
509 assert_eq!(err.code, Code::ProtocolOrderError);
510
511 r.table_event(TableEvent::End).unwrap();
512 for ev in [
513 TableEvent::End,
514 TableEvent::Row(&[s("x")]),
515 TableEvent::Schema(&columns),
516 ] {
517 let err = r.table_event(ev).unwrap_err();
518 assert_eq!(err.code, Code::ProtocolOrderError);
519 assert!(err.committed_output);
520 }
521 assert_eq!(r.rows(), 0);
522 assert_eq!(r.into_inner(), vec![ArrayStart, ArrayEnd, End]);
523 }
524
525 #[test]
526 fn a_sink_that_fails_on_end_leaves_the_stage_not_done() {
527 let sink = FnSink(|ev: JsonEvent<'_>| {
528 if let JsonEvent::End = ev {
529 Err(Fail::output("closed"))
530 } else {
531 Ok(Flow::Continue)
532 }
533 });
534 let mut r = RecordsToJson::new(sink);
535 let columns = cols(&["a"]);
536 r.table_event(TableEvent::Schema(&columns)).unwrap();
537 let err = r.table_event(TableEvent::End).unwrap_err();
538 assert_eq!(err.code, Code::OutputFailed);
539 assert!(!r.is_done());
540 }
541
542 #[test]
543 fn keys_borrow_the_schemas_labels_rather_than_copying_them_per_cell() {
544 let mut seen: Vec<usize> = Vec::new();
549 let sink = FnSink(|ev: JsonEvent<'_>| {
550 if let JsonEvent::Key(k) = ev {
551 seen.push(k.as_ptr() as usize);
552 }
553 Ok(Flow::Continue)
554 });
555 let mut r = RecordsToJson::new(sink);
556 let columns = cols(&["first", "second"]);
557 r.table_event(TableEvent::Schema(&columns)).unwrap();
558 let row = [s("x"), Cell::Null];
559 r.table_event(TableEvent::Row(&row)).unwrap();
560 r.table_event(TableEvent::Row(&row)).unwrap();
561 let labels: Vec<usize> = r.labels.iter().map(|l| l.as_ptr() as usize).collect();
562 drop(r);
563 assert_eq!(seen, [labels.clone(), labels].concat());
564 }
565
566 #[test]
567 fn a_stop_from_the_sink_stops_the_row() {
568 let mut seen = 0;
569 let sink = FnSink(|_ev: JsonEvent<'_>| {
570 seen += 1;
571 Ok(if seen == 3 {
572 Flow::Stop
573 } else {
574 Flow::Continue
575 })
576 });
577 let mut r = RecordsToJson::new(sink);
578 let columns = cols(&["a", "b"]);
579 assert_eq!(
580 r.table_event(TableEvent::Schema(&columns)).unwrap(),
581 Flow::Continue
582 );
583 assert_eq!(
584 r.table_event(TableEvent::Row(&[s("x"), s("y")])).unwrap(),
585 Flow::Stop
586 );
587 drop(r);
588 assert_eq!(
589 seen, 3,
590 "ArrayStart, ObjectStart, the first key, then no more"
591 );
592 }
593
594 #[test]
595 fn records_render_as_json_text() {
596 let renderer = JsonRenderer::new(StringOut::new(), JsonOptions::default());
597 let mut r = RecordsToJson::new(renderer);
598 let columns = cols(&["name", "age"]);
599 r.table_event(TableEvent::Schema(&columns)).unwrap();
600 r.table_event(TableEvent::Row(&[
601 s("ada"),
602 Cell::Number {
603 value: 36.0,
604 lexeme: Some("36".into()),
605 },
606 ]))
607 .unwrap();
608 r.table_event(TableEvent::Row(&[s("lin"), Cell::Missing]))
609 .unwrap();
610 r.table_event(TableEvent::End).unwrap();
611 assert_eq!(r.rows(), 2);
612 assert_eq!(
613 r.into_inner().into_inner().as_str(),
614 r#"[{"name":"ada","age":36},{"name":"lin"}]"#
615 );
616 }
617}