1use std::sync::{Arc, Mutex};
36
37use tabnas::Tabnas;
38
39use crate::error::Code;
40use crate::error::Fail;
41use crate::event::JsonEvent;
42use crate::limits::{AbortFlag, Limits, Metrics};
43use crate::sink::{Flow, Sink};
44use crate::source::guard::Guarded;
45use crate::source::rule_events::{self, Adapter, Status, GUARD};
46use crate::source::{capability, engine_failure, walk_value, Prune, Source, SourceMode};
47
48pub struct ParserSource<'s> {
50 parser: Tabnas,
51 text: &'s str,
52 mode: SourceMode,
53 limits: Limits,
54 abort: AbortFlag,
55 metrics: Arc<Metrics>,
56 grammar: Option<Box<str>>,
57 unverified: bool,
58}
59
60impl<'s> ParserSource<'s> {
61 pub fn new(parser: Tabnas, text: &'s str) -> ParserSource<'s> {
64 ParserSource {
65 parser,
66 text,
67 mode: SourceMode::Materialize,
68 limits: Limits::default(),
69 abort: AbortFlag::new(),
70 metrics: Metrics::new(),
71 grammar: None,
72 unverified: false,
73 }
74 }
75
76 pub fn mode(mut self, mode: SourceMode) -> Self {
77 self.mode = mode;
78 self
79 }
80
81 pub fn grammar(mut self, name: &str) -> Self {
87 self.grammar = Some(name.into());
88 self
89 }
90
91 pub fn unverified(mut self) -> Self {
97 self.unverified = true;
98 self
99 }
100
101 fn gate(&self) -> Option<Fail> {
103 if self.unverified {
104 return None;
105 }
106 match self.grammar.as_deref() {
107 Some(name) if capability::incremental(name) => None,
108 Some(name) => Some(Fail::new(
109 Code::StreamabilityUnknown,
110 format!(
111 "grammar {name:?} is not in capability::incremental: the differential suite \
112 has not verified that its rule events stream as the walk does; run it with \
113 SourceMode::Materialize"
114 ),
115 )),
116 None => Some(Fail::new(
117 Code::StreamabilityUnknown,
118 "SourceMode::Incremental needs the grammar's name (ParserSource::grammar) to \
119 check capability::incremental; without one, run SourceMode::Materialize",
120 )),
121 }
122 }
123
124 pub fn limits(mut self, limits: Limits) -> Self {
125 self.limits = limits;
126 self
127 }
128
129 pub fn abort(mut self, abort: AbortFlag) -> Self {
130 self.abort = abort;
131 self
132 }
133
134 pub fn metrics(mut self, metrics: Arc<Metrics>) -> Self {
135 self.metrics = metrics;
136 self
137 }
138
139 pub fn run_owned<S: Sink + Send + 'static>(self, sink: S) -> (Result<Flow, Fail>, S) {
142 let (outcome, sink, _) = self.run_owned_with_value(sink);
143 (outcome, sink)
144 }
145
146 pub fn run_owned_with_value<S: Sink + Send + 'static>(
153 self,
154 sink: S,
155 ) -> (Result<Flow, Fail>, S, Option<tabnas::Value>) {
156 match &self.mode {
157 SourceMode::Materialize => {
158 let mut guarded =
159 Guarded::new(sink, &self.limits, self.abort.clone(), self.metrics.clone());
160 let (outcome, value) =
161 materialize(self.parser, self.text, &self.abort, &mut guarded);
162 (outcome, guarded.into_inner(), value)
163 }
164 SourceMode::Incremental { prune } => {
165 if let Some(refused) = self.gate() {
166 return (Err(refused), sink, None);
167 }
168 incremental(
169 self.parser,
170 self.text,
171 &self.limits,
172 self.abort,
173 self.metrics,
174 prune,
175 sink,
176 )
177 }
178 }
179 }
180
181 pub fn run_boxed(
183 self,
184 sink: Box<dyn Sink + Send>,
185 ) -> (Result<Flow, Fail>, Box<dyn Sink + Send>) {
186 self.run_owned(sink)
187 }
188}
189
190impl Source for ParserSource<'_> {
191 fn run(self, sink: &mut dyn Sink) -> Result<Flow, Fail> {
196 let mut guarded = Guarded::new(sink, &self.limits, self.abort.clone(), self.metrics);
197 let (outcome, _) = materialize(self.parser, self.text, &self.abort, &mut guarded);
198 guarded.flush();
199 outcome
200 }
201}
202
203fn materialize<S: Sink>(
205 mut parser: Tabnas,
206 text: &str,
207 abort: &AbortFlag,
208 guarded: &mut Guarded<S>,
209) -> (Result<Flow, Fail>, Option<tabnas::Value>) {
210 let flag = abort.clone();
211 parser.parse_guard(GUARD, move |_ctx| !flag.is_aborted());
212 let value = match parser.parse(text) {
213 Ok(value) => value,
214 Err(e) => return (Err(engine_failure(&e, abort)), None),
215 };
216 drop(parser);
217 let outcome = match walk_value(&value, guarded) {
218 Ok(Flow::Continue) => guarded.event(JsonEvent::End),
219 other => other,
220 };
221 (outcome, Some(value))
222}
223
224fn incremental<S: Sink + Send + 'static>(
225 mut parser: Tabnas,
226 text: &str,
227 limits: &Limits,
228 abort: AbortFlag,
229 metrics: Arc<Metrics>,
230 prune: &Prune,
231 sink: S,
232) -> (Result<Flow, Fail>, S, Option<tabnas::Value>) {
233 let stop = AbortFlag::new();
234 let adapter = Adapter::new(sink, limits, abort.clone(), metrics, prune, stop.clone());
235 let shared = Arc::new(Mutex::new(adapter));
236 Adapter::install(
237 &mut parser,
238 Arc::downgrade(&shared),
239 abort.clone(),
240 stop.clone(),
241 );
242 let parsed = parser.parse(text);
243 drop(parser);
244 let mut adapter = rule_events::take(shared);
245 let outcome = match adapter.status() {
246 Status::Failed(_) | Status::Stopped => Ok(Flow::Continue),
247 Status::Running => match &parsed {
248 Ok(_) if adapter.complete() => adapter.send(JsonEvent::End),
249 Ok(value) if adapter.idle() => match adapter.walk_whole(value) {
250 Ok(Flow::Continue) => adapter.send(JsonEvent::End),
251 other => other,
252 },
253 Ok(_) => Err(rule_events::not_streamable()),
254 Err(e) => Err(engine_failure(e, &abort)),
255 },
256 };
257 let (status, sink) = adapter.finish();
258 let outcome = match status {
259 Status::Failed(fail) => Err(fail),
260 Status::Stopped => Ok(Flow::Stop),
261 Status::Running => outcome,
262 };
263 (outcome, sink, parsed.ok())
264}
265
266#[cfg(test)]
267mod tests {
268 use super::*;
269 use crate::error::Code;
270 use crate::event::OwnedJsonEvent;
271 use crate::selector::Selector;
272 use crate::sink::FnSink;
273
274 fn incremental_mode() -> SourceMode {
275 SourceMode::Incremental {
276 prune: Prune::Never,
277 }
278 }
279
280 fn record(mode: SourceMode, src: &str) -> (Result<Flow, Fail>, Vec<OwnedJsonEvent>) {
281 ParserSource::new(tabnas_json::make(), src)
282 .grammar("json")
283 .mode(mode)
284 .run_owned(Vec::new())
285 }
286
287 #[test]
288 fn incremental_mode_needs_a_verified_grammar_name_and_emits_nothing_without_one() {
289 let (r, events) = ParserSource::new(tabnas_json::make(), DOC)
290 .mode(incremental_mode())
291 .run_owned(Vec::<OwnedJsonEvent>::new());
292 let err = r.unwrap_err();
293 assert_eq!(err.code, Code::StreamabilityUnknown);
294 assert!(err.message.contains("ParserSource::grammar"), "{err}");
295 assert!(events.is_empty());
296
297 let (r, events) = ParserSource::new(tabnas_csv::make(), "a,b\n1,2\n")
298 .grammar("csv")
299 .mode(incremental_mode())
300 .run_owned(Vec::<OwnedJsonEvent>::new());
301 let err = r.unwrap_err();
302 assert_eq!(err.code, Code::StreamabilityUnknown);
303 assert!(err.message.contains("\"csv\""), "{err}");
304 assert!(events.is_empty(), "refused before the parse");
305
306 let (r, walked) = ParserSource::new(tabnas_csv::make(), "a,b\n1,2\n")
312 .run_owned(Vec::<OwnedJsonEvent>::new());
313 r.unwrap();
314 assert_eq!(walked.last(), Some(&OwnedJsonEvent::End));
315 let (r, streamed) = ParserSource::new(tabnas_csv::make(), "a,b\n1,2\n")
316 .unverified()
317 .mode(incremental_mode())
318 .run_owned(Vec::<OwnedJsonEvent>::new());
319 let err = r.unwrap_err();
320 assert_eq!(err.code, Code::StreamabilityUnknown);
321 assert!(
322 err.message
323 .contains("the incremental source cannot follow a grammar that builds"),
324 "{err}"
325 );
326 assert!(
327 !streamed.is_empty(),
328 "refused during the parse, not before it"
329 );
330 assert!(!streamed.contains(&OwnedJsonEvent::End));
331 }
332
333 fn without_lexemes(events: &[OwnedJsonEvent]) -> Vec<OwnedJsonEvent> {
334 events
335 .iter()
336 .map(|e| match e {
337 OwnedJsonEvent::Number { value, .. } => OwnedJsonEvent::Number {
338 value: *value,
339 lexeme: None,
340 },
341 other => other.clone(),
342 })
343 .collect()
344 }
345
346 const DOC: &str = r#"{"a":[1,2.50,"x",{"b":null}],"c":{},"d":[],"e":1e21,"f":true}"#;
347
348 #[test]
349 fn incremental_events_equal_the_walk_and_carry_lexemes() {
350 let (r1, inc) = record(incremental_mode(), DOC);
351 let (r2, mat) = record(SourceMode::Materialize, DOC);
352 assert_eq!(r1.unwrap(), Flow::Continue);
353 assert_eq!(r2.unwrap(), Flow::Continue);
354 assert_eq!(without_lexemes(&inc), mat);
355 assert_eq!(inc.last(), Some(&OwnedJsonEvent::End));
356 let lexemes: Vec<Option<&str>> = inc
357 .iter()
358 .filter_map(|e| match e {
359 OwnedJsonEvent::Number { lexeme, .. } => Some(lexeme.as_deref()),
360 _ => None,
361 })
362 .collect();
363 assert_eq!(lexemes, [Some("1"), Some("2.50"), Some("1e21")]);
364 assert!(mat.iter().all(|e| !matches!(
365 e,
366 OwnedJsonEvent::Number {
367 lexeme: Some(_),
368 ..
369 }
370 )));
371 }
372
373 #[test]
374 fn a_root_scalar_is_one_event_then_end() {
375 let (r, inc) = record(incremental_mode(), " 42 ");
376 assert_eq!(r.unwrap(), Flow::Continue);
377 assert_eq!(
378 inc,
379 vec![
380 OwnedJsonEvent::Number {
381 value: 42.0,
382 lexeme: Some("42".into())
383 },
384 OwnedJsonEvent::End
385 ]
386 );
387 let (r, inc) = record(incremental_mode(), r#""s""#);
388 assert_eq!(r.unwrap(), Flow::Continue);
389 assert_eq!(
390 inc,
391 vec![OwnedJsonEvent::String("s".into()), OwnedJsonEvent::End]
392 );
393 }
394
395 #[test]
396 fn the_borrowed_run_materializes_in_either_mode() {
397 for mode in [SourceMode::Materialize, incremental_mode()] {
398 let mut rec: Vec<OwnedJsonEvent> = Vec::new();
399 let r = ParserSource::new(tabnas_json::make(), DOC)
400 .grammar("json")
401 .mode(mode)
402 .run(&mut rec);
403 assert_eq!(r.unwrap(), Flow::Continue);
404 assert_eq!(rec, record(SourceMode::Materialize, DOC).1);
405 }
406 }
407
408 #[test]
409 fn a_stop_from_the_sink_stops_the_parse_and_returns_the_sink() {
410 for mode in [SourceMode::Materialize, incremental_mode()] {
411 let seen = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
412 let counter = seen.clone();
413 let sink = FnSink(move |_ev: JsonEvent<'_>| {
414 let n = counter.fetch_add(1, std::sync::atomic::Ordering::Relaxed) + 1;
415 Ok(if n == 3 { Flow::Stop } else { Flow::Continue })
416 });
417 let (r, _sink) = ParserSource::new(tabnas_json::make(), DOC)
418 .grammar("json")
419 .mode(mode)
420 .run_owned(sink);
421 assert_eq!(r.unwrap(), Flow::Stop);
422 assert_eq!(seen.load(std::sync::atomic::Ordering::Relaxed), 3);
423 }
424 }
425
426 #[test]
427 fn a_sink_failure_comes_back_unchanged() {
428 for mode in [SourceMode::Materialize, incremental_mode()] {
429 let sink = FnSink(|ev: JsonEvent<'_>| {
430 if ev == JsonEvent::Key("c") {
431 Err(Fail::output("disk full").at_path(".c"))
432 } else {
433 Ok(Flow::Continue)
434 }
435 });
436 let (r, _) = ParserSource::new(tabnas_json::make(), DOC)
437 .grammar("json")
438 .mode(mode)
439 .run_owned(sink);
440 let err = r.unwrap_err();
441 assert_eq!(err.code, Code::OutputFailed);
442 assert_eq!(err.path.as_deref(), Some(".c"));
443 }
444 }
445
446 #[test]
447 fn an_aborted_flag_cancels_the_parse_as_aborted() {
448 for mode in [SourceMode::Materialize, incremental_mode()] {
449 let abort = AbortFlag::new();
450 abort.abort();
451 let (r, _) = ParserSource::new(tabnas_json::make(), DOC)
452 .grammar("json")
453 .mode(mode)
454 .abort(abort)
455 .run_owned(Vec::<OwnedJsonEvent>::new());
456 assert_eq!(r.unwrap_err().code, Code::Aborted);
457 }
458 }
459
460 #[test]
461 fn a_parse_error_is_invalid_input_with_its_position() {
462 for mode in [SourceMode::Materialize, incremental_mode()] {
463 let (r, _) = record(mode, "{\"a\": 1,\n \"b\": }");
464 let err = r.unwrap_err();
465 assert_eq!(err.code, Code::InputInvalid);
466 assert!(err.message.starts_with("unexpected"), "{}", err.message);
467 assert_eq!(err.row, Some(2));
468 assert_eq!(err.column, Some(7));
469 }
470 }
471
472 #[test]
476 fn a_grammars_own_guard_is_invalid_input_that_names_the_grammar() {
477 let src = format!("{}1{}", "[".repeat(200), "]".repeat(200));
478 for mode in [SourceMode::Materialize, incremental_mode()] {
479 let (r, _) = record(mode, &src);
480 let err = r.unwrap_err();
481 assert_eq!(err.code, Code::InputInvalid);
482 assert!(
483 err.message.starts_with("the grammar stopped the parse"),
484 "{err}"
485 );
486 assert!(err.message.contains("cancel"), "{err}");
487 assert_eq!(err.column, Some(128));
488 }
489 }
490
491 #[test]
492 fn source_limits_apply_in_both_modes_by_name() {
493 for mode in [SourceMode::Materialize, incremental_mode()] {
494 let limits = Limits {
495 max_key_bytes: 1,
496 ..Limits::default()
497 };
498 let (r, _) = ParserSource::new(tabnas_json::make(), r#"{"ab":1}"#)
499 .grammar("json")
500 .mode(mode.clone())
501 .limits(limits)
502 .run_owned(Vec::<OwnedJsonEvent>::new());
503 let err = r.unwrap_err();
504 assert_eq!(err.code, Code::ResourceLimitExceeded);
505 assert_eq!(err.limit.as_ref().unwrap().name, "max_key_bytes");
506
507 let limits = Limits {
508 max_depth: 2,
509 ..Limits::default()
510 };
511 let (r, _) = ParserSource::new(tabnas_json::make(), "[[[1]]]")
512 .grammar("json")
513 .mode(mode.clone())
514 .limits(limits)
515 .run_owned(Vec::<OwnedJsonEvent>::new());
516 assert_eq!(r.unwrap_err().limit.unwrap().name, "max_depth");
517
518 let limits = Limits {
519 max_scalar_bytes: 2,
520 ..Limits::default()
521 };
522 let (r, _) = ParserSource::new(tabnas_json::make(), r#"["abc"]"#)
523 .grammar("json")
524 .mode(mode)
525 .limits(limits)
526 .run_owned(Vec::<OwnedJsonEvent>::new());
527 assert_eq!(r.unwrap_err().limit.unwrap().name, "max_scalar_bytes");
528 }
529 }
530
531 #[test]
532 fn metrics_count_the_source_events() {
533 for mode in [SourceMode::Materialize, incremental_mode()] {
534 let metrics = Metrics::new();
535 let (r, events) = ParserSource::new(tabnas_json::make(), DOC)
536 .grammar("json")
537 .mode(mode)
538 .metrics(metrics.clone())
539 .run_owned(Vec::<OwnedJsonEvent>::new());
540 r.unwrap();
541 assert_eq!(Metrics::get(&metrics.events), events.len() as u64);
542 assert_eq!(Metrics::get(&metrics.keys), 6);
543 assert_eq!(Metrics::get(&metrics.scalars), 6);
544 }
545 }
546
547 #[test]
548 fn pruning_leaves_the_events_untouched() {
549 let (_, plain) = record(incremental_mode(), DOC);
550 for prune in [
551 Prune::AllArrays,
552 Prune::Under(Selector::root().property("a").each_index()),
553 Prune::Under(Selector::root().property("a")),
554 ] {
555 let (r, pruned) = record(SourceMode::Incremental { prune }, DOC);
556 r.unwrap();
557 assert_eq!(pruned, plain);
558 }
559 }
560}