Skip to main content

tabnas_alchemy/shared/
sink.rs

1//! The push boundary between stages.
2//!
3//! A pipeline is a chain of sinks. The source calls the first sink once per
4//! event, synchronously, on the thread that parses; each stage does its work
5//! and calls the next. Nothing is queued between stages, so a slow writer at
6//! the end slows the parser at the start: that is the backpressure, and it
7//! costs no buffer. A stage that must stop early (a `take`) answers
8//! [`Flow::Stop`], which the source turns into a cancelled parse.
9
10use indexmap::IndexSet;
11
12use crate::shared::error::{Code, Fail};
13use crate::shared::event::{JsonEvent, OwnedJsonEvent};
14use crate::shared::selector::{Path, Segment};
15
16/// What a stage wants next.
17#[derive(Clone, Copy, Debug, PartialEq, Eq)]
18pub enum Flow {
19    /// Keep sending.
20    Continue,
21    /// The stage has all it needs; the source should stop. Not an error:
22    /// the source stops the parse, releases what it holds, and reports
23    /// nothing further. Whether the rest of the input is validated first
24    /// is the source's documented policy.
25    Stop,
26}
27
28/// A consumer of `JsonEvents/1`.
29pub trait Sink {
30    /// One event. An `Err` aborts the run; the source stops the parse and
31    /// the error reaches the caller unchanged.
32    fn event(&mut self, ev: JsonEvent<'_>) -> Result<Flow, Fail>;
33}
34
35impl Sink for Vec<OwnedJsonEvent> {
36    fn event(&mut self, ev: JsonEvent<'_>) -> Result<Flow, Fail> {
37        self.push(ev.to_owned());
38        Ok(Flow::Continue)
39    }
40}
41
42impl<S: Sink + ?Sized> Sink for &mut S {
43    fn event(&mut self, ev: JsonEvent<'_>) -> Result<Flow, Fail> {
44        (**self).event(ev)
45    }
46}
47
48impl<S: Sink + ?Sized> Sink for Box<S> {
49    fn event(&mut self, ev: JsonEvent<'_>) -> Result<Flow, Fail> {
50        (**self).event(ev)
51    }
52}
53
54/// A sink made of a closure.
55pub struct FnSink<F>(pub F);
56
57impl<F> Sink for FnSink<F>
58where
59    F: FnMut(JsonEvent<'_>) -> Result<Flow, Fail>,
60{
61    fn event(&mut self, ev: JsonEvent<'_>) -> Result<Flow, Fail> {
62        (self.0)(ev)
63    }
64}
65
66/// A sink that counts events and drops them: the cheapest consumer, for
67/// measuring a source on its own.
68#[derive(Debug, Default)]
69pub struct CountSink {
70    pub events: u64,
71}
72
73impl Sink for CountSink {
74    fn event(&mut self, _ev: JsonEvent<'_>) -> Result<Flow, Fail> {
75        self.events += 1;
76        Ok(Flow::Continue)
77    }
78}
79
80/// A tree's events, held to their contract in front of a sink that takes
81/// them as one (a render that writes a document from them): one root
82/// value, and in each object a key and then its value, each key once. A
83/// value walked from a parsed tree keeps it by construction. A parse
84/// streamed as it proceeds may not: it hands on a member its grammar reads
85/// twice (JSON's `{"a":1,"a":2}`, whose value keeps the last) where a tree
86/// has one, and the rule-event adapter refuses most shapes it cannot
87/// follow but not every one a grammar can produce. A repeated key in one
88/// object is refused with `DUPLICATE_MEMBER`, and events no tree has (a
89/// value where a key is due, a key outside an object, a close with nothing
90/// open, a second root) with `STREAMABILITY_UNKNOWN`, each at the path of
91/// the object concerned; the event is not passed on, and what to do then
92/// is the host's (aless falls back once to the parsed value when nothing
93/// has been written). Each open object keeps the keys it has had, dropped
94/// when it closes, so the cost is a lookup per key and the keys of the
95/// objects open at once, which a parse holds in its tree already.
96pub struct TreeContract<S> {
97    open: Vec<Open>,
98    inner: S,
99    /// A whole root value has passed.
100    root_done: bool,
101}
102
103/// A container [`TreeContract`] has open.
104enum Open {
105    /// The keys the object has had, in order, the last of them the member
106    /// whose value is due unless `key_due`; and whether a key is due next
107    /// rather than a value.
108    Object {
109        keys: IndexSet<Box<str>>,
110        key_due: bool,
111    },
112    /// The index of the element due next.
113    Array { next: usize },
114}
115
116impl<S> TreeContract<S> {
117    pub fn new(inner: S) -> TreeContract<S> {
118        TreeContract {
119            open: Vec::new(),
120            inner,
121            root_done: false,
122        }
123    }
124
125    pub fn inner(&self) -> &S {
126        &self.inner
127    }
128
129    pub fn into_inner(self) -> S {
130        self.inner
131    }
132
133    /// The path of the value due next: the open containers, each by the
134    /// member or element open in it, and in the innermost object, its
135    /// last key when that member's value is due.
136    fn path(&self) -> Path {
137        let mut segments = Vec::with_capacity(self.open.len());
138        let last = self.open.len().wrapping_sub(1);
139        for (i, open) in self.open.iter().enumerate() {
140            match open {
141                Open::Object { keys, key_due } => {
142                    // A container open inside this object is its last
143                    // key's value, whatever `key_due` says: the member was
144                    // counted as taken when its value opened.
145                    if i != last || !*key_due {
146                        if let Some(key) = keys.last() {
147                            segments.push(Segment::Key(key.clone()));
148                        }
149                    }
150                }
151                Open::Array { next } => segments.push(Segment::Index(*next)),
152            }
153        }
154        Path(segments)
155    }
156
157    /// The failure for events no tree has.
158    fn not_a_tree(&self, what: &str) -> Fail {
159        Fail::new(
160            Code::StreamabilityUnknown,
161            format!(
162                "the stream holds {what}, which a tree's events never do, so it is not a tree's"
163            ),
164        )
165        .at_path(self.path().to_string())
166    }
167}
168
169impl<S: Sink> Sink for TreeContract<S> {
170    fn event(&mut self, ev: JsonEvent<'_>) -> Result<Flow, Fail> {
171        match &ev {
172            JsonEvent::Key(key) => match self.open.last_mut() {
173                Some(Open::Object { keys, key_due }) if *key_due => {
174                    if !keys.insert((*key).into()) {
175                        let mut path = self.path();
176                        path.0.push(Segment::Key((*key).into()));
177                        return Err(Fail::new(
178                            Code::DuplicateMember,
179                            format!(
180                                "member {key:?} appears twice in one object, and a tree's events \
181                                 hold each key once"
182                            ),
183                        )
184                        .at_path(path.to_string()));
185                    }
186                    *key_due = false;
187                }
188                Some(Open::Object { .. }) => {
189                    return Err(self.not_a_tree("a key where a value is due"));
190                }
191                _ => return Err(self.not_a_tree("a key outside an object")),
192            },
193            JsonEvent::ObjectEnd => match self.open.last() {
194                Some(Open::Object { key_due: true, .. }) => {
195                    self.open.pop();
196                    self.closed();
197                }
198                Some(Open::Object { .. }) => {
199                    return Err(self.not_a_tree("an object's end where a value is due"));
200                }
201                _ => return Err(self.not_a_tree("an object's end where none is due")),
202            },
203            JsonEvent::ArrayEnd => match self.open.last() {
204                Some(Open::Array { .. }) => {
205                    self.open.pop();
206                    self.closed();
207                }
208                _ => return Err(self.not_a_tree("an array's end where none is due")),
209            },
210            JsonEvent::End => {
211                if !self.open.is_empty() {
212                    return Err(self.not_a_tree("its end inside an open container"));
213                }
214            }
215            // A value: in an object, only once its key is in; at the root,
216            // only once.
217            value => {
218                match self.open.last_mut() {
219                    Some(Open::Object { key_due, .. }) => {
220                        if *key_due {
221                            return Err(self.not_a_tree("a value where a key is due"));
222                        }
223                        *key_due = true;
224                    }
225                    Some(Open::Array { .. }) => {}
226                    None => {
227                        if self.root_done {
228                            return Err(self.not_a_tree("a second root value"));
229                        }
230                    }
231                }
232                match value {
233                    JsonEvent::ObjectStart => self.open.push(Open::Object {
234                        keys: IndexSet::new(),
235                        key_due: true,
236                    }),
237                    JsonEvent::ArrayStart => self.open.push(Open::Array { next: 0 }),
238                    _ => self.closed(),
239                }
240            }
241        }
242        self.inner.event(ev)
243    }
244}
245
246impl<S> TreeContract<S> {
247    /// A value is complete: the next in its array is due, or the root is.
248    fn closed(&mut self) {
249        match self.open.last_mut() {
250            Some(Open::Array { next }) => *next += 1,
251            Some(Open::Object { .. }) => {}
252            None => self.root_done = true,
253        }
254    }
255}
256
257/// Replay a recording into a sink, stopping where the sink stops.
258pub fn replay(events: &[OwnedJsonEvent], sink: &mut dyn Sink) -> Result<Flow, Fail> {
259    for ev in events {
260        if sink.event(ev.as_event())? == Flow::Stop {
261            return Ok(Flow::Stop);
262        }
263    }
264    Ok(Flow::Continue)
265}
266
267#[cfg(test)]
268mod tests {
269    use super::*;
270
271    #[test]
272    fn a_vector_records_and_replays() {
273        let mut rec: Vec<OwnedJsonEvent> = Vec::new();
274        rec.event(JsonEvent::ArrayStart).unwrap();
275        rec.event(JsonEvent::Bool(true)).unwrap();
276        rec.event(JsonEvent::ArrayEnd).unwrap();
277        rec.event(JsonEvent::End).unwrap();
278        let mut count = CountSink::default();
279        assert_eq!(replay(&rec, &mut count).unwrap(), Flow::Continue);
280        assert_eq!(count.events, 4);
281    }
282
283    #[test]
284    fn stop_ends_a_replay() {
285        let rec = vec![
286            OwnedJsonEvent::ArrayStart,
287            OwnedJsonEvent::Null,
288            OwnedJsonEvent::ArrayEnd,
289            OwnedJsonEvent::End,
290        ];
291        let mut seen = 0;
292        let mut stopper = FnSink(|_ev: JsonEvent<'_>| {
293            seen += 1;
294            Ok(if seen == 2 {
295                Flow::Stop
296            } else {
297                Flow::Continue
298            })
299        });
300        assert_eq!(replay(&rec, &mut stopper).unwrap(), Flow::Stop);
301        assert_eq!(seen, 2);
302    }
303
304    fn events(json: &str) -> Vec<OwnedJsonEvent> {
305        let value: serde_json::Value = serde_json::from_str(json).unwrap();
306        let mut rec = Vec::new();
307        fn walk(v: &serde_json::Value, rec: &mut Vec<OwnedJsonEvent>) {
308            use OwnedJsonEvent as E;
309            match v {
310                serde_json::Value::Null => rec.push(E::Null),
311                serde_json::Value::Bool(b) => rec.push(E::Bool(*b)),
312                serde_json::Value::Number(n) => rec.push(E::Number {
313                    value: n.as_f64().unwrap(),
314                    lexeme: None,
315                }),
316                serde_json::Value::String(s) => rec.push(E::String(s.as_str().into())),
317                serde_json::Value::Array(items) => {
318                    rec.push(E::ArrayStart);
319                    items.iter().for_each(|i| walk(i, rec));
320                    rec.push(E::ArrayEnd);
321                }
322                serde_json::Value::Object(members) => {
323                    rec.push(E::ObjectStart);
324                    for (k, v) in members {
325                        rec.push(E::Key(k.as_str().into()));
326                        walk(v, rec);
327                    }
328                    rec.push(E::ObjectEnd);
329                }
330            }
331        }
332        walk(&value, &mut rec);
333        rec.push(OwnedJsonEvent::End);
334        rec
335    }
336
337    /// A tree's events pass the contract unchanged, the same key in two
338    /// objects included.
339    #[test]
340    fn a_trees_events_pass_the_tree_contract_unchanged() {
341        let doc = r#"{"a":{"b":1,"a":2},"b":[{"a":3},{"a":4},[],{}],"c":[],"d":{}}"#;
342        let mut guard = TreeContract::new(Vec::new());
343        assert_eq!(replay(&events(doc), &mut guard).unwrap(), Flow::Continue);
344        assert_eq!(guard.inner(), &events(doc));
345        for root in ["1", "\"x\"", "null", "[]", "[1,[2,[3]]]"] {
346            let mut guard = TreeContract::new(Vec::new());
347            replay(&events(root), &mut guard).unwrap();
348            assert_eq!(guard.into_inner(), events(root), "{root}");
349        }
350    }
351
352    /// A key repeated in one object is `DUPLICATE_MEMBER` at the member's
353    /// path, and the event that repeats it is not passed on.
354    #[test]
355    fn a_repeated_key_in_one_object_is_a_duplicate_member_at_its_path() {
356        use OwnedJsonEvent as E;
357        let key = |k: &str| E::Key(k.into());
358        let stream = vec![
359            E::ObjectStart,
360            key("a"),
361            E::Null,
362            key("b"),
363            E::ArrayStart,
364            E::Null,
365            E::ObjectStart,
366            key("x"),
367            E::Null,
368            key("x"),
369        ];
370        let mut guard = TreeContract::new(Vec::new());
371        let fail = replay(&stream, &mut guard).unwrap_err();
372        assert_eq!(fail.code, Code::DuplicateMember, "{fail}");
373        assert!(fail.message.starts_with("member \"x\""), "{fail}");
374        assert_eq!(fail.path.as_deref(), Some(".b[1].x"));
375        assert_eq!(guard.inner().len(), stream.len() - 1);
376        let mut guard = TreeContract::new(Vec::new());
377        let fail = replay(
378            &[E::ObjectStart, key("a b"), E::Null, key("a b")],
379            &mut guard,
380        )
381        .unwrap_err();
382        assert_eq!(fail.path.as_deref(), Some(r#"."a b""#));
383    }
384
385    /// Events no tree has are `STREAMABILITY_UNKNOWN`, naming what was
386    /// met and the path of the value due.
387    #[test]
388    fn events_no_tree_has_are_refused_as_not_a_tree() {
389        use OwnedJsonEvent as E;
390        let key = |k: &str| E::Key(k.into());
391        for (bad, why, path) in [
392            (
393                vec![E::ObjectStart, E::Null],
394                "a value where a key is due",
395                ".",
396            ),
397            (
398                vec![E::ObjectStart, key("a"), E::ObjectStart, E::ArrayStart],
399                "a value where a key is due",
400                ".a",
401            ),
402            (
403                vec![E::ObjectStart, key("a"), key("b")],
404                "a key where a value is due",
405                ".a",
406            ),
407            (
408                vec![E::ArrayStart, key("a")],
409                "a key outside an object",
410                "[0]",
411            ),
412            (vec![key("a")], "a key outside an object", "."),
413            (
414                vec![E::ObjectStart, key("a"), E::ObjectEnd],
415                "an object's end where a value is due",
416                ".a",
417            ),
418            (
419                vec![E::ArrayStart, E::Null, E::ObjectEnd],
420                "an object's end where none is due",
421                "[1]",
422            ),
423            (vec![E::ObjectEnd], "an object's end where none is due", "."),
424            (
425                vec![E::ObjectStart, E::ArrayEnd],
426                "an array's end where none is due",
427                ".",
428            ),
429            (vec![E::Null, E::Null], "a second root value", "."),
430            (
431                vec![E::ArrayStart, E::ArrayEnd, E::ObjectStart],
432                "a second root value",
433                ".",
434            ),
435            (
436                vec![E::ArrayStart, E::End],
437                "its end inside an open container",
438                "[0]",
439            ),
440        ] {
441            let mut guard = TreeContract::new(Vec::new());
442            let fail = replay(&bad, &mut guard).unwrap_err();
443            assert_eq!(fail.code, Code::StreamabilityUnknown, "{bad:?}");
444            assert!(fail.message.contains(why), "{bad:?}: {fail}");
445            assert_eq!(fail.path.as_deref(), Some(path), "{bad:?}: {fail}");
446            assert_eq!(guard.inner().len(), bad.len() - 1, "{bad:?}");
447        }
448    }
449}