Skip to main content

tabnas_render/
records.rs

1//! `TableRows/1` as `JsonEvents/1`: an array of records keyed by label.
2//!
3//! A table has a natural JSON form, one object per row with the column
4//! labels as member names, and producing it as events rather than text
5//! means the JSON renderer, and every other `JsonEvents/1` consumer, gets
6//! it for free. The stage retains the labels and nothing else: each row is
7//! emitted as it arrives and forgotten. Labels are data from the source's
8//! metadata, so a repeated label is not an error here (the CSV renderer
9//! allows it too); it is resolved the way a JSON reader would resolve a
10//! repeated member in the un-deduplicated record, by keeping the last
11//! value that is there, and the output then carries each label once.
12
13use tabnas_transduce::{
14    Cell, Code, Fail, Flow, JsonEvent, Number, PublicColumn, Sink, TableEvent, TableSink,
15};
16
17/// What a [`Cell::Missing`] becomes in a record.
18#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
19pub enum MissingRecord {
20    /// Leave the member out: the record says nothing where the source had
21    /// nothing, which is what an absent path meant.
22    #[default]
23    Skip,
24    /// Write the member with a `null` value.
25    Null,
26    /// Fail the run with `MISSING_VALUE`.
27    Error,
28}
29
30#[derive(Clone, Copy, Debug, PartialEq, Eq)]
31enum Phase {
32    BeforeSchema,
33    Rows,
34    Done,
35}
36
37/// Turns `TableRows/1` into `JsonEvents/1`.
38///
39/// `Schema` opens the array, each `Row` is one object whose members are the
40/// labels in schema order, `End` closes the array and ends the document.
41/// The protocol is validated as the CSV renderer validates it: one schema
42/// first, rows of the schema's width, one end, `PROTOCOL_ORDER_ERROR`
43/// otherwise. A schema with no columns is allowed, since an empty object
44/// is a JSON value; the CSV renderer's refusal is about CSV. When a label
45/// repeats, each record carries it once, from the last column whose cell
46/// contributes a member (under `Skip` a `Missing` cell contributes none),
47/// in that column's position: the value a reader of the un-deduplicated
48/// record would keep, since a member that was never written cannot win.
49/// A failure found after events were forwarded is marked as having
50/// committed output, since the stage downstream may have rendered them.
51pub struct RecordsToJson<S: Sink> {
52    sink: S,
53    missing: MissingRecord,
54    phase: Phase,
55    labels: Vec<Box<str>>,
56    /// Per column, the next later column with the same label, so a row
57    /// can find the column that carries the label's value; `None` for the
58    /// common case of a label that does not repeat.
59    next_same: Vec<Option<usize>>,
60    rows: u64,
61    forwarded: bool,
62}
63
64/// Whether a cell contributes a member to its record under this policy:
65/// every cell but a `Missing` that is skipped.
66fn 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    /// Rows emitted so far.
89    pub fn rows(&self) -> u64 {
90        self.rows
91    }
92
93    /// Whether `End` has been forwarded.
94    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            // Before any of the row is forwarded, so a row is emitted whole
145            // or not at all, as the CSV renderer renders it.
146            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        // The sink copies what it keeps, as the protocol says, so a key is
158        // the label borrowed and the row costs no allocation. The fields
159        // are borrowed apart from one another for that: the sink and the
160        // forwarded flag mutably for the sends, the labels for the keys.
161        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            // A repeated label is written from the last column whose cell
187            // contributes; an earlier column's value is superseded only by
188            // a member that will actually be there.
189            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                    // `Skip` does not contribute and was passed over above;
210                    // `Error` was rejected before the row began.
211                    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        // Done only once `End` has been taken downstream: a sink that
232        // failed on it has not seen the document end.
233        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_transduce::{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        // Row two's last "a" is Missing and skipped, so the reader of the
422        // un-deduplicated record would keep "4"; row three has no "a" at
423        // all. Under `Null` the Missing member is there, and wins.
424        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        // Three columns with one label: the middle one wins when the last
468        // is absent.
469        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        // Two live allocations cannot share an address, so a key that is a
545        // per-cell copy of the label would point elsewhere than the label
546        // the stage retains; the same address on every row proves the
547        // borrow (and, with it, the absence of the allocation).
548        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}