From 584f3538a74f72827dc83ee20ff4f706770f2ad2 Mon Sep 17 00:00:00 2001 From: Eric Rodrigues Pires Date: Fri, 28 Nov 2025 18:19:17 -0300 Subject: [PATCH] Add processors, formatter, and assemble pipeline in main --- Cargo.lock | 244 ++++++++++++++++++++++++++++++ duper/src/ast.rs | 40 +++++ duper/src/parser/mod.rs | 3 +- duperq/Cargo.toml | 4 + duperq/src/accessor.rs | 23 +++ duperq/src/filter.rs | 66 +++++---- duperq/src/formatter.rs | 238 +++++++++++++++++++++++++++++ duperq/src/lib.rs | 4 + duperq/src/main.rs | 104 ++++++++++++- duperq/src/processor.rs | 94 ++++++++++++ duperq/src/query.rs | 321 ++++++++++++++++++++++++++-------------- 11 files changed, 995 insertions(+), 146 deletions(-) create mode 100644 duperq/src/formatter.rs create mode 100644 duperq/src/processor.rs diff --git a/Cargo.lock b/Cargo.lock index 7c1f04a..d73b59d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -200,6 +200,72 @@ dependencies = [ "winnow", ] +[[package]] +name = "async-channel" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "924ed96dd52d1b75e9c1a3e6275715fd320f5f9439fb5a4a11fa51f4221158d2" +dependencies = [ + "concurrent-queue", + "event-listener-strategy", + "futures-core", + "pin-project-lite", +] + +[[package]] +name = "async-executor" +version = "1.13.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "497c00e0fd83a72a79a39fcbd8e3e2f055d6f6c7e025f3b3d91f4f8e76527fb8" +dependencies = [ + "async-task", + "concurrent-queue", + "fastrand", + "futures-lite", + "pin-project-lite", + "slab", +] + +[[package]] +name = "async-fs" +version = "2.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8034a681df4aed8b8edbd7fbe472401ecf009251c8b40556b304567052e294c5" +dependencies = [ + "async-lock", + "blocking", + "futures-lite", +] + +[[package]] +name = "async-io" +version = "2.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "456b8a8feb6f42d237746d4b3e9a178494627745c3c56c6ea55d92ba50d026fc" +dependencies = [ + "autocfg", + "cfg-if", + "concurrent-queue", + "futures-io", + "futures-lite", + "parking", + "polling", + "rustix 1.1.2", + "slab", + "windows-sys 0.61.2", +] + +[[package]] +name = "async-lock" +version = "3.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5fd03604047cee9b6ce9de9f70c6cd540a0520c813cbd49bae61f33ab80ed1dc" +dependencies = [ + "event-listener", + "event-listener-strategy", + "pin-project-lite", +] + [[package]] name = "async-lsp" version = "0.2.2" @@ -220,6 +286,76 @@ dependencies = [ "waitpid-any", ] +[[package]] +name = "async-net" +version = "2.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b948000fad4873c1c9339d60f2623323a0cfd3816e5181033c6a5cb68b2accf7" +dependencies = [ + "async-io", + "blocking", + "futures-lite", +] + +[[package]] +name = "async-process" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc50921ec0055cdd8a16de48773bfeec5c972598674347252c0399676be7da75" +dependencies = [ + "async-channel", + "async-io", + "async-lock", + "async-signal", + "async-task", + "blocking", + "cfg-if", + "event-listener", + "futures-lite", + "rustix 1.1.2", +] + +[[package]] +name = "async-signal" +version = "0.2.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "43c070bbf59cd3570b6b2dd54cd772527c7c3620fce8be898406dd3ed6adc64c" +dependencies = [ + "async-io", + "async-lock", + "atomic-waker", + "cfg-if", + "futures-core", + "futures-io", + "rustix 1.1.2", + "signal-hook-registry", + "slab", + "windows-sys 0.61.2", +] + +[[package]] +name = "async-task" +version = "4.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8b75356056920673b02621b35afd0f7dda9306d03c79a30f5c56c44cf256e3de" + +[[package]] +name = "async-trait" +version = "0.1.89" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9035ad2d096bed7955a320ee7e2230574d28fd3c3a0f186cbea1ff3c7eed5dbb" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.110", +] + +[[package]] +name = "atomic-waker" +version = "1.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" + [[package]] name = "autocfg" version = "1.5.0" @@ -343,6 +479,19 @@ dependencies = [ "wyz", ] +[[package]] +name = "blocking" +version = "1.6.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e83f8d02be6967315521be875afa792a316e28d57b5a2d401897e2a7921b7f21" +dependencies = [ + "async-channel", + "async-task", + "futures-io", + "futures-lite", + "piper", +] + [[package]] name = "borsh" version = "1.5.7" @@ -583,6 +732,15 @@ dependencies = [ "windows-sys 0.45.0", ] +[[package]] +name = "concurrent-queue" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4ca0197aee26d1ae37445ee532fefce43251d24cc7c166799f4d46817f1d3973" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "console" version = "0.15.11" @@ -865,11 +1023,15 @@ dependencies = [ name = "duperq" version = "0.1.0" dependencies = [ + "anyhow", + "async-trait", "chumsky", "clap", "duper", + "glob", "owo-colors", "regex", + "smol", "temporal_rs", "tinyvec", "tracing", @@ -925,6 +1087,27 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "event-listener" +version = "5.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e13b66accf52311f30a0db42147dadea9850cb48cd070028831ae5f5d4b856ab" +dependencies = [ + "concurrent-queue", + "parking", + "pin-project-lite", +] + +[[package]] +name = "event-listener-strategy" +version = "0.5.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8be9f3dfaaffdae2972880079a491a1a8bb7cbed0b8dd7a347f668b4150a3b93" +dependencies = [ + "event-listener", + "pin-project-lite", +] + [[package]] name = "fallible-iterator" version = "0.3.0" @@ -1027,6 +1210,19 @@ version = "0.3.31" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9e5c1b78ca4aae1ac06c48a526a655760685149f0d465d21f37abfe57ce075c6" +[[package]] +name = "futures-lite" +version = "2.6.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f78e10609fe0e0b3f4157ffab1876319b5b0db102a2c60dc4626306dc46b44ad" +dependencies = [ + "fastrand", + "futures-core", + "futures-io", + "parking", + "pin-project-lite", +] + [[package]] name = "futures-macro" version = "0.3.31" @@ -1776,6 +1972,12 @@ version = "4.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9c6901729fa79e91a0913333229e9ca5dc725089d1c363b2f4b4760709dc4a52" +[[package]] +name = "parking" +version = "2.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f38d5652c16fde515bb1ecef450ab0f6a219d619a7274976324d5e377f7dceba" + [[package]] name = "parking_lot" version = "0.12.5" @@ -1823,12 +2025,37 @@ version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8b870d8c151b6f2fb93e84a13146138f05d02ed11c7e7c54f8826aaaf7c9f184" +[[package]] +name = "piper" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "96c8c490f422ef9a4efd2cb5b42b76c8613d7e7dfc1caf667b8a3350a5acc066" +dependencies = [ + "atomic-waker", + "fastrand", + "futures-io", +] + [[package]] name = "plain" version = "0.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b4596b6d070b27117e987119b4dac604f3c58cfb0b191112e24771b2faeac1a6" +[[package]] +name = "polling" +version = "3.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5d0e4f59085d47d8241c88ead0f274e8a0cb551f3625263c05eb8dd897c34218" +dependencies = [ + "cfg-if", + "concurrent-queue", + "hermit-abi", + "pin-project-lite", + "rustix 1.1.2", + "windows-sys 0.61.2", +] + [[package]] name = "portable-atomic" version = "1.11.1" @@ -2501,6 +2728,23 @@ version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b7c388c1b5e93756d0c740965c41e8822f866621d41acbdf6336a6a168f8840c" +[[package]] +name = "smol" +version = "2.0.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a33bd3e260892199c3ccfc487c88b2da2265080acb316cd920da72fdfd7c599f" +dependencies = [ + "async-channel", + "async-executor", + "async-fs", + "async-io", + "async-lock", + "async-net", + "async-process", + "blocking", + "futures-lite", +] + [[package]] name = "socket2" version = "0.6.1" diff --git a/duper/src/ast.rs b/duper/src/ast.rs index 19d3a19..ed56bfa 100644 --- a/duper/src/ast.rs +++ b/duper/src/ast.rs @@ -283,6 +283,46 @@ impl<'a> DuperValue<'a> { DuperInner::Null => visitor.visit_null(self.identifier.as_ref()), } } + + /// Create a clone of this DuperValue with a static lifetime. + pub fn static_clone(&self) -> DuperValue<'static> { + DuperValue { + identifier: self + .identifier + .as_ref() + .map(|identifier| identifier.static_clone()), + inner: match &self.inner { + DuperInner::Object(object) => DuperInner::Object(DuperObject( + object + .iter() + .map(|(key, value)| { + ( + DuperKey(Cow::Owned(key.0.clone().into_owned())), + value.static_clone(), + ) + }) + .collect(), + )), + DuperInner::Array(array) => DuperInner::Array(DuperArray( + array.iter().map(|element| element.static_clone()).collect(), + )), + DuperInner::Tuple(tuple) => DuperInner::Tuple(DuperTuple( + tuple.iter().map(|element| element.static_clone()).collect(), + )), + DuperInner::String(string) => { + DuperInner::String(DuperString(Cow::Owned(string.0.clone().into_owned()))) + } + DuperInner::Bytes(bytes) => { + DuperInner::Bytes(DuperBytes(Cow::Owned(bytes.0.clone().into_owned()))) + } + DuperInner::Temporal(temporal) => DuperInner::Temporal(temporal.static_clone()), + DuperInner::Integer(integer) => DuperInner::Integer(*integer), + DuperInner::Float(float) => DuperInner::Float(*float), + DuperInner::Boolean(boolean) => DuperInner::Boolean(*boolean), + DuperInner::Null => DuperInner::Null, + }, + } + } } impl<'a> TryFrom<&'a str> for DuperValue<'a> { diff --git a/duper/src/parser/mod.rs b/duper/src/parser/mod.rs index 3b92386..ad76aaa 100644 --- a/duper/src/parser/mod.rs +++ b/duper/src/parser/mod.rs @@ -19,7 +19,6 @@ impl DuperParser { /// /// A pretty-printed version of the error can be obtained from the `prettify_error` method. pub fn parse_duper_trunk<'a>(input: &'a str) -> Result, Vec>> { - // duper_trunk().parse(input).into_result() let value = duper_trunk().parse(input).into_result()?; match &value.inner { DuperInner::Object(_) | DuperInner::Array(_) | DuperInner::Tuple(_) => Ok(value), @@ -116,7 +115,7 @@ pub(crate) fn identifier_lossy<'a>() .map(|string| DuperIdentifier(Cow::Owned(string))) } -pub(crate) fn identifier<'a>() +pub fn identifier<'a>() -> impl Parser<'a, &'a str, DuperIdentifier<'a>, extra::Err>> + Clone { one_of('A'..='Z') .labelled("ASCII uppercase letter") diff --git a/duperq/Cargo.toml b/duperq/Cargo.toml index 5376dbd..bead6b4 100644 --- a/duperq/Cargo.toml +++ b/duperq/Cargo.toml @@ -11,11 +11,15 @@ homepage = "https://duper.dev.br" readme = "README.md" [dependencies] +anyhow = "1.0.100" +async-trait = "0.1.89" chumsky = "0.11.2" clap = { version = "4.5.53", features = ["derive"] } duper = { version = "0.4.2", path = "../duper" } +glob = "0.3.3" owo-colors = "4.2.3" regex = "1.12.2" +smol = "2.0.2" temporal_rs = "0.1.2" tinyvec = { version = "1.10.0", features = ["alloc"] } tracing = "0.1.41" diff --git a/duperq/src/accessor.rs b/duperq/src/accessor.rs index 505805d..44bce65 100644 --- a/duperq/src/accessor.rs +++ b/duperq/src/accessor.rs @@ -2,6 +2,8 @@ use std::{iter, ops::Bound}; use duper::{DuperInner, DuperValue}; +use crate::filter::DuperFilter; + pub(crate) trait DuperAccessor { fn access<'accessor: 'value, 'value>( &'accessor self, @@ -132,3 +134,24 @@ impl DuperAccessor for AnyAccessor { } } } + +pub(crate) struct FilterAccessor(pub(crate) Box); + +impl DuperAccessor for FilterAccessor { + fn access<'accessor: 'value, 'value>( + &'accessor self, + value: &'value DuperValue<'value>, + ) -> Box> + 'value> { + if let DuperInner::Array(array) = &value.inner { + Box::new(array.iter().filter_map(|value| { + if self.0.filter(value) { + Some(value) + } else { + None + } + })) + } else { + Box::new(iter::empty()) + } + } +} diff --git a/duperq/src/filter.rs b/duperq/src/filter.rs index b46321c..b0ffb43 100644 --- a/duperq/src/filter.rs +++ b/duperq/src/filter.rs @@ -15,14 +15,6 @@ pub(crate) trait DuperFilter { // Branchless filters -pub(crate) struct FalseFilter; - -impl DuperFilter for FalseFilter { - fn filter<'v>(&self, _: &'v DuperValue<'v>) -> bool { - false - } -} - pub(crate) struct TrueFilter; impl DuperFilter for TrueFilter { @@ -128,6 +120,7 @@ impl From for TryFromDuperValueError { } pub(crate) enum EqValue { + Identifier(Option), Len(usize), Tuple(Vec), String(String), @@ -152,8 +145,8 @@ impl EqValue { epsilon: Option, ) -> Result { match value.inner { - DuperInner::Object(_) => Err(TryFromDuperValueError::InvalidType("object")), - DuperInner::Array(_) => Err(TryFromDuperValueError::InvalidType("array")), + DuperInner::Object(_) => Err(TryFromDuperValueError::InvalidType("Object")), + DuperInner::Array(_) => Err(TryFromDuperValueError::InvalidType("Array")), DuperInner::Tuple(tuple) => { let vec: Result, _> = tuple .into_inner() @@ -210,6 +203,13 @@ pub(crate) struct EqFilter(pub(crate) EqValue); impl DuperFilter for EqFilter { fn filter<'v>(&self, value: &'v DuperValue<'v>) -> bool { match (&self.0, &value.inner) { + (EqValue::Identifier(this), _) => match this { + Some(this) => value + .identifier + .as_ref() + .is_some_and(|that| this == that.as_ref()), + None => value.identifier.is_none(), + }, (EqValue::Len(this), DuperInner::Object(that)) => *this == that.len(), (EqValue::Len(this), DuperInner::Array(that)) => *this == that.len(), (EqValue::Len(this), DuperInner::String(that)) => *this == that.as_ref().len(), @@ -270,6 +270,13 @@ pub(crate) struct NeFilter(pub(crate) EqValue); impl DuperFilter for NeFilter { fn filter<'v>(&self, value: &'v DuperValue<'v>) -> bool { match (&self.0, &value.inner) { + (EqValue::Identifier(this), _) => match this { + Some(this) => value + .identifier + .as_ref() + .is_none_or(|that| this != that.as_ref()), + None => value.identifier.is_some(), + }, (EqValue::Len(this), DuperInner::Object(that)) => *this != that.len(), (EqValue::Len(this), DuperInner::Array(that)) => *this != that.len(), (EqValue::Len(this), DuperInner::String(that)) => *this != that.as_ref().len(), @@ -348,7 +355,6 @@ pub(crate) enum CmpValue { TemporalPlainTime(PlainTime), TemporalPlainDateTime(PlainDateTime), TemporalPlainYearMonth(PlainYearMonth), - TemporalPlainMonthDay(PlainMonthDay), TemporalDuration(Duration), Integer(i64), Float(f64), @@ -359,11 +365,11 @@ impl TryFrom> for CmpValue { fn try_from(value: DuperValue<'_>) -> Result { match value.inner { - DuperInner::Object(_) => Err(TryFromDuperValueError::InvalidType("object")), - DuperInner::Array(_) => Err(TryFromDuperValueError::InvalidType("array")), - DuperInner::Tuple(_) => Err(TryFromDuperValueError::InvalidType("tuple")), - DuperInner::String(_) => Err(TryFromDuperValueError::InvalidType("string")), - DuperInner::Bytes(_) => Err(TryFromDuperValueError::InvalidType("bytes")), + DuperInner::Object(_) => Err(TryFromDuperValueError::InvalidType("Object")), + DuperInner::Array(_) => Err(TryFromDuperValueError::InvalidType("Array")), + DuperInner::Tuple(_) => Err(TryFromDuperValueError::InvalidType("Tuple")), + DuperInner::String(_) => Err(TryFromDuperValueError::InvalidType("String")), + DuperInner::Bytes(_) => Err(TryFromDuperValueError::InvalidType("Bytes")), DuperInner::Temporal(temporal) => match value.identifier { Some(identifier) if identifier.as_ref() == "Instant" => Ok( CmpValue::TemporalInstant(Instant::from_str(temporal.as_ref())?), @@ -387,9 +393,9 @@ impl TryFrom> for CmpValue { Some(identifier) if identifier.as_ref() == "PlainYearMonth" => Ok( CmpValue::TemporalPlainYearMonth(PlainYearMonth::from_str(temporal.as_ref())?), ), - Some(identifier) if identifier.as_ref() == "PlainMonthDay" => Ok( - CmpValue::TemporalPlainMonthDay(PlainMonthDay::from_str(temporal.as_ref())?), - ), + Some(identifier) if identifier.as_ref() == "PlainMonthDay" => { + Err(TryFromDuperValueError::InvalidType("PlainMonthDay")) + } Some(identifier) if identifier.as_ref() == "Duration" => Ok( CmpValue::TemporalDuration(Duration::from_str(temporal.as_ref())?), ), @@ -397,8 +403,8 @@ impl TryFrom> for CmpValue { }, DuperInner::Integer(integer) => Ok(CmpValue::Integer(integer)), DuperInner::Float(float) => Ok(CmpValue::Float(float)), - DuperInner::Boolean(_) => Err(TryFromDuperValueError::InvalidType("boolean")), - DuperInner::Null => Err(TryFromDuperValueError::InvalidType("null")), + DuperInner::Boolean(_) => Err(TryFromDuperValueError::InvalidType("Boolean")), + DuperInner::Null => Err(TryFromDuperValueError::InvalidType("Null")), } } } @@ -558,18 +564,14 @@ impl DuperFilter for RegexFilter { } } -pub(crate) struct FieldExistsFilter(pub(crate) String); +pub(crate) struct RegexIdentifierFilter(pub(crate) regex::Regex); -impl DuperFilter for FieldExistsFilter { - fn filter<'v>(&self, value: &'v DuperValue<'_>) -> bool { - if let DuperInner::Object(object) = &value.inner { - object - .iter() - .find(|(key, _)| key.as_ref() == self.0) - .is_some() - } else { - false - } +impl DuperFilter for RegexIdentifierFilter { + fn filter<'v>(&self, value: &'v DuperValue<'v>) -> bool { + value + .identifier + .as_ref() + .is_some_and(|identifier| self.0.find(identifier.as_ref()).is_some()) } } diff --git a/duperq/src/formatter.rs b/duperq/src/formatter.rs new file mode 100644 index 0000000..53ab542 --- /dev/null +++ b/duperq/src/formatter.rs @@ -0,0 +1,238 @@ +// duperq 'span.tagged && span[0]name == sp0001 | "[${level}] ${span[0]time} - ${span[0]status} ${telemetry.duration:ms}"' + +use duper::{ + DuperArray, DuperBytes, DuperIdentifier, DuperObject, DuperString, DuperTemporal, DuperTuple, + DuperValue, + format::{ + format_boolean, format_duper_bytes, format_duper_string, format_float, format_integer, + format_key, format_null, format_temporal, + }, + visitor::DuperVisitor, +}; + +use crate::accessor::DuperAccessor; + +pub(crate) enum FormatterAtom { + Fixed(String), + Dynamic(Box), +} + +pub(crate) struct Formatter { + atoms: Vec, + visitor: FormatterVisitor, +} + +impl Formatter { + pub(crate) fn new(atoms: Vec) -> Self { + Self { + atoms, + visitor: FormatterVisitor { buf: String::new() }, + } + } + + pub(crate) fn format(&mut self, value: DuperValue<'static>) -> String { + let mut buf = String::new(); + for atom in &self.atoms { + match atom { + FormatterAtom::Fixed(fixed) => buf.push_str(&fixed), + FormatterAtom::Dynamic(duper_accessor) => { + match duper_accessor.access(&value).into_iter().next() { + Some(value) => buf.push_str(&self.visitor.visit(value)), + None => buf.push_str("-MISSING-"), + } + } + } + } + buf + } +} + +struct FormatterVisitor { + buf: String, +} + +impl FormatterVisitor { + pub fn visit<'a>(&mut self, value: &'a DuperValue<'a>) -> String { + self.buf.clear(); + value.accept(self); + std::mem::take(&mut self.buf) + } +} + +impl DuperVisitor for FormatterVisitor { + type Value = (); + + fn visit_object<'a>( + &mut self, + identifier: Option<&DuperIdentifier<'a>>, + object: &DuperObject<'a>, + ) -> Self::Value { + let len = object.len(); + + if let Some(identifier) = identifier { + self.buf.push_str(identifier.as_ref()); + self.buf.push_str("({"); + for (i, (key, value)) in object.iter().enumerate() { + self.buf.push_str(&format_key(key)); + self.buf.push_str(": "); + value.accept(self); + if i < len - 1 { + self.buf.push_str(", "); + } + } + self.buf.push_str("})"); + } else { + self.buf.push('{'); + for (i, (key, value)) in object.iter().enumerate() { + self.buf.push_str(&format_key(key)); + self.buf.push_str(": "); + value.accept(self); + if i < len - 1 { + self.buf.push_str(", "); + } + } + self.buf.push('}'); + } + } + + fn visit_array<'a>( + &mut self, + identifier: Option<&DuperIdentifier<'a>>, + array: &DuperArray<'a>, + ) -> Self::Value { + let len = array.len(); + + if let Some(identifier) = identifier { + self.buf.push_str(identifier.as_ref()); + self.buf.push_str("(["); + for (i, value) in array.iter().enumerate() { + value.accept(self); + if i < len - 1 { + self.buf.push_str(", "); + } + } + self.buf.push_str("])"); + } else { + self.buf.push('['); + for (i, value) in array.iter().enumerate() { + value.accept(self); + if i < len - 1 { + self.buf.push_str(", "); + } + } + self.buf.push(']'); + } + } + + fn visit_tuple<'a>( + &mut self, + identifier: Option<&DuperIdentifier<'a>>, + tuple: &DuperTuple<'a>, + ) -> Self::Value { + let len = tuple.len(); + + if let Some(identifier) = identifier { + self.buf.push_str(identifier.as_ref()); + self.buf.push_str("(("); + for (i, value) in tuple.iter().enumerate() { + value.accept(self); + if i < len - 1 { + self.buf.push_str(", "); + } + } + self.buf.push_str("))"); + } else { + self.buf.push('('); + for (i, value) in tuple.iter().enumerate() { + value.accept(self); + if i < len - 1 { + self.buf.push_str(", "); + } + } + self.buf.push(')'); + } + } + + fn visit_string<'a>( + &mut self, + identifier: Option<&DuperIdentifier<'a>>, + value: &DuperString<'a>, + ) -> Self::Value { + if let Some(identifier) = identifier { + let value = format_duper_string(value); + self.buf.push_str(&format!("{identifier}({value})")); + } else { + self.buf.push_str(value.as_ref()); + } + } + + fn visit_bytes<'a>( + &mut self, + identifier: Option<&DuperIdentifier<'a>>, + bytes: &DuperBytes<'a>, + ) -> Self::Value { + if let Some(identifier) = identifier { + let bytes = format_duper_bytes(bytes); + self.buf.push_str(&format!("{identifier}({bytes})")); + } else { + self.buf.push_str(&format_duper_bytes(bytes)); + } + } + + fn visit_temporal<'a>( + &mut self, + identifier: Option<&DuperIdentifier<'a>>, + temporal: &DuperTemporal<'a>, + ) -> Self::Value { + if let Some(identifier) = identifier { + let value = format_temporal(temporal); + self.buf.push_str(&format!("{identifier}({value})")); + } else { + self.buf.push_str(&format_temporal(temporal)); + } + } + + fn visit_integer( + &mut self, + identifier: Option<&DuperIdentifier<'_>>, + integer: i64, + ) -> Self::Value { + if let Some(identifier) = identifier { + let value = format_integer(integer); + self.buf.push_str(&format!("{identifier}({value})")); + } else { + self.buf.push_str(&format_integer(integer)); + } + } + + fn visit_float(&mut self, identifier: Option<&DuperIdentifier<'_>>, float: f64) -> Self::Value { + if let Some(identifier) = identifier { + let value = format_float(float); + self.buf.push_str(&format!("{identifier}({value})")); + } else { + self.buf.push_str(&format_float(float)); + } + } + + fn visit_boolean( + &mut self, + identifier: Option<&DuperIdentifier<'_>>, + boolean: bool, + ) -> Self::Value { + if let Some(identifier) = identifier { + let value = format_boolean(boolean); + self.buf.push_str(&format!("{identifier}({value})")); + } else { + self.buf.push_str(format_boolean(boolean)); + } + } + + fn visit_null(&mut self, identifier: Option<&DuperIdentifier<'_>>) -> Self::Value { + if let Some(identifier) = identifier { + let value = format_null(); + self.buf.push_str(&format!("{identifier}({value})")); + } else { + self.buf.push_str(format_null()); + } + } +} diff --git a/duperq/src/lib.rs b/duperq/src/lib.rs index 129e5d8..c1d418c 100644 --- a/duperq/src/lib.rs +++ b/duperq/src/lib.rs @@ -1,3 +1,7 @@ mod accessor; mod filter; +mod formatter; +mod processor; mod query; + +pub use query::query; diff --git a/duperq/src/main.rs b/duperq/src/main.rs index b52bbaf..fc8f19c 100644 --- a/duperq/src/main.rs +++ b/duperq/src/main.rs @@ -1,4 +1,13 @@ +use chumsky::Parser as _; use clap::Parser; +use duper::DuperParser; +use duperq::query; +use glob::glob; +use smol::{ + LocalExecutor, Unblock, + io::{AsyncBufReadExt, AsyncWriteExt, BufReader}, + stream::StreamExt, +}; #[derive(Parser)] #[command(version, about, long_about = None)] @@ -6,10 +15,97 @@ struct Cli { /// Query to run. query: String, - /// Files to read from. If missing, defaults to stdin. - blob: Option, + /// Glob of files to read from. If missing, defaults to stdin. + glob: Option, + + /// If set, disables logs about parsing errors from being printed to stderr. + #[arg(short = 'E', long)] + disable_stderr: bool, } -fn main() { - println!("Hello, world!"); +fn main() -> anyhow::Result<()> { + let cli = Cli::parse(); + smol::block_on(async move { + let mut stderr = Unblock::new(std::io::stderr()); + let (pipeline_fns, output) = match query().parse(&cli.query).into_result() { + Ok(pipeline) => pipeline, + Err(errors) => { + return Err(anyhow::anyhow!(DuperParser::prettify_error( + &cli.query, &errors, None + )?)); + } + }; + + let executor = LocalExecutor::new(); + let mut tasks = Vec::with_capacity(pipeline_fns.len()); + let mut sink = pipeline_fns + .into_iter() + .rfold(output, |mut output, pipeline_fn| { + let (sender, receiver) = smol::channel::bounded(1024); + tasks.push(executor.spawn(async move { + while let Ok(value) = receiver.recv().await { + output.process(value).await; + } + })); + (pipeline_fn)(sender) + }); + + if let Some(duper_glob) = cli.glob { + // Read from files + for entry in glob(&duper_glob)? { + match entry { + Ok(path) => match smol::fs::read_to_string(&path).await { + Ok(input) => match DuperParser::parse_duper_trunk(&input) { + Ok(trunk) => sink.process(trunk.static_clone()).await, + Err(errors) => { + if !cli.disable_stderr { + if let Ok(parse_error) = DuperParser::prettify_error( + &input, + &errors, + Some(path.to_string_lossy().as_ref()), + ) { + let _ = stderr.write_all(parse_error.as_bytes()).await; + } + } + } + }, + Err(error) => { + if !cli.disable_stderr { + let _ = stderr.write_all(error.to_string().as_bytes()).await; + } + } + }, + Err(error) => { + if !cli.disable_stderr { + let _ = stderr.write_all(error.to_string().as_bytes()).await; + } + } + } + } + } else { + // Read from stdin + let stdin = BufReader::new(Unblock::new(std::io::stdin())); + let mut lines = stdin.lines(); + while let Some(Ok(line)) = lines.next().await { + match DuperParser::parse_duper_trunk(&line) { + Ok(trunk) => sink.process(trunk.static_clone()).await, + Err(errors) => { + if !cli.disable_stderr { + if let Ok(parse_error) = + DuperParser::prettify_error(&line, &errors, None) + { + let _ = stderr.write_all(parse_error.as_bytes()).await; + } + } + } + } + } + } + + for task in tasks.into_iter().rev() { + task.await; + } + + Ok(()) + }) } diff --git a/duperq/src/processor.rs b/duperq/src/processor.rs new file mode 100644 index 0000000..3fcdca0 --- /dev/null +++ b/duperq/src/processor.rs @@ -0,0 +1,94 @@ +use async_trait::async_trait; +use duper::DuperValue; +use smol::{Unblock, channel, io::AsyncWriteExt}; + +use crate::filter::DuperFilter; + +#[async_trait(?Send)] +pub trait Processor { + async fn process(&mut self, value: DuperValue<'static>); + + async fn close(&mut self) {} +} + +pub(crate) struct FilterProcessor { + filter: Box, + sender: channel::Sender>, + is_open: bool, +} + +impl FilterProcessor { + pub(crate) fn new( + sender: channel::Sender>, + filter: Box, + ) -> Self { + Self { + is_open: true, + sender, + filter, + } + } +} + +#[async_trait(?Send)] +impl Processor for FilterProcessor { + async fn process(&mut self, value: DuperValue<'static>) { + if self.is_open && self.filter.filter(&value) { + if self.sender.send(value).await.is_err() { + self.is_open = false; + } + } + } +} + +pub(crate) struct TakeProcessor { + available: usize, + sender: channel::Sender>, +} + +impl TakeProcessor { + pub(crate) fn new(sender: channel::Sender>, available: usize) -> Self { + Self { sender, available } + } +} + +#[async_trait(?Send)] +impl Processor for TakeProcessor { + async fn process(&mut self, value: DuperValue<'static>) { + if self.available > 0 { + if self.sender.send(value).await.is_err() { + self.available = 0; + } else { + self.available = self.available.saturating_sub(1); + } + } + } +} + +pub(crate) struct OutputProcessor { + stdout: Unblock, + printer: Box) -> String>, +} + +impl OutputProcessor { + pub(crate) fn new(printer: Box) -> String>) -> Self { + Self { + stdout: Unblock::new(std::io::stdout()), + printer, + } + } +} + +#[async_trait(?Send)] +impl Processor for OutputProcessor { + async fn process(&mut self, value: DuperValue<'static>) { + self.stdout + .write_all((self.printer)(value).as_bytes()) + .await + .expect("stdout was closed"); + self.stdout + .write_all(b"\n") + .await + .expect("stdout was closed"); + } +} diff --git a/duperq/src/query.rs b/duperq/src/query.rs index 2ed27fd..7500fb8 100644 --- a/duperq/src/query.rs +++ b/duperq/src/query.rs @@ -1,44 +1,78 @@ -// duperq 'span.tagged && span[0]name == sp0001 | "[${level}] ${span[0]time} - ${span[0]status} ${telemetry.duration:ms}"' -// duperq 'metadata.tags[created_at >= Instant(2025-11-22T00:00:00-03:00)]' - use chumsky::prelude::*; use duper::{ - DuperInner, - parser::{duper_value, integer, object_key}, + DuperInner, DuperValue, PrettyPrinter, Serializer, + escape::unescape_str, + parser::{duper_value, identifier, integer, object_key}, }; +use smol::channel; use crate::{ accessor::{ - AnyAccessor, DuperAccessor, FieldAccessor, FlattenedAccessor, IndexAccessor, - RangeIndexAccessor, ReverseIndexAccessor, + AnyAccessor, DuperAccessor, FieldAccessor, FilterAccessor, FlattenedAccessor, + IndexAccessor, RangeIndexAccessor, ReverseIndexAccessor, }, filter::{ AccessorFilter, AndFilter, CmpValue, DuperFilter, EqFilter, EqValue, GeFilter, GtFilter, IsFilter, IsTruthyFilter, LeFilter, LtFilter, NeFilter, NotFilter, OrFilter, RegexFilter, - TryFromDuperValueError, + RegexIdentifierFilter, TrueFilter, TryFromDuperValueError, }, + formatter::{Formatter, FormatterAtom}, + processor::{FilterProcessor, OutputProcessor, Processor, TakeProcessor}, }; -fn query<'a>() --> impl Parser<'a, &'a str, (Box, Option<()>), extra::Err>> { - choice((just("filter").padded().ignore_then(filter()),)) - .separated_by(just('|')) - .collect::() - .map(|filter| Box::new(filter) as Box) - .padded() - .then( - just('|') - .padded() - .ignore_then(just("format").padded()) - .ignore_then(fmt().padded()) - .or_not(), - ) +pub(crate) type CreateProcessorFn = + Box>) -> Box>; + +pub fn query<'a>() +-> impl Parser<'a, &'a str, (Vec, Box), extra::Err>> +{ + choice(( + just("filter").padded().ignore_then(filter()).map(|filter| { + Box::new(move |sender| { + Box::new(FilterProcessor::new(sender, filter)) as Box + }) as CreateProcessorFn + }), + just("take") + .padded() + .ignore_then(integer()) + .try_map(|take, span| { + if take > 0 { + Ok(Box::new(move |sender| { + Box::new(TakeProcessor::new(sender, take as usize)) as Box + }) as CreateProcessorFn) + } else { + Err(Rich::custom( + span, + "take parameter must be greater than zero", + )) + } + }), + )) + .separated_by(just('|')) + .collect() + .then( + just('|') + .padded() + .ignore_then(just("format").padded().ignore_then(fmt().padded()).or( + just("pretty-print").padded().map(|_| { + let mut pretty_printer = PrettyPrinter::default(); + OutputProcessor::new(Box::new(move |value| pretty_printer.pretty_print(value))) + }), + )) + .or_not() + .map(|processor| { + Box::new(processor.unwrap_or_else(|| { + let mut serializer = Serializer::default(); + OutputProcessor::new(Box::new(move |value| serializer.serialize(value))) + })) as Box + }), + ) } fn filter<'a>() -> impl Parser<'a, &'a str, Box, extra::Err>> + Clone { recursive(|filter| { - let atom = leaf_filter() + let atom = leaf_filter(accessor()) .or(filter.delimited_by(just('('), just(')'))) .padded(); @@ -74,84 +108,93 @@ fn filter<'a>() -> impl Parser<'a, &'a str, Box, extra::Err() -> impl Parser<'a, &'a str, Box, extra::Err>> + Clone { - let access = just('.').or_not().ignore_then(choice(( - object_key().padded().map(|key: duper::DuperKey<'a>| { - Box::new(FieldAccessor(key.as_ref().into())) as Box - }), - integer() - .or_not() - .padded() - .then_ignore(just("..")) - .then(just('=').or_not()) - .then(integer().or_not().padded()) - .delimited_by(just('['), just(']')) - .padded() - .try_map(|((start, end_inclusive), end), span| match (start, end) { - (Some(start), _) if start < 0 => { - Err(Rich::custom(span, "range start must be positive")) - } - (_, Some(end)) if end < 0 => Err(Rich::custom(span, "range end must be positive")), - (None, None) => Ok(RangeIndexAccessor { - start: std::ops::Bound::Unbounded, - end: std::ops::Bound::Unbounded, - }), - (Some(start), None) => Ok(RangeIndexAccessor { - start: std::ops::Bound::Included(start as usize), - end: std::ops::Bound::Unbounded, - }), - (None, Some(end)) => Ok(RangeIndexAccessor { - start: std::ops::Bound::Unbounded, - end: if end_inclusive.is_some() { - std::ops::Bound::Included(end as usize) - } else { - std::ops::Bound::Excluded(end as usize) - }, - }), - (Some(start), Some(end)) => Ok(RangeIndexAccessor { - start: std::ops::Bound::Included(start as usize), - end: if end_inclusive.is_some() { - std::ops::Bound::Included(end as usize) + recursive(|accessor| { + let access = just('.').or_not().ignore_then(choice(( + object_key().padded().map(|key: duper::DuperKey<'a>| { + Box::new(FieldAccessor(key.as_ref().into())) as Box + }), + integer() + .or_not() + .padded() + .then_ignore(just("..")) + .then(just('=').or_not()) + .then(integer().or_not().padded()) + .delimited_by(just('['), just(']')) + .padded() + .try_map(|((start, end_inclusive), end), span| match (start, end) { + (Some(start), _) if start < 0 => { + Err(Rich::custom(span, "range start must be positive")) + } + (_, Some(end)) if end < 0 => { + Err(Rich::custom(span, "range end must be positive")) + } + (None, None) => Ok(RangeIndexAccessor { + start: std::ops::Bound::Unbounded, + end: std::ops::Bound::Unbounded, + }), + (Some(start), None) => Ok(RangeIndexAccessor { + start: std::ops::Bound::Included(start as usize), + end: std::ops::Bound::Unbounded, + }), + (None, Some(end)) => Ok(RangeIndexAccessor { + start: std::ops::Bound::Unbounded, + end: if end_inclusive.is_some() { + std::ops::Bound::Included(end as usize) + } else { + std::ops::Bound::Excluded(end as usize) + }, + }), + (Some(start), Some(end)) => Ok(RangeIndexAccessor { + start: std::ops::Bound::Included(start as usize), + end: if end_inclusive.is_some() { + std::ops::Bound::Included(end as usize) + } else { + std::ops::Bound::Excluded(end as usize) + }, + }), + }) + .map(|accessor| Box::new(accessor) as Box), + integer() + .padded() + .delimited_by(just('['), just(']')) + .padded() + .map(|int| { + if int < 0 { + Box::new(ReverseIndexAccessor(int.unsigned_abs() as usize)) + as Box } else { - std::ops::Bound::Excluded(end as usize) - }, + Box::new(IndexAccessor(int as usize)) as Box + } }), - }) - .map(|accessor| Box::new(accessor) as Box), - integer() - .padded() - .delimited_by(just('['), just(']')) - .padded() - .map(|int| { - if int < 0 { - Box::new(ReverseIndexAccessor(int.unsigned_abs() as usize)) - as Box - } else { - Box::new(IndexAccessor(int as usize)) as Box - } - }), - text::whitespace() - .delimited_by(just('['), just(']')) - .padded() - .map(|_| Box::new(AnyAccessor) as Box), - ))); + leaf_filter(accessor) + .padded() + .delimited_by(just('['), just(']')) + .padded() + .map(|filter| Box::new(FilterAccessor(filter)) as Box), + text::whitespace() + .delimited_by(just('['), just(']')) + .padded() + .map(|_| Box::new(AnyAccessor) as Box), + ))); - access - .clone() - .repeated() - .at_least(2) - .collect::>() - .map(|vec| Box::new(FlattenedAccessor(vec)) as Box) - .or(access) + access + .clone() + .repeated() + .at_least(2) + .collect::>() + .map(|vec| Box::new(FlattenedAccessor(vec)) as Box) + .or(access) + }) } -fn leaf_filter<'a>() --> impl Parser<'a, &'a str, Box, extra::Err>> + Clone { - let accessor = accessor(); - +fn leaf_filter<'a>( + accessor: impl Parser<'a, &'a str, Box, extra::Err>> + Clone, +) -> impl Parser<'a, &'a str, Box, extra::Err>> + Clone { let eq_op = just("==").ignored().or(just('=').ignored()).padded(); let ne_op = just("!=").ignored().or(just("<>").ignored()).padded(); let lt_op = just("<").ignored().padded(); @@ -161,7 +204,7 @@ fn leaf_filter<'a>() let re_op = just("=~").ignored().padded(); let is_op = just("is").ignored().padded(); - just("len") + let len_filter = just("len") .ignore_then(accessor.clone().delimited_by(just('('), just(')'))) .then(choice(( eq_op @@ -248,8 +291,56 @@ fn leaf_filter<'a>() )) } }), - ))) - .or(accessor.clone().then(choice(( + ))); + + let identifier_filter = just("identifier") + .ignore_then(accessor.clone().delimited_by(just('('), just(')'))) + .then(choice(( + eq_op + .clone() + .ignore_then( + identifier() + .map(|identifier| Some(identifier.to_string())) + .or(just("null").map(|_| None)) + .padded(), + ) + .map(|value| { + Box::new(EqFilter(EqValue::Identifier(value))) as Box + }), + ne_op + .clone() + .ignore_then( + identifier() + .map(|identifier| Some(identifier.to_string())) + .or(just("null").map(|_| None)) + .padded(), + ) + .map(|value| { + Box::new(NeFilter(EqValue::Identifier(value))) as Box + }), + re_op + .clone() + .ignore_then(duper_value().padded()) + .try_map(|value, span| match value.inner { + DuperInner::String(string) => regex::Regex::new(string.as_ref()) + .map(|regex| Box::new(RegexIdentifierFilter(regex)) as Box) + .map_err(|error| Rich::custom(span, error)), + _ => Err(Rich::custom( + span, + "can only use regex operator =~ with string", + )), + }), + ))); + + let exists_filter = just("exists") + .ignore_then(accessor.clone().delimited_by(just('('), just(')'))) + .map(|accessor| (accessor, Box::new(TrueFilter) as Box)); + + choice(( + len_filter, + identifier_filter, + exists_filter, + accessor.clone().then(choice(( eq_op .ignore_then(duper_value().padded()) .try_map(|value, span| { @@ -329,21 +420,35 @@ fn leaf_filter<'a>() .padded(), ) .map(|value| Box::new(value) as Box), - )))) - .map(|(accessor, filter)| { - Box::new(AccessorFilter { accessor, filter }) as Box - }) - .or(accessor.map(|accessor| { - Box::new(AccessorFilter { - accessor, - filter: Box::new(IsTruthyFilter), - }) as Box - })) + ))), + )) + .map(|(accessor, filter)| Box::new(AccessorFilter { accessor, filter }) as Box) + .or(accessor.map(|accessor| { + Box::new(AccessorFilter { + accessor, + filter: Box::new(IsTruthyFilter), + }) as Box + })) } -fn fmt<'a>() -> impl Parser<'a, &'a str, (), extra::Err>> { - any() +fn fmt<'a>() -> impl Parser<'a, &'a str, OutputProcessor, extra::Err>> { + just('$') + .ignore_then(accessor().padded().delimited_by(just('{'), just('}'))) + .map(|accessor| FormatterAtom::Dynamic(accessor)) + .or(any() + .and_is(just("${").not()) + .repeated() + .at_least(1) + .to_slice() + .try_map(|slice: &str, span| match unescape_str(slice) { + Ok(unescaped) => Ok(FormatterAtom::Fixed(unescaped.clone().into_owned())), + Err(error) => Err(Rich::custom(span, error.to_string())), + })) .repeated() + .collect::>() .delimited_by(just('"'), just('"')) - .ignored() + .map(|atoms| { + let mut formatter = Formatter::new(atoms); + OutputProcessor::new(Box::new(move |value| formatter.format(value))) + }) } -- 2.51.2