use anyhow::Context;
use clap::Parser;
use quarb::{AstAdapter, NodeId, QueryResult, Value};
use quarb_athena::AthenaAdapter;
use quarb_atrep::AtrepAdapter;
use quarb_bigquery::BigqueryAdapter;
use quarb_compose::ComposeAdapter;
use quarb_csv::CsvAdapter;
use quarb_duckdb::DuckdbAdapter;
use quarb_html::HtmlAdapter;
use quarb_json::JsonAdapter;
#[cfg(feature = "kuzu")]
use quarb_kuzu::KuzuAdapter;
use quarb_mount::{Latent, Mount, MountAdapter};
use quarb_mysql::MysqlAdapter;
use quarb_open::*;
pub use quarb_open::{Mounted, OpenOpts, SourceCtx, open_target};
use quarb_oracle::OracleAdapter;
use quarb_postgres::PostgresAdapter;
use quarb_sqlite::SqliteAdapter;
use quarb_xml::XmlAdapter;
use std::io::IsTerminal;
use std::path::{Path, PathBuf};
use std::rc::Rc;
#[derive(Parser, Default)]
#[command(version, about)]
struct Cli {
query: String,
paths: Vec<PathBuf>,
#[arg(long)]
hidden: bool,
#[arg(long = "no-ignore")]
no_ignore: bool,
#[arg(long)]
xpath: bool,
#[arg(long, conflicts_with = "xpath")]
jq: bool,
#[arg(long, conflicts_with_all = ["xpath", "jq"])]
sql: bool,
#[arg(long)]
kaiv: bool,
#[arg(long, value_name = "N", default_value_t = 8, requires = "kaiv")]
kaiv_origins: usize,
#[arg(long)]
reproducible: bool,
#[arg(long, conflicts_with_all = ["jsonl", "kaiv"])]
json: bool,
#[arg(long, conflicts_with_all = ["json", "kaiv"])]
jsonl: bool,
#[arg(long, conflicts_with_all = ["json", "jsonl", "kaiv", "csv"])]
table: bool,
#[arg(long, conflicts_with_all = ["json", "jsonl", "kaiv", "table"])]
csv: bool,
#[arg(long, value_name = "FILE")]
defs: Option<PathBuf>,
#[arg(long)]
expand: bool,
#[arg(long = "expand-1", conflicts_with = "expand")]
expand_1: bool,
#[arg(long = "no-pushdown")]
no_pushdown: bool,
#[arg(long)]
explain: bool,
#[arg(long)]
plan: bool,
#[arg(long, value_name = "BOUNDS")]
budget: Option<String>,
#[arg(long, value_name = "TAG")]
locale: Option<String>,
#[arg(long)]
grouping: bool,
#[arg(long = "highlight-html", hide = true)]
highlight_html: bool,
#[arg(long, value_name = "FILE")]
save: Option<PathBuf>,
#[arg(long = "as", value_name = "NAME", default_value = "result")]
save_as: String,
#[arg(long)]
graft: bool,
#[arg(long = "no-graft", conflicts_with = "graft")]
no_graft: bool,
#[arg(long, value_name = "FILE")]
refs: Option<PathBuf>,
#[arg(long, value_name = "FILE")]
model: Option<PathBuf>,
#[arg(long, value_name = "FILE")]
desm: Vec<PathBuf>,
#[arg(long, value_name = "N")]
quantifier_bound: Option<usize>,
#[arg(long)]
allow_shell: bool,
#[arg(long, value_name = "ISO")]
now: Option<String>,
#[arg(long)]
resident: bool,
#[arg(long, value_name = "SECS", default_value_t = 1800)]
resident_ttl: u64,
#[arg(long, hide = true)]
resident_serve: bool,
#[arg(long)]
highlight: bool,
#[arg(long)]
cache: bool,
#[arg(long, value_name = "DIR")]
cache_dir: Option<PathBuf>,
}
#[derive(Clone, Copy, PartialEq, Eq)]
enum Output {
Quarb,
Json,
Jsonl,
Table,
Csv,
}
#[derive(Clone)]
struct OutputCtx {
output: Output,
expand: bool,
expand_1: bool,
kaiv_origins: usize,
save: Option<(PathBuf, String)>,
quant_bound: Option<usize>,
allow_shell: bool,
environments: Vec<quarb::Environment>,
reproducible: bool,
now: (i64, u32),
explain: bool,
plan: bool,
budget: Option<quarb::plan::Budget>,
model: Option<quarb_model::Model>,
locale: Option<(quarb::locale::OutputLocale, &'static str)>,
resident: Option<(PathBuf, u64, bool)>,
encoding: Option<(String, quarb::source::Escape)>,
}
fn split_scheme_query(q: &str) -> Option<(&'static str, &str)> {
for scheme in ["github:", "gitlab:", "k8s:", "kubernetes:"] {
if let Some(rest) = q.strip_prefix(scheme)
&& rest.starts_with('/')
{
return Some((scheme, rest));
}
}
None
}
impl From<&Cli> for OpenOpts {
fn from(cli: &Cli) -> Self {
OpenOpts {
hidden: cli.hidden,
no_ignore: cli.no_ignore,
graft: cli.graft,
no_graft: cli.no_graft,
refs: cli.refs.clone(),
allow_shell: cli.allow_shell,
}
}
}
fn is_prefixed_target(s: &str) -> bool {
s.split_once(':').is_some_and(|(pre, _)| {
pre.len() > 1
&& pre
.chars()
.all(|c| c.is_ascii_alphanumeric() || "+.-".contains(c))
})
}
fn is_relational_target(s: &str) -> bool {
s.starts_with("mssql://")
|| s.starts_with("oracle://")
|| s.starts_with("bigquery://")
|| s.starts_with("athena:")
|| s.starts_with("mysql://")
|| is_pg_config(s)
}
pub fn cli_main() -> anyhow::Result<()> {
#[cfg(unix)]
unsafe {
libc::signal(libc::SIGPIPE, libc::SIG_DFL);
}
if std::env::var_os("KAIV_OFFLINE").is_none() {
unsafe {
std::env::set_var("KAIV_OFFLINE", "1");
}
}
let mut cli = Cli::parse();
let environments = take_environments(&mut cli)?;
let locale = set_output_locale(&cli)?;
let output = if cli.json {
Output::Json
} else if cli.jsonl {
Output::Jsonl
} else if cli.table {
Output::Table
} else if cli.csv {
Output::Csv
} else {
Output::Quarb
};
if cli.highlight_html {
use std::io::BufRead;
for line in std::io::stdin().lock().lines() {
println!("{}", quarb::highlight::highlight_html(&line?));
}
return Ok(());
}
if cli.paths.is_empty()
&& let Some((scheme, query)) = split_scheme_query(&cli.query)
{
cli.paths.push(PathBuf::from(scheme));
cli.query = query.to_string();
}
if cli.xpath {
let translation = quarb_xpath::translate(&cli.query).context("translating XPath")?;
for note in &translation.notes {
eprintln!("note: {note}");
}
cli.query = translation.query;
}
if cli.jq {
let translation = quarb_jq::translate(&cli.query).context("translating jq")?;
for note in &translation.notes {
eprintln!("note: {note}");
}
cli.query = translation.query;
}
if cli.sql {
let translation = quarb_sql::translate(&cli.query).context("translating SQL")?;
for note in &translation.notes {
eprintln!("note: {note}");
}
cli.query = translation.query;
}
if cli.highlight {
if std::env::var_os("NO_COLOR").is_some() {
println!("{}", cli.query);
} else {
println!("{}", quarb::highlight::highlight_ansi(&cli.query));
}
return Ok(());
}
if let Some(defs_path) = &cli.defs {
let text = std::fs::read_to_string(defs_path)
.with_context(|| format!("reading {}", defs_path.display()))?;
let text = text.strip_prefix('\u{feff}').unwrap_or(&text).to_owned();
quarb::parse_defs(&text)
.with_context(|| format!("parsing definitions in {}", defs_path.display()))?;
cli.query = format!("{}\n{}", quarb::strip_defs_comments(&text), cli.query);
}
if cli.expand {
if cli.paths.is_empty() && cli.model.is_none() {
println!(
"{}",
quarb::expand(&cli.query, &quarb::Defs::default())
.context("expanding the query")?
);
return Ok(());
}
}
if cli.expand_1 {
if cli.paths.is_empty() {
for t in quarb::expand_first(&cli.query, &quarb::Defs::default())
.context("expanding the query")?
{
println!("{t}");
}
return Ok(());
}
}
if let Some(n) = cli.quantifier_bound {
anyhow::ensure!(n >= 1, "--quantifier-bound must be at least 1");
}
let now = match &cli.now {
Some(text) => {
let (secs, nanos, _) = quarb::temporal::parse_iso(text)
.ok_or_else(|| anyhow::anyhow!("--now needs an ISO-8601 instant, got '{text}'"))?;
(secs, nanos)
}
None => {
let since = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default();
(since.as_secs() as i64, since.subsec_nanos())
}
};
quarb::set_invocation_instant(now.0, now.1);
if !cli.desm.is_empty() {
let mut bonds = syndesmos::Syndesmos::empty();
for f in &cli.desm {
bonds.extend(
syndesmos::Syndesmos::load(f)
.map_err(|e| anyhow::anyhow!("reading {}: {e}", f.display()))?,
);
}
quarb::set_sentence_bonds(Some(Rc::new(bonds)));
}
let mut parsed_model: Option<quarb_model::Model> = None;
if let Some(model_path) = &cli.model {
let model = quarb_model::parse_model_file(model_path)
.map_err(|e| anyhow::anyhow!("parsing model {}: {e}", model_path.display()))?;
let base_dir = model_path.parent();
for m in &model.mounts {
let target = quarb_model::resolve_mount_target(&m.target, base_dir);
if quarb::Environment::is_environment_target(&target) {
anyhow::bail!(
"model {}: mount '{}' names an environment ({target}); environments \
are mounted from the command line only",
model_path.display(),
m.name
);
}
cli.paths
.insert(0, PathBuf::from(format!("{}={}", m.name, target)));
}
if !model.defs_text.trim().is_empty() {
quarb::parse_defs(&model.defs_text)
.with_context(|| format!("parsing definitions in {}", model_path.display()))?;
cli.query = format!("{}\n{}", model.defs_text, cli.query);
}
parsed_model = Some(model);
if cli.expand && cli.paths.is_empty() {
println!(
"{}",
quarb::expand(&cli.query, &quarb::Defs::default())
.context("expanding the query")?
);
return Ok(());
}
}
if cli.cache || cli.cache_dir.is_some() {
let dir = cli
.cache_dir
.clone()
.unwrap_or_else(quarb_tree_sitter::Cache::default_dir);
quarb_tree_sitter::set_cache(Some(quarb_tree_sitter::Cache::new(dir)));
}
if cli.resident || cli.resident_serve {
anyhow::ensure!(
!cli.kaiv && cli.save.is_none() && !cli.expand && !cli.expand_1,
"--resident does not combine with --kaiv/--save/--expand"
);
anyhow::ensure!(
!cli.paths.is_empty(),
"--resident needs file/directory inputs (stdin has no session identity)"
);
#[cfg(not(unix))]
anyhow::bail!("--resident needs Unix domain sockets (unavailable on this platform)");
}
let mut oc = OutputCtx {
output,
expand: cli.expand,
expand_1: cli.expand_1,
kaiv_origins: cli.kaiv_origins,
save: cli.save.as_ref().map(|p| (p.clone(), cli.save_as.clone())),
quant_bound: cli.quantifier_bound,
allow_shell: cli.allow_shell,
environments: environments.clone(),
reproducible: cli.reproducible,
now,
explain: cli.explain,
plan: cli.plan,
budget: cli.budget.as_deref().map(parse_budget).transpose()?,
model: parsed_model,
locale,
resident: None,
encoding: None,
};
#[cfg(unix)]
{
if cli.resident && !cli.resident_serve {
return resident_client(&cli, &oc);
}
if cli.resident_serve {
let sock = resident_socket(&cli, &oc)?;
oc.resident = Some((sock, cli.resident_ttl, cli.now.is_some()));
}
}
execute(&cli, &cli.query, &oc)
}
#[cfg(unix)]
fn resident_socket(cli: &Cli, oc: &OutputCtx) -> anyhow::Result<PathBuf> {
use std::hash::{Hash, Hasher};
let mut h = std::collections::hash_map::DefaultHasher::new();
for p in &cli.paths {
std::fs::canonicalize(p)
.unwrap_or_else(|_| p.clone())
.hash(&mut h);
}
(
cli.graft,
cli.no_graft,
cli.hidden,
cli.no_ignore,
cli.allow_shell,
cli.quantifier_bound,
&cli.now,
&cli.refs,
&cli.defs,
&cli.model,
cli.no_pushdown,
oc.locale.map(|(l, _)| l).map(|l| l.tag),
)
.hash(&mut h);
(
cli.json,
cli.jsonl,
cli.csv,
cli.table,
cli.kaiv,
cli.grouping,
)
.hash(&mut h);
std::env::current_dir().ok().hash(&mut h);
let dir = resident_dir()?;
Ok(dir.join(format!("quarb-{:016x}.sock", h.finish())))
}
#[cfg(unix)]
fn resident_dir() -> anyhow::Result<PathBuf> {
use std::os::unix::fs::{MetadataExt as _, PermissionsExt as _};
if let Some(d) = std::env::var_os("XDG_RUNTIME_DIR") {
return Ok(PathBuf::from(d));
}
let uid = unsafe { libc::getuid() };
let d = std::env::temp_dir().join(format!("quarb-{uid}"));
let _ = std::fs::create_dir(&d);
let _ = std::fs::set_permissions(&d, std::fs::Permissions::from_mode(0o700));
let ok = std::fs::symlink_metadata(&d).is_ok_and(|m| {
m.file_type().is_dir() && m.uid() == uid && m.permissions().mode() & 0o777 == 0o700
});
anyhow::ensure!(
ok,
"{} is not a private directory owned by this user \
(another user may have created it); remove it or set \
XDG_RUNTIME_DIR to use resident sessions",
d.display()
);
Ok(d)
}
#[cfg(unix)]
fn resident_client(cli: &Cli, oc: &OutputCtx) -> anyhow::Result<()> {
use std::io::Write as _;
let sock = resident_socket(cli, oc)?;
let mut stream = match std::os::unix::net::UnixStream::connect(&sock) {
Ok(s) => s,
Err(_) => spawn_resident(&sock)?,
};
let q = cli.query.as_bytes();
stream.write_all(format!("Q {}\n", q.len()).as_bytes())?;
stream.write_all(q)?;
stream.flush()?;
let mut reader = std::io::BufReader::new(stream);
let mut header = String::new();
std::io::BufRead::read_line(&mut reader, &mut header)?;
let mut parts = header.trim_end().split(' ');
anyhow::ensure!(
parts.next() == Some("R"),
"bad session response: {header:?}"
);
let len: usize = parts
.next()
.and_then(|s| s.parse().ok())
.context("bad session response length")?;
let status: u8 = parts
.next()
.and_then(|s| s.parse().ok())
.context("bad session response status")?;
let mut body = vec![0u8; len];
std::io::Read::read_exact(&mut reader, &mut body)?;
if status == 0 {
std::io::stdout().write_all(&body)?;
Ok(())
} else {
anyhow::bail!("{}", String::from_utf8_lossy(&body));
}
}
#[cfg(unix)]
fn spawn_resident(sock: &std::path::Path) -> anyhow::Result<std::os::unix::net::UnixStream> {
use std::os::unix::process::CommandExt as _;
let log = sock.with_extension("log");
let logfile =
std::fs::File::create(&log).with_context(|| format!("creating {}", log.display()))?;
let exe = std::env::current_exe().context("resolving qua binary")?;
let mut cmd = std::process::Command::new(exe);
if std::env::var_os("QUARB_RESIDENT_CHILD").is_some() {
anyhow::bail!("a resident session does not spawn a session of its own");
}
cmd.arg("--resident-serve")
.args(std::env::args_os().skip(1))
.env("QUARB_RESIDENT_CHILD", "1")
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::from(logfile));
unsafe {
cmd.pre_exec(|| {
libc::setsid();
Ok(())
});
}
let mut child = cmd.spawn().context("starting resident session")?;
eprintln!(
"resident session starting (first query pays materialization; \
log: {})",
log.display()
);
let started = std::time::Instant::now();
let mut last_note = 0u64;
loop {
if let Ok(s) = std::os::unix::net::UnixStream::connect(sock) {
return Ok(s);
}
if let Some(status) = child.try_wait()? {
if let Ok(s) = std::os::unix::net::UnixStream::connect(sock) {
return Ok(s);
}
let tail = std::fs::read_to_string(&log).unwrap_or_default();
let tail = tail.lines().rev().take(5).collect::<Vec<_>>();
anyhow::bail!(
"resident session exited ({status}) before binding its socket:\n{}",
tail.into_iter().rev().collect::<Vec<_>>().join("\n")
);
}
let elapsed = started.elapsed().as_secs();
if elapsed >= last_note + 15 {
eprintln!(" … materializing ({elapsed}s)");
last_note = elapsed;
}
std::thread::sleep(std::time::Duration::from_millis(200));
}
}
#[cfg(unix)]
const RESIDENT_MAX_QUERY: usize = 1 << 20;
#[cfg(unix)]
fn resident_serve_loop<A: AstAdapter>(
adapter: &A,
render: impl Fn(NodeId) -> String,
sock: &std::path::Path,
ttl: u64,
now_pinned: bool,
oc: &OutputCtx,
) -> anyhow::Result<()> {
use std::io::Write as _;
unsafe {
libc::signal(libc::SIGPIPE, libc::SIG_IGN);
}
let listener = match std::os::unix::net::UnixListener::bind(sock) {
Ok(l) => l,
Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => {
if std::os::unix::net::UnixStream::connect(sock).is_ok() {
eprintln!("resident session already live; deferring to it");
return Ok(());
}
let _ = std::fs::remove_file(sock);
std::os::unix::net::UnixListener::bind(sock)
.with_context(|| format!("binding {}", sock.display()))?
}
Err(e) => return Err(e).with_context(|| format!("binding {}", sock.display())),
};
let _ = std::fs::set_permissions(sock, {
use std::os::unix::fs::PermissionsExt as _;
std::fs::Permissions::from_mode(0o600)
});
listener.set_nonblocking(true)?;
let mut idle = std::time::Instant::now();
loop {
match listener.accept() {
Ok((mut conn, _)) => {
idle = std::time::Instant::now();
conn.set_nonblocking(false)?;
let _ = conn.set_read_timeout(Some(std::time::Duration::from_secs(30)));
let _ = conn.set_write_timeout(Some(std::time::Duration::from_secs(30)));
let mut reader = std::io::BufReader::new(conn.try_clone()?);
let mut header = String::new();
if std::io::BufRead::read_line(&mut reader, &mut header).is_err() {
continue;
}
let len: usize = match header
.trim_end()
.strip_prefix("Q ")
.and_then(|s| s.parse().ok())
{
Some(n) if n <= RESIDENT_MAX_QUERY => n,
Some(_) => {
let msg = b"query exceeds the resident frame limit";
let _ = conn.write_all(format!("R {} 1\n", msg.len()).as_bytes());
let _ = conn.write_all(msg);
continue;
}
None => continue,
};
let mut qbytes = vec![0u8; len];
if std::io::Read::read_exact(&mut reader, &mut qbytes).is_err() {
continue;
}
let query = String::from_utf8_lossy(&qbytes).into_owned();
let (query, encoding) = match take_output_encoding(&query) {
Ok(q) => q,
Err(e) => {
let body = format!("{e:#}").into_bytes();
let _ = conn.write_all(format!("R {} 1\n", body.len()).as_bytes());
let _ = conn.write_all(&body);
continue;
}
};
let mut now = oc.now;
if !now_pinned {
let since = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default();
now = (since.as_secs() as i64, since.subsec_nanos());
quarb::set_invocation_instant(now.0, now.1);
}
let per_query = OutputCtx {
encoding,
now,
..oc.clone()
};
let (result, output) =
with_stdout_capture(|| run_wrapped(&query, adapter, &render, None, &per_query));
let (status, body) = match result {
Ok(()) => (0u8, output),
Err(e) => (1u8, format!("{e:#}").into_bytes()),
};
let _ = conn.write_all(format!("R {} {}\n", body.len(), status).as_bytes());
let _ = conn.write_all(&body);
let _ = conn.flush();
}
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
if idle.elapsed().as_secs() >= ttl {
break;
}
std::thread::sleep(std::time::Duration::from_millis(200));
}
Err(e) => {
eprintln!("resident session accept error: {e}");
std::thread::sleep(std::time::Duration::from_millis(200));
}
}
}
let _ = std::fs::remove_file(sock);
Ok(())
}
#[cfg(unix)]
fn with_stdout_capture<R>(f: impl FnOnce() -> R) -> (R, Vec<u8>) {
use std::io::{Read as _, Seek as _, Write as _};
use std::os::fd::AsRawFd as _;
let _ = std::io::stdout().flush();
let mut tmp = match tempfile_in_temp() {
Ok(t) => t,
Err(_) => return (f(), Vec::new()),
};
let saved = unsafe { libc::dup(1) };
if saved < 0 {
return (f(), Vec::new());
}
unsafe { libc::dup2(tmp.as_raw_fd(), 1) };
let r = f();
let _ = std::io::stdout().flush();
unsafe {
libc::dup2(saved, 1);
libc::close(saved);
}
let mut out = Vec::new();
let _ = tmp.seek(std::io::SeekFrom::Start(0));
let _ = tmp.read_to_end(&mut out);
(r, out)
}
#[cfg(unix)]
fn tempfile_in_temp() -> std::io::Result<std::fs::File> {
let path = std::env::temp_dir().join(format!(
"quarb-capture-{}-{:x}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.subsec_nanos())
.unwrap_or(0)
));
let f = std::fs::OpenOptions::new()
.create_new(true)
.read(true)
.write(true)
.open(&path)?;
let _ = std::fs::remove_file(&path);
Ok(f)
}
fn preflight(cli: &Cli, query: &str, oc: &OutputCtx) -> anyhow::Result<bool> {
let cap = oc.budget.as_ref().map_or(quarb::Cap::NONE, |b| b.cap());
let mut catalogs = Vec::new();
for (i, p) in cli.paths.iter().enumerate() {
let (name, target) = match split_alias(p) {
Some(alias) => alias,
None => (
p.file_stem()
.map(|s| s.to_string_lossy().into_owned())
.unwrap_or_else(|| format!("doc{i}")),
p.clone(),
),
};
let t = target.to_string_lossy().into_owned();
if quarb_open::catalog(&t, cap).is_none() {
return Ok(false);
}
let lookup: quarb::CatalogLookup = Box::new(move |ask| quarb_open::catalog_for(&t, ask));
catalogs.push((name, target.display().to_string(), lookup));
}
if catalogs.is_empty() {
return Ok(false);
}
let run = oc.run_opts();
let defs = quarb::Defs::default();
let plan = if catalogs.len() >= 2 || cli.paths.iter().any(|p| split_alias(p).is_some()) {
let mut m = MountAdapter::new(Vec::new());
for (name, target, lookup) in catalogs {
m.declare(Latent {
name,
target: Some(target),
opener: Box::new(|| Err("planned, not opened".into())),
catalog: Some(lookup),
})
.map_err(|e| anyhow::anyhow!("mount: {e}"))?;
}
let q = quarb::Query::parse_with(query, &defs, &m, &run)?;
match &oc.budget {
Some(b) => q.plan_within(&m, b),
None => q.plan(&m),
}
} else {
let a = quarb::Unopened(catalogs.pop().expect("one input").2);
let q = quarb::Query::parse_with(query, &defs, &a, &run)?;
match &oc.budget {
Some(b) => q.plan_within(&a, b),
None => q.plan(&a),
}
};
if oc.plan {
print!("{}", plan.to_kaiv(query));
return Ok(true);
}
if let Some(budget) = &oc.budget
&& let Some(why) = budget.exceeded_by(&plan)
{
return Err(quarb::QuarbError::refused(why).into());
}
Ok(false)
}
fn take_environments(cli: &mut Cli) -> anyhow::Result<Vec<quarb::Environment>> {
let mut environments: Vec<quarb::Environment> = Vec::new();
let mut kept = Vec::new();
for p in std::mem::take(&mut cli.paths) {
match quarb::Environment::parse_mount(&p.to_string_lossy()) {
None => kept.push(p),
Some(Ok(e)) => {
if environments.iter().any(|x| x.alias == e.alias) {
anyhow::bail!(
"two environments mounted as '{}'; give each its own alias",
e.alias
);
}
environments.push(e);
}
Some(Err(e)) => anyhow::bail!("{e}"),
}
}
cli.paths = kept;
if cli.allow_shell && !environments.iter().any(|e| e.alias == "sh") {
environments.push(quarb::Environment {
alias: "sh".to_string(),
kind: quarb::EnvKind::Shell,
target: String::new(),
});
}
if environments.iter().any(|e| e.kind == quarb::EnvKind::Shell) {
cli.allow_shell = true;
}
Ok(environments)
}
fn execute(cli: &Cli, query: &str, oc: &OutputCtx) -> anyhow::Result<()> {
let ctx = &SourceCtx::default();
let (query, encoding) = if cli.expand || cli.expand_1 {
(query.to_string(), None)
} else {
take_output_encoding(query)?
};
let query = &query;
let oc = &OutputCtx {
encoding,
..oc.clone()
};
if cli.refs.is_some() {
let consumes = |p: &PathBuf| {
let target = split_alias(p).map(|(_, t)| t).unwrap_or_else(|| p.clone());
is_sqlite(&target)
|| target
.to_str()
.is_some_and(|s| s.starts_with("firebase://"))
};
if !cli.paths.iter().any(consumes) {
eprintln!(
"qua: --refs: no target consumes a declared-references document \
(SQLite databases and firebase:// do); ignoring it"
);
}
}
if (oc.plan || oc.budget.is_some()) && preflight(cli, query, oc)? {
return Ok(());
}
if cli.paths.len() >= 2 || cli.paths.iter().any(|p| split_alias(p).is_some()) {
let mut mounts: Vec<Mount> = Vec::new();
let mut renders: Vec<Box<dyn Fn(NodeId) -> String>> = Vec::new();
for (i, p) in cli.paths.iter().enumerate() {
let (name, target) = match split_alias(p) {
Some(alias) => alias,
None => (
p.file_stem()
.map(|s| s.to_string_lossy().into_owned())
.unwrap_or_else(|| format!("doc{i}")),
p.clone(),
),
};
if mounts.iter().any(|m| m.name == name) {
anyhow::bail!(
"input '{}' mounts as '{name}', colliding with an earlier input of the \
same name; give one an explicit alias (NAME=TARGET)",
p.display()
);
}
let Mounted {
adapter, render, ..
} = open_mount(&target, &OpenOpts::from(cli))?;
mounts.push(Mount {
name,
target: Some(target.display().to_string()),
adapter,
});
renders.push(render);
}
let sources = cli
.paths
.iter()
.map(|p| p.display().to_string())
.collect::<Vec<_>>()
.join(", ");
let adapter = MountAdapter::try_new(mounts).map_err(|e| anyhow::anyhow!("mount: {e}"))?;
return run(
query,
&adapter,
|n| match adapter.decode(n) {
None => "/".to_string(),
Some((m, inner)) => {
format!("/{}{}", adapter.mount_name(m), renders[m](inner))
}
},
cli.kaiv.then_some(sources.as_str()),
ctx,
oc,
);
}
let path = match cli.paths.first() {
Some(p) => Some(take_encoding_option(p, ctx)?),
None => None,
};
if let Some(s) = path.as_ref().and_then(|p| p.to_str())
&& is_prefixed_target(s)
&& !is_relational_target(s)
&& !s.starts_with("web:")
{
let src = s.to_string();
let m = open_mount(cli.paths.first().expect("a target"), &OpenOpts::from(cli))?;
return run(
query,
&m.adapter,
|n| (m.render)(n),
cli.kaiv.then_some(src.as_str()),
ctx,
oc,
);
}
if let Some(p) = &path
&& let Some(rest) = p.to_str().and_then(|s| s.strip_prefix("web:"))
&& !rest.is_empty()
{
let src = rest.to_string();
return match web_level(rest, ctx)? {
WebSite::Memory(adapter) => run(
query,
&adapter,
|n| adapter.locator(n),
cli.kaiv.then_some(src.as_str()),
ctx,
oc,
),
WebSite::Sqlite(store) => run_web_store(cli, query, store, &src, ctx, oc),
WebSite::Postgres(store) => run_web_store(cli, query, store, &src, ctx, oc),
};
}
#[cfg(feature = "kuzu")]
if let Some(s) = path.as_ref().and_then(|p| p.to_str())
&& s.starts_with("mssql://")
{
if let Some(plan) = pushdown_plan(cli, query, Some(quarb_sql::Dialect::Mssql)) {
match quarb_mssql::raw_query(
s,
&plan.sql,
plan.order_table.as_deref(),
plan.join_left
.as_ref()
.map(|(t, c)| (t.as_str(), c.as_slice())),
) {
Ok((cols, rows)) => {
print_raw(&cols, rows, oc)?;
return Ok(());
}
Err(e) => {
if cli.explain {
eprintln!("pushdown: {}", plan.sql);
eprintln!("pushdown: plan not executed ({e}); scanning");
}
}
}
}
let adapter = match partial_plan_with(cli, query, Some(quarb_sql::Dialect::Mssql), true) {
Some(pl) => {
let a = MssqlAdapter::connect_filtered(s, &pl.table, &pl.where_sql)
.context("connecting to SQL Server")?;
match a.prefetch(&pl.table) {
Ok(()) => a,
Err(e) => {
if cli.explain {
eprintln!("partial pushdown: prefilter rejected ({e}); scanning");
}
MssqlAdapter::connect(s).context("connecting to SQL Server")?
}
}
}
None => MssqlAdapter::connect(s).context("connecting to SQL Server")?,
};
return run_relational(
adapter,
cli.no_graft,
query,
|a, n| a.locator(n),
cli.kaiv.then_some(s),
ctx,
oc,
);
}
if let Some(s) = path.as_ref().and_then(|p| p.to_str())
&& s.starts_with("oracle://")
{
if let Some(plan) = pushdown_plan(cli, query, Some(quarb_sql::Dialect::Oracle)) {
match quarb_oracle::raw_query(
s,
&plan.sql,
plan.order_table.as_deref(),
plan.join_left
.as_ref()
.map(|(t, c)| (t.as_str(), c.as_slice())),
) {
Ok((cols, rows)) => {
print_raw(&cols, rows, oc)?;
return Ok(());
}
Err(e) => {
if cli.explain {
eprintln!("pushdown: {}", plan.sql);
eprintln!("pushdown: plan not executed ({e}); scanning");
}
}
}
}
let adapter = match partial_plan_with(cli, query, Some(quarb_sql::Dialect::Oracle), true) {
Some(pl) => {
let a = OracleAdapter::connect_filtered(s, &pl.table, &pl.where_sql)
.context("connecting to Oracle")?;
match a.prefetch(&pl.table) {
Ok(()) => a,
Err(e) => {
if cli.explain {
eprintln!("partial pushdown: prefilter rejected ({e}); scanning");
}
OracleAdapter::connect(s).context("connecting to Oracle")?
}
}
}
None => OracleAdapter::connect(s).context("connecting to Oracle")?,
};
return run_relational(
adapter,
cli.no_graft,
query,
|a, n| a.locator(n),
cli.kaiv.then_some(s),
ctx,
oc,
);
}
if let Some(s) = path.as_ref().and_then(|p| p.to_str())
&& s.starts_with("bigquery://")
{
if let Some(plan) = pushdown_plan(cli, query, None) {
match quarb_bigquery::raw_query(
s,
&plan.sql,
plan.order_table.as_deref(),
plan.join_left
.as_ref()
.map(|(t, c)| (t.as_str(), c.as_slice())),
) {
Ok((cols, rows)) => {
print_raw(&cols, rows, oc)?;
return Ok(());
}
Err(e) => {
if cli.explain {
eprintln!("pushdown: {}", plan.sql);
eprintln!("pushdown: plan not executed ({e}); scanning");
}
}
}
}
let adapter = match partial_plan(cli, query, None) {
Some(pl) => {
let a = BigqueryAdapter::connect_filtered(s, &pl.table, &pl.where_sql)
.context("connecting to BigQuery")?;
match a.prefetch(&pl.table) {
Ok(()) => a,
Err(e) => {
if cli.explain {
eprintln!("partial pushdown: prefilter rejected ({e}); scanning");
}
BigqueryAdapter::connect(s).context("connecting to BigQuery")?
}
}
}
None => BigqueryAdapter::connect(s).context("connecting to BigQuery")?,
};
return run_relational(
adapter,
cli.no_graft,
query,
|a, n| a.locator(n),
cli.kaiv.then_some(s),
ctx,
oc,
);
}
if let Some(s) = path.as_ref().and_then(|p| p.to_str())
&& s.starts_with("athena:")
{
if let Some(plan) = pushdown_plan(cli, query, None) {
match quarb_athena::raw_query(
s,
&plan.sql,
plan.order_table.as_deref(),
plan.join_left
.as_ref()
.map(|(t, c)| (t.as_str(), c.as_slice())),
) {
Ok((cols, rows)) => {
print_raw(&cols, rows, oc)?;
return Ok(());
}
Err(e) => {
if cli.explain {
eprintln!("pushdown: {}", plan.sql);
eprintln!("pushdown: plan not executed ({e}); scanning");
}
}
}
}
let adapter = match partial_plan(cli, query, None) {
Some(pl) => {
let a = AthenaAdapter::connect_filtered(s, &pl.table, &pl.where_sql)
.context("connecting to Athena")?;
match a.prefetch(&pl.table) {
Ok(()) => a,
Err(e) => {
if cli.explain {
eprintln!("partial pushdown: prefilter rejected ({e}); scanning");
}
AthenaAdapter::connect(s).context("connecting to Athena")?
}
}
}
None => AthenaAdapter::connect(s).context("connecting to Athena")?,
};
return run_relational(
adapter,
cli.no_graft,
query,
|a, n| a.locator(n),
cli.kaiv.then_some(s),
ctx,
oc,
);
}
if let Some(s) = path.as_ref().and_then(|p| p.to_str())
&& s.starts_with("mysql://")
{
if let Some(plan) =
pushdown_plan(cli, query, Some(quarb_sql::Dialect::MySql)).filter(|plan| {
column_types_hold(cli, "pushdown", &plan.column_types, || {
let tables: Vec<&str> = plan.tables.iter().map(String::as_str).collect();
Ok(quarb_mysql::column_catalog(s, &tables)?)
})
})
{
match quarb_mysql::raw_query(
s,
&plan.sql,
plan.order_table.as_deref(),
plan.join_left
.as_ref()
.map(|(t, c)| (t.as_str(), c.as_slice())),
) {
Ok((cols, rows)) => {
print_raw(&cols, rows, oc)?;
return Ok(());
}
Err(e) => {
if cli.explain {
eprintln!("pushdown: {}", plan.sql);
eprintln!("pushdown: plan not executed ({e}); scanning");
}
}
}
}
let adapter = match partial_plan_with(cli, query, Some(quarb_sql::Dialect::MySql), true)
.filter(|pl| {
column_types_hold(cli, "partial pushdown", &pl.column_types, || {
Ok(quarb_mysql::column_catalog(s, &[pl.table.as_str()])?)
})
}) {
Some(pl) => {
let a = MysqlAdapter::connect_filtered(s, &pl.table, &pl.where_sql)
.context("connecting to MySQL")?;
match a.prefetch(&pl.table) {
Ok(()) => a,
Err(e) => {
if cli.explain {
eprintln!("partial pushdown: prefilter rejected ({e}); scanning");
}
MysqlAdapter::connect(s).context("connecting to MySQL")?
}
}
}
None => MysqlAdapter::connect(s).context("connecting to MySQL")?,
};
return run_relational(
adapter,
cli.no_graft,
query,
|a, n| a.locator(n),
cli.kaiv.then_some(s),
ctx,
oc,
);
}
if let Some(s) = path.as_ref().and_then(|p| p.to_str())
&& is_pg_config(s)
{
if let Some(plan) = pushdown_plan(cli, query, Some(quarb_sql::Dialect::Postgres)) {
match quarb_postgres::raw_query(
s,
&plan.sql,
plan.order_table.as_deref(),
plan.join_left
.as_ref()
.map(|(t, c)| (t.as_str(), c.as_slice())),
) {
Ok((cols, rows)) => {
print_raw(&cols, rows, oc)?;
return Ok(());
}
Err(e) => {
if cli.explain {
eprintln!("pushdown: {}", plan.sql);
eprintln!("pushdown: plan not executed ({e}); scanning");
}
}
}
}
let adapter = match partial_plan_with(cli, query, Some(quarb_sql::Dialect::Postgres), true)
{
Some(pl) => {
let a = PostgresAdapter::connect_filtered(s, &pl.table, &pl.where_sql)
.context("connecting to PostgreSQL")?;
match a.prefetch(&pl.table) {
Ok(()) => a,
Err(e) => {
if cli.explain {
eprintln!("partial pushdown: prefilter rejected ({e}); scanning");
}
PostgresAdapter::connect(s).context("connecting to PostgreSQL")?
}
}
}
None => PostgresAdapter::connect(s).context("connecting to PostgreSQL")?,
};
return run_relational(
adapter,
cli.no_graft,
query,
|a, n| a.locator(n),
cli.kaiv.then_some(s),
ctx,
oc,
);
}
if let Some(p) = &path
&& p.extension()
.and_then(|e| e.to_str())
.is_some_and(|e| e.eq_ignore_ascii_case("duckdb") || e.eq_ignore_ascii_case("ddb"))
{
if let Some(plan) = pushdown_plan(cli, query, None) {
match quarb_duckdb::raw_query(
p,
&plan.sql,
plan.order_table.as_deref(),
plan.join_left
.as_ref()
.map(|(t, c)| (t.as_str(), c.as_slice())),
) {
Ok((cols, rows)) => {
print_raw(&cols, rows, oc)?;
return Ok(());
}
Err(e) => {
if cli.explain {
eprintln!("pushdown: {}", plan.sql);
eprintln!("pushdown: plan not executed ({e}); scanning");
}
}
}
}
let adapter = DuckdbAdapter::open(p).context("opening DuckDB database")?;
let src = p.display().to_string();
return run_relational(
adapter,
cli.no_graft,
query,
|a, n| a.locator(n),
cli.kaiv.then_some(src.as_str()),
ctx,
oc,
);
}
if let Some(p) = &path
&& is_sqlite(p)
{
if let Some(plan) =
pushdown_plan(cli, query, Some(quarb_sql::Dialect::Sqlite)).filter(|plan| {
column_types_hold(cli, "pushdown", &plan.column_types, || {
let tables: Vec<&str> = plan.tables.iter().map(String::as_str).collect();
Ok(quarb_sqlite::column_catalog(p, &tables)?)
})
})
{
match quarb_sqlite::raw_query(
p,
&plan.sql,
plan.order_table.as_deref(),
plan.join_left
.as_ref()
.map(|(t, c)| (t.as_str(), c.as_slice())),
) {
Ok((cols, rows)) => {
print_raw(&cols, rows, oc)?;
return Ok(());
}
Err(e) => {
if cli.explain {
eprintln!("pushdown: {}", plan.sql);
eprintln!("pushdown: plan not executed ({e}); scanning");
}
}
}
}
let refs = relational_refs(&cli.refs)?;
let adapter = match partial_plan_with(cli, query, Some(quarb_sql::Dialect::Sqlite), true)
.filter(|pl| {
column_types_hold(cli, "partial pushdown", &pl.column_types, || {
Ok(quarb_sqlite::column_catalog(p, &[pl.table.as_str()])?)
})
}) {
Some(pl) => {
let a = SqliteAdapter::open_filtered_with_refs(p, &pl.table, &pl.where_sql, &refs)
.context("opening SQLite database")?;
match a.prefetch(&pl.table) {
Ok(()) => a,
Err(e) => {
if cli.explain {
eprintln!("partial pushdown: prefilter rejected ({e}); scanning");
}
SqliteAdapter::open_with_refs(p, &refs)
.context("opening SQLite database")?
}
}
}
None => SqliteAdapter::open_with_refs(p, &refs).context("opening SQLite database")?,
};
let src = p.display().to_string();
return run_relational(
adapter,
cli.no_graft,
query,
|a, n| a.locator(n),
cli.kaiv.then_some(src.as_str()),
ctx,
oc,
);
}
if let Some(p) = &path {
let src = p.display().to_string();
let m = open_mount(cli.paths.first().expect("a target"), &OpenOpts::from(cli))?;
return run(
query,
&m.adapter,
|n| (m.render)(n),
cli.kaiv.then_some(src.as_str()),
ctx,
oc,
);
}
let text = match &path {
Some(_) => unreachable!("a target opened above"),
None if std::io::stdin().is_terminal() => {
if !quarb::is_calculator(&cli.query) {
quarb::expand(&cli.query, &quarb::Defs::default()).context("parsing the query")?;
anyhow::bail!(
"no input: give a directory, a file, or pipe a document to \
stdin — an expression head '= expr' runs without one; a \
--model file opens its own sources with `mount NAME: target;`"
);
}
"{}".to_owned()
}
None => read_prose(Path::new("-"), ctx)?,
};
let path: Option<&Path> = None;
let text = match text.strip_prefix('\u{feff}') {
Some(rest) => rest.to_owned(),
None => text,
};
let source = path.map_or_else(|| "stdin".to_string(), |p| p.display().to_string());
let kaiv = cli.kaiv.then_some(source.as_str());
if is_quarb(path) {
let adapter = quarb::reflect::QueryArbor::parse(&text).context("parsing Quarb query")?;
return run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc);
}
if let Some(delim) = csv_delimiter(path) {
let adapter = CsvAdapter::parse_declared(&text, delim, ctx.csv()).context("parsing CSV")?;
return run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc);
}
if let Some(ext) = path.and_then(|p| p.extension()).and_then(|e| e.to_str()) {
let ext = ext.to_ascii_lowercase();
let ext = ext.as_str();
if matches!(ext, "yaml" | "yml") {
let adapter = quarb_yaml::parse(&text).context("parsing YAML")?;
return run(query, &adapter, |n| adapter.pointer(n), kaiv, ctx, oc);
}
if ext == "toml" {
let adapter = quarb_toml::parse(&text).context("parsing TOML")?;
return run(query, &adapter, |n| adapter.pointer(n), kaiv, ctx, oc);
}
if matches!(ext, "md" | "markdown") {
let adapter = quarb_markdown::parse(&text);
return run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc);
}
if ext == "txt" {
let adapter = quarb_text::TextModel::parse_plain(&text);
return run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc);
}
if matches!(ext, "conllu" | "conllup") {
let adapter = quarb_text::TextModel::parse_conllu_text(&text)
.map_err(|e| anyhow::anyhow!("reading CoNLL-U: {e}"))?;
return run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc);
}
if matches!(ext, "jsonl" | "ndjson") {
let adapter = JsonAdapter::parse_lines(&text).context("parsing JSONL")?;
return run(query, &adapter, |n| adapter.pointer(n), kaiv, ctx, oc);
}
if matches!(ext, "daiv" | "kaiv" | "raiv") {
let dir = path.and_then(|p| p.parent());
let adapter = parse_kaiv_ext(ext, &text, dir)?;
return run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc);
}
if matches!(ext, "atd" | "atk" | "usfm" | "sfm") {
let dir = path.and_then(|p| p.parent()).unwrap_or(Path::new("."));
let adapter = AtrepAdapter::parse_str(&text, dir).context("parsing atrep document")?;
return run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc);
}
}
if is_atrep(&text) {
let dir = path
.and_then(|p| p.parent())
.unwrap_or_else(|| Path::new("."));
let adapter = AtrepAdapter::parse_str(&text, dir).context("parsing atrep document")?;
return run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc);
}
if is_xml(path, &text) {
let adapter = XmlAdapter::parse(&text).context("parsing XML")?;
run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc)
} else if is_html(path, &text) {
let adapter = HtmlAdapter::parse(&text);
run(query, &adapter, |n| adapter.locator(n), kaiv, ctx, oc)
} else {
let adapter = match JsonAdapter::parse(&text) {
Ok(a) => a,
Err(e) => match JsonAdapter::parse_lines(&text) {
Ok(a) => a,
Err(_) => return Err(e).context("parsing JSON"),
},
};
run(query, &adapter, |n| adapter.pointer(n), kaiv, ctx, oc)
}
}
fn pushdown_applies(cli: &Cli) -> bool {
!cli.no_pushdown
&& !cli.plan
&& cli.budget.is_none()
&& !cli.kaiv
&& !cli.expand
&& !cli.expand_1
&& cli.save.is_none()
&& !cli.resident
&& !cli.resident_serve
}
fn partial_plan(
cli: &Cli,
query: &str,
dialect: Option<quarb_sql::Dialect>,
) -> Option<quarb_sql::Partial> {
partial_plan_with(cli, query, dialect, false)
}
fn column_types_hold(
cli: &Cli,
what: &str,
types: &quarb_sql::ColumnTypes,
catalog: impl FnOnce() -> anyhow::Result<Vec<(String, bool)>>,
) -> bool {
if types.is_empty() {
return true;
}
let unmet = match catalog() {
Ok(c) => types.unmet(&c),
Err(_) => types.numeric.iter().chain(&types.text).cloned().collect(),
};
if !unmet.is_empty() && cli.explain {
eprintln!(
"{what}: not used (the comparison on {} would follow the column's type where the scan follows the value); scanning",
unmet.join(", ")
);
}
unmet.is_empty()
}
fn partial_plan_with(
cli: &Cli,
query: &str,
dialect: Option<quarb_sql::Dialect>,
keyed_hops: bool,
) -> Option<quarb_sql::Partial> {
if !pushdown_applies(cli) {
return None;
}
let plan = if keyed_hops {
quarb_sql::partial_pushdown_keyed(query, dialect)
} else {
quarb_sql::partial_pushdown_explained(query, dialect)
};
match plan {
Ok(p) => {
if cli.explain {
eprintln!(
"partial pushdown: WHERE {} on {}; the rest scans the filtered set",
p.where_sql, p.table
);
}
Some(p)
}
Err(e) => {
if cli.explain {
eprintln!("partial pushdown refused: {e}; scanning");
}
None
}
}
}
fn pushdown_plan(
cli: &Cli,
query: &str,
dialect: Option<quarb_sql::Dialect>,
) -> Option<quarb_sql::Pushdown> {
if !pushdown_applies(cli) {
if cli.explain {
eprintln!("pushdown: disabled; scanning");
}
return None;
}
match quarb_sql::pushdown_explained(query, dialect) {
Ok(plan) => {
if cli.explain {
}
Some(plan)
}
Err(e) => {
if cli.explain {
eprintln!("pushdown refused: {e}; scanning");
}
None
}
}
}
fn print_raw(cols: &[String], rows: Vec<Vec<Value>>, oc: &OutputCtx) -> anyhow::Result<()> {
if oc.explain
&& let Some(sql) = quarb_relational::take_executed()
{
eprintln!("pushdown: {sql}");
}
let values: Vec<Value> = rows
.into_iter()
.flat_map(|row| {
if cols.len() <= 1 {
row
} else {
vec![Value::Record(cols.iter().cloned().zip(row).collect())]
}
})
.collect();
let mut out = OutSink::new(oc.encoding.clone());
emit_values(&mut out, values, oc.output, oc)?;
out.finish()
}
struct OutSink {
buf: Vec<u8>,
encoding: Option<(String, quarb::source::Escape)>,
}
impl OutSink {
fn new(encoding: OutputEncoding) -> Self {
OutSink {
buf: Vec::new(),
encoding,
}
}
fn finish(self) -> anyhow::Result<()> {
use std::io::Write as _;
let stdout = std::io::stdout();
let mut out = stdout.lock();
match self.encoding {
Some((label, escape)) => {
let text = String::from_utf8(self.buf).context("the output is not UTF-8")?;
let bytes = quarb::source::encode_with(&text, &label, escape)
.map_err(|e| anyhow::anyhow!("encoding({label}): {e}"))?;
out.write_all(&bytes)?;
}
None => out.write_all(&self.buf)?,
}
out.flush()?;
Ok(())
}
}
impl std::io::Write for OutSink {
fn write(&mut self, data: &[u8]) -> std::io::Result<usize> {
self.buf.extend_from_slice(data);
Ok(data.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
type OutputEncoding = Option<(String, quarb::source::Escape)>;
fn take_output_encoding(query: &str) -> anyhow::Result<(String, OutputEncoding)> {
let trimmed = query.trim_end();
let Some(head) = trimmed.strip_suffix(')') else {
return Ok((query.to_string(), None));
};
let Some(open) = head.rfind('(') else {
return Ok((query.to_string(), None));
};
let (before, label) = head.split_at(open);
let before = before.trim_end();
let Some(rest) = before.strip_suffix("encoding") else {
return Ok((query.to_string(), None));
};
let rest = rest.trim_end();
let Some(rest) = rest.strip_suffix("@|") else {
return Ok((query.to_string(), None));
};
let label = label[1..].trim().trim_matches('"');
quarb::source::encoding_for_label(label).ok_or_else(|| {
anyhow::anyhow!(
"encoding: unknown label {label:?} (labels are the Encoding Standard's: utf-8, \
windows-1251, koi8-r, iso-8859-2, …)"
)
})?;
let rest = rest.trim_end();
let escape = match rest.rsplit("@|").next().map(str::trim) {
Some(command)
if command == "json" || command == "jsonl" || command.starts_with("json(") =>
{
quarb::source::Escape::Json
}
_ => quarb::source::Escape::Refuse,
};
Ok((rest.to_string(), Some((label.to_string(), escape))))
}
fn run_web_store<S: quarb_web::db::SqlStore + 'static>(
cli: &Cli,
query: &str,
store: S,
src: &str,
ctx: &SourceCtx,
oc: &OutputCtx,
) -> anyhow::Result<()> {
use quarb_web::plan::{Plan, Rung};
let model = oc.model.is_some();
let plan = if pushdown_applies(cli) {
quarb_web::plan::plan(query, store.dialect(), model)
} else {
Plan {
rung: Rung::Scan,
where_sql: String::new(),
params: Vec::new(),
reason: "disabled".into(),
unverified: Vec::new(),
limit: None,
}
};
if cli.explain {
eprintln!("web: {} — {}", web_rung_name(&plan.rung), plan.reason);
if plan.rung != Rung::Scan {
eprintln!("web: WHERE {}", plan.where_sql);
if let Some(e) = store.estimate_where(&plan.where_sql, &plan.params) {
eprintln!("web: plan {e}");
}
}
}
match plan.rung {
Rung::Full => {
let n = store.count_where(&plan.where_sql, &plan.params);
println!("{n}");
Ok(())
}
Rung::Prefilter => {
let keys = match &plan.limit {
Some(l) => store.keys_where_top(
&plan.where_sql,
&plan.params,
&l.column,
l.descending,
l.n,
),
None => store.keys_where(&plan.where_sql, &plan.params),
};
if cli.explain {
eprintln!("web: {} candidate(s)", keys.len());
}
let adapter = quarb_web::WebAdapter::new(store).with_scope(keys);
run(
query,
&adapter,
|n| adapter.locator(n),
cli.kaiv.then_some(src),
ctx,
oc,
)
}
Rung::Scan => {
let adapter = quarb_web::WebAdapter::new(store);
run(
query,
&adapter,
|n| adapter.locator(n),
cli.kaiv.then_some(src),
ctx,
oc,
)
}
}
}
fn web_rung_name(r: &quarb_web::plan::Rung) -> &'static str {
match r {
quarb_web::plan::Rung::Full => "full",
quarb_web::plan::Rung::Prefilter => "prefilter",
quarb_web::plan::Rung::Scan => "scan",
}
}
fn run_relational<A: AstAdapter>(
inner: A,
no_graft: bool,
query: &str,
outer_loc: impl Fn(&A, NodeId) -> String,
kaiv_source: Option<&str>,
ctx: &SourceCtx,
oc: &OutputCtx,
) -> anyhow::Result<()> {
if no_graft {
return run(
query,
&inner,
|n| outer_loc(&inner, n),
kaiv_source,
ctx,
oc,
);
}
let adapter = ComposeAdapter::new(inner)
.with_decoding(current_decoding(ctx)?)
.with_csv_options(ctx.csv());
run(
query,
&adapter,
|n| adapter.locator(n, |o| outer_loc(adapter.outer(), o)),
kaiv_source,
ctx,
oc,
)
}
fn run<A: AstAdapter>(
query: &str,
adapter: &A,
render: impl Fn(NodeId) -> String,
kaiv_source: Option<&str>,
ctx: &SourceCtx,
oc: &OutputCtx,
) -> anyhow::Result<()> {
if let Some((encoding, replaced)) = ctx.take_record() {
let sourced = quarb::source::Sourced {
inner: quarb_model::Borrowed(adapter),
encoding,
replaced,
};
return run_modelled(query, &sourced, render, kaiv_source, oc);
}
run_modelled(query, adapter, render, kaiv_source, oc)
}
fn run_modelled<A: AstAdapter>(
query: &str,
adapter: &A,
render: impl Fn(NodeId) -> String,
kaiv_source: Option<&str>,
oc: &OutputCtx,
) -> anyhow::Result<()> {
if let Some(model) = oc.model.clone() {
let enriched = quarb_model::ModelAdapter::new(quarb_model::Borrowed(adapter), model);
enriched
.check()
.map_err(|e| anyhow::anyhow!("reading the model: {e}"))?;
let base_render = &render;
let model_render = |n: NodeId| enriched.locator(n, base_render);
return run_dispatch(query, &enriched, model_render, kaiv_source, oc);
}
run_dispatch(query, adapter, render, kaiv_source, oc)
}
fn run_dispatch<A: AstAdapter>(
query: &str,
adapter: &A,
render: impl Fn(NodeId) -> String,
kaiv_source: Option<&str>,
oc: &OutputCtx,
) -> anyhow::Result<()> {
#[cfg(unix)]
if let Some((sock, ttl, pinned)) = oc.resident.clone() {
return resident_serve_loop(adapter, render, &sock, ttl, pinned, oc);
}
run_wrapped(query, adapter, &render, kaiv_source, oc)
}
impl OutputCtx {
fn run_opts(&self) -> quarb::Run {
quarb::Run {
now: Some(self.now),
allow_shell: self.allow_shell,
environments: self.environments.clone(),
quantifier_bound: self.quant_bound,
reproducible: self.reproducible,
collation: self.locale.map(|(l, _)| l.tag.to_string()),
budget: self.budget.clone(),
plan: self.plan,
..quarb::Run::ambient()
}
}
fn explain_refs(&self, refs: &[String]) {
if self.explain {
for r in refs {
eprintln!(
"unresolved external reference: {r} (mount it, or declare a tree for it)"
);
}
}
}
}
fn parse_budget(text: &str) -> anyhow::Result<quarb::plan::Budget> {
quarb::plan::Budget::parse(text).map_err(|e| anyhow::anyhow!("--{e}"))
}
fn run_wrapped<A: AstAdapter>(
query: &str,
adapter: &A,
render: impl Fn(NodeId) -> String,
kaiv_source: Option<&str>,
oc: &OutputCtx,
) -> anyhow::Result<()> {
let run = oc.run_opts();
if oc.expand_1 {
for t in quarb::expand_first_with(query, &quarb::Defs::default(), adapter, &run)
.context("expanding the query")?
{
println!("{t}");
}
return Ok(());
}
if oc.expand {
println!(
"{}",
quarb::expand_with(query, &quarb::Defs::default(), adapter, &run)
.context("expanding the query")?
);
return Ok(());
}
let parsed = quarb::Query::parse_with(query, &quarb::Defs::default(), adapter, &run)?;
if oc.plan {
print!("{}", parsed.plan(adapter).to_kaiv(query));
return Ok(());
}
if let Some(source) = kaiv_source {
let rows = parsed.run_traced_prov(adapter, &run)?;
print!(
"{}",
emit_kaiv(
&rows,
source,
&render,
|n| quarb::resolved_provenance(adapter, n),
oc.kaiv_origins,
)?
);
return Ok(());
}
let save = oc.save.clone();
if let Some((path, table)) = save {
let out = parsed.run(adapter, &run)?;
oc.explain_refs(&out.refs);
let values = match out.result {
QueryResult::Values(vs) => vs,
QueryResult::Nodes(ns) => ns.into_iter().map(|n| Value::Str(render(n))).collect(),
};
let n = values.len();
save_result(&path, &table, values)?;
eprintln!("saved {n} row(s) to {}", path.display());
return Ok(());
}
let out = OutSink::new(oc.encoding.clone());
let mode = oc.output;
let outcome = match oc.locale.map(|(l, _)| l) {
Some(_) if mode == Output::Quarb => parsed.for_rendering().run(adapter, &run)?,
_ => parsed.run(adapter, &run)?,
};
let quarb::Outcome {
result,
closing,
refs,
plan: _,
} = outcome;
let printed = print_result(out, result, closing, render, oc);
oc.explain_refs(&refs);
printed
}
fn print_result(
mut out: OutSink,
result: QueryResult,
closing: Option<quarb::ClosingTable>,
render: impl Fn(NodeId) -> String,
oc: &OutputCtx,
) -> anyhow::Result<()> {
use std::io::Write as _;
let mode = oc.output;
if let (Some(t), Some(l)) = (&closing, oc.locale.map(|(l, _)| l)) {
let rows = match result {
QueryResult::Values(values) => values,
QueryResult::Nodes(nodes) => nodes.into_iter().map(|n| Value::Str(render(n))).collect(),
};
let columns = (!t.columns.is_empty()).then_some(t.columns.as_slice());
writeln!(
out,
"{}",
quarb::tabulate_in(
&rows,
columns,
t.csv,
t.delimiter.unwrap_or(l.csv_delimiter),
&|v| l.cell(v)
)
)?;
return out.finish();
}
let values = match result {
QueryResult::Nodes(nodes) => {
if mode == Output::Quarb {
for node in nodes {
writeln!(out, "{}", render(node))?;
}
return out.finish();
}
nodes.into_iter().map(|n| Value::Str(render(n))).collect()
}
QueryResult::Values(values) => values,
};
emit_values(&mut out, values, mode, oc)?;
out.finish()
}
fn emit_values(
out: &mut impl std::io::Write,
values: Vec<Value>,
mode: Output,
oc: &OutputCtx,
) -> anyhow::Result<()> {
match mode {
Output::Quarb => {
let locale = oc.locale.map(|(l, _)| l);
for value in values {
match &locale {
Some(l) => writeln!(out, "{}", l.display(&value))?,
None => writeln!(out, "{}", value.display_form())?,
}
}
}
Output::Json => {
let items: Vec<String> = values.iter().map(Value::to_json).collect();
writeln!(out, "[{}]", items.join(", "))?;
}
Output::Jsonl => {
for value in values {
writeln!(out, "{}", value.to_json())?;
}
}
Output::Table => write_table(out, &values, false, oc)?,
Output::Csv => write_table(out, &values, true, oc)?,
}
Ok(())
}
fn write_table(
out: &mut impl std::io::Write,
values: &[Value],
csv: bool,
oc: &OutputCtx,
) -> anyhow::Result<()> {
match oc.locale.map(|(l, _)| l) {
Some(l) => writeln!(
out,
"{}",
quarb::tabulate_in(values, None, csv, l.csv_delimiter, &|v| l.cell(v))
)?,
None => writeln!(out, "{}", quarb::tabulate(values, None, csv))?,
}
Ok(())
}
fn set_output_locale(
cli: &Cli,
) -> anyhow::Result<Option<(quarb::locale::OutputLocale, &'static str)>> {
let layers: [(&'static str, Option<String>); 4] = [
("--locale", cli.locale.clone()),
("QUARB_LOCALE", std::env::var("QUARB_LOCALE").ok()),
("config", config_locale()),
(
"system",
std::env::var("LC_ALL")
.ok()
.filter(|v| !v.is_empty())
.or_else(|| std::env::var("LANG").ok()),
),
];
let chosen = layers
.iter()
.find_map(|(layer, tag)| tag.as_deref().map(|t| (*layer, t.to_string())));
let Some((layer, tag)) = chosen else {
return Ok(None);
};
if quarb::locale::is_canonical(&tag) {
return Ok(None);
}
match quarb::locale::for_tag(&tag) {
Some(mut l) => {
l.grouping = cli.grouping
|| std::env::var("QUARB_GROUPING").is_ok_and(|v| v == "1" || v == "true")
|| config_output("grouping").is_some_and(|v| v.as_bool() == Some(true));
if cli.explain {
eprintln!("locale: {tag} (from {layer})");
}
return Ok(Some((l, layer)));
}
None => {
if layer == "--locale" {
anyhow::bail!(
"--locale {tag}: no rendering rules for this locale (shipped: {}; C is the canonical form)",
quarb::locale::SHIPPED.join(", ")
);
}
if cli.explain {
eprintln!("locale: {tag} (from {layer}) has no rendering rules; canonical");
}
}
}
Ok(None)
}
fn config_locale() -> Option<String> {
config_output("locale")?.as_str().map(str::to_string)
}
fn config_output(key: &str) -> Option<toml::Value> {
let path = match std::env::var_os("QUARB_CONFIG") {
Some(p) => PathBuf::from(p),
None => std::env::var_os("HOME").map(|h| PathBuf::from(h).join(".quarb/config.toml"))?,
};
let text = std::fs::read_to_string(path).ok()?;
let doc: toml::Table = toml::from_str(&text).ok()?;
doc.get("output")?.get(key).cloned()
}
fn save_result(path: &Path, table: &str, values: Vec<Value>) -> anyhow::Result<()> {
let ext = path
.extension()
.and_then(|e| e.to_str())
.unwrap_or("")
.to_ascii_lowercase();
if ext == "json" {
use std::io::Write as _;
let mut f = match std::fs::OpenOptions::new()
.create_new(true)
.write(true)
.open(path)
{
Ok(f) => f,
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => {
anyhow::bail!("{} already exists (refusing to overwrite)", path.display())
}
Err(e) => return Err(e).with_context(|| format!("creating {}", path.display())),
};
let items: Vec<String> = values.iter().map(|v| v.to_json()).collect();
f.write_all(
format!(
"[{}]
",
items.join(
",
"
)
)
.as_bytes(),
)?;
return Ok(());
}
let mut conn =
rusqlite::Connection::open(path).with_context(|| format!("opening {}", path.display()))?;
let exists: i64 = conn.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name = ?1",
[table],
|r| r.get(0),
)?;
if exists > 0 {
anyhow::bail!(
"table '{table}' already exists in {} (refusing to overwrite)",
path.display()
);
}
let mut columns: Vec<String> = Vec::new();
let all_records = values.iter().all(|v| matches!(v, Value::Record(_)));
if all_records {
for v in &values {
if let Value::Record(fields) = v {
for (k, _) in fields {
if !columns.contains(k) {
columns.push(k.clone());
}
}
}
}
}
if columns.is_empty() {
columns.push("value".to_string());
}
let ident = |name: &str| format!("\"{}\"", name.replace('"', "\"\""));
let decl: Vec<String> = columns.iter().map(|c| ident(c)).collect();
let tx = conn.transaction()?;
tx.execute(
&format!("CREATE TABLE {} ({})", ident(table), decl.join(", ")),
[],
)?;
let placeholders: Vec<String> = (1..=columns.len()).map(|i| format!("?{i}")).collect();
{
let mut stmt = tx.prepare(&format!(
"INSERT INTO {} ({}) VALUES ({})",
ident(table),
decl.join(", "),
placeholders.join(", ")
))?;
for v in values {
let row: Vec<rusqlite::types::Value> = if all_records {
let Value::Record(fields) = &v else {
unreachable!()
};
columns
.iter()
.map(|c| {
fields
.iter()
.find(|(k, _)| k == c)
.map(|(_, v)| sqlite_value(v))
.unwrap_or(rusqlite::types::Value::Null)
})
.collect()
} else {
vec![sqlite_value(&v)]
};
stmt.execute(rusqlite::params_from_iter(row))?;
}
}
tx.commit()?;
Ok(())
}
fn sqlite_value(v: &Value) -> rusqlite::types::Value {
match v {
Value::Null => rusqlite::types::Value::Null,
Value::Bool(b) => rusqlite::types::Value::Integer(*b as i64),
Value::Int(n) => rusqlite::types::Value::Integer(*n),
Value::Float(f) => rusqlite::types::Value::Real(*f),
other => rusqlite::types::Value::Text(other.to_string()),
}
}
fn emit_kaiv(
rows: &[quarb::Traced],
source: &str,
render: impl Fn(NodeId) -> String,
prov_of: impl Fn(NodeId) -> quarb::ProvenanceList,
cap: usize,
) -> anyhow::Result<String> {
use kaiv::{KaivBuilder, ProvEntry, Provenance};
use quarb::Prov;
use quarb::kaiv_out::{ProvCtx, ident_of, kaiv_put};
let mut b = KaivBuilder::new();
b.declare_source("q", source).map_err(kaiv_err)?;
fn origin_nodes(p: &Prov, out: &mut Vec<NodeId>) {
out.extend(p.origins.nodes());
for (_, q) in p.parts.iter() {
origin_nodes(q, out);
}
}
let mut nodes: Vec<NodeId> = Vec::new();
for r in rows {
nodes.push(r.node);
origin_nodes(&r.prov, &mut nodes);
}
let mut source_ids: std::collections::HashMap<String, String> =
std::collections::HashMap::new();
for n in &nodes {
for p in prov_of(*n).entries {
let Some(src) = p.source else { continue };
if source_ids.contains_key(&src) {
continue;
}
let id = format!("src{}", source_ids.len() + 1);
let id = match b.declare_source(&id, &src) {
Ok(()) => id,
Err(_) => "q".to_string(),
};
source_ids.insert(src, id);
}
}
let entry_of = |n: NodeId, p: &quarb::Provenance| -> ProvEntry {
ProvEntry {
source: p
.source
.as_ref()
.and_then(|s| source_ids.get(s).cloned())
.unwrap_or_else(|| "q".to_string()),
timestamp: p
.instant
.map(|(secs, _, _)| quarb::temporal::format_instant(secs, 0, Some(0)))
.filter(|t| t.len() == 20),
dpid: p.dpid.as_deref().map(ident_of).or_else(|| {
let loc = match p.path.as_deref() {
Some(path) if path != "/" => ident_of(path),
_ => ident_of(&render(n)),
};
match &p.source {
Some(s) if ident_of(s) == loc => None,
_ => Some(loc),
}
}),
}
};
for (i, row) in rows.iter().enumerate() {
let resolve = |o: &quarb::Origins| -> Option<Provenance> {
let mut entries: Vec<ProvEntry> = Vec::new();
let mut elided = o.more;
let ns: Vec<NodeId> = if o.is_empty() {
vec![row.node]
} else {
o.nodes().collect()
};
for n in ns {
let list = prov_of(n);
elided = elided.saturating_add(list.elided);
let ps = if list.entries.is_empty() {
vec![quarb::Provenance::default()]
} else {
list.entries
};
for p in ps {
let e = entry_of(n, &p);
match entries
.iter()
.position(|q| q.source == e.source && q.dpid == e.dpid)
{
Some(at) => {
if e.timestamp > entries[at].timestamp {
entries[at].timestamp = e.timestamp;
}
}
None if entries.len() < cap => entries.push(e),
None => elided = elided.saturating_add(1),
}
}
}
Some(Provenance { entries, elided })
};
let ctx = ProvCtx {
prov: Some(&row.prov),
resolve: &resolve,
};
let base = format!("/@results/{i}");
let mut used = std::collections::HashSet::new();
match &row.topic {
None => {
let loc = render(row.node);
kaiv_put(&mut b, &base, "node", &Value::Str(loc), &ctx, &mut used)
.map_err(anyhow::Error::msg)?;
}
Some(Value::Record(fields)) => {
for (k, v) in fields {
let part = row.prov.part(k);
let fc = ProvCtx {
prov: Some(&part),
resolve: &resolve,
};
kaiv_put(&mut b, &base, k, v, &fc, &mut used).map_err(anyhow::Error::msg)?;
}
}
Some(v) => {
kaiv_put(&mut b, &base, "value", v, &ctx, &mut used).map_err(anyhow::Error::msg)?
}
}
}
b.finish().map_err(kaiv_err)
}
fn kaiv_err(e: kaiv::PipelineError) -> anyhow::Error {
anyhow::anyhow!("emitting kaiv: {e}")
}
#[cfg(test)]
mod tests {
use super::{split_alias, split_scheme_query};
use std::path::{Path, PathBuf};
#[test]
fn mount_aliases_split() {
assert_eq!(
split_alias(Path::new("ga=bigquery://p/quarb_ga?account=a@b.c")),
Some((
"ga".to_string(),
PathBuf::from("bigquery://p/quarb_ga?account=a@b.c")
))
);
assert_eq!(
split_alias(Path::new("raw_2026-06=events.jsonl")),
Some(("raw_2026-06".to_string(), PathBuf::from("events.jsonl")))
);
assert_eq!(split_alias(Path::new("events.jsonl")), None);
assert_eq!(split_alias(Path::new("ga=")), None);
assert_eq!(
split_alias(Path::new("толстой=corpus:anna-karenina.atd")),
Some((
"толстой".to_string(),
PathBuf::from("corpus:anna-karenina.atd")
))
);
assert_eq!(split_alias(Path::new("2ga=x.json")), None);
assert_eq!(split_alias(Path::new("a/b=x.json")), None);
}
#[test]
fn scheme_prefixed_queries_split() {
assert_eq!(
split_scheme_query("github:/torvalds/linux::stars"),
Some(("github:", "/torvalds/linux::stars"))
);
assert_eq!(
split_scheme_query("gitlab:/tesslab//*<repo>"),
Some(("gitlab:", "/tesslab//*<repo>"))
);
assert_eq!(
split_scheme_query("k8s:/namespaces/*"),
Some(("k8s:", "/namespaces/*"))
);
assert_eq!(split_scheme_query("github:torvalds/linux"), None);
assert_eq!(split_scheme_query("git:/repo"), None);
assert_eq!(split_scheme_query("/a/b::c"), None);
}
fn row(n: u64, topic: Option<quarb::Value>) -> quarb::Traced {
quarb::Traced {
node: quarb::NodeId(n),
topic,
prov: quarb::Prov::default(),
}
}
#[test]
fn emit_kaiv_provenance_per_row() {
use quarb::{NodeId, Provenance, ProvenanceList, Value};
let rows = vec![
row(1, Some(Value::Int(7))),
row(2, Some(Value::Int(9))),
row(3, Some(Value::Int(11))),
];
let render = |n: NodeId| format!("/row/{}", n.0);
let (secs, _, _) = quarb::temporal::parse_iso("2026-07-17T12:00:00Z").unwrap();
let prov_of = move |n: NodeId| match n.0 {
1 => Provenance {
source: Some("https://sensors.example.com/1".into()),
instant: Some((secs, 0, Some(0))),
dpid: Some("req-42".into()),
..Default::default()
},
2 => Provenance {
source: Some("https://sensors.example.com/1".into()),
..Default::default()
},
_ => Provenance::default(),
};
let out = super::emit_kaiv(
&rows,
"a.daiv, b.csv",
render,
|n| ProvenanceList::single(prov_of(n)),
8,
)
.unwrap();
assert!(out.contains(".?q a.daiv, b.csv\n"));
assert_eq!(out.matches("sensors.example.com").count(), 1);
assert!(out.contains(".?src1 https://sensors.example.com/1\n"));
assert!(
out.contains("!int?src1@2026-07-17T12:00:00Z#req-42\nvalue=7"),
"{out}"
);
assert!(out.contains("!int?src1#row-2\nvalue=9"));
assert!(out.contains("!int?q#row-3\nvalue=11"));
let fs = super::emit_kaiv(
&[row(1, Some(Value::Int(1)))],
".",
|_| "/a/b.txt".to_string(),
|_| {
ProvenanceList::single(Provenance {
source: Some("/a/b.txt".into()),
..Default::default()
})
},
8,
)
.unwrap();
assert!(fs.contains("!int?src1\nvalue=1"), "{fs}");
assert!(!fs.contains('#'), "{fs}");
let nested = super::emit_kaiv(
&[row(
1,
Some(Value::Record(vec![
(
"r".to_string(),
Value::Record(vec![
("n".to_string(), Value::Str("grc.atd".into())),
(
"size".to_string(),
Value::Quantity {
value: 4025440.0,
base: "B".into(),
written: None,
},
),
]),
),
(
"tags".to_string(),
Value::list(vec![Value::Str("a".into()), Value::Str("b".into())]),
),
])),
)],
"u.json",
|_| "/0".to_string(),
|_| ProvenanceList::default(),
8,
)
.unwrap();
assert!(nested.contains("(/r)"), "{nested}");
assert!(nested.contains("\nn=grc.atd\n"), "{nested}");
assert!(nested.contains("!float:B?q#0\nsize=4025440\n"), "{nested}");
assert!(nested.contains("@tags"), "{nested}");
assert!(!nested.contains("{\"n\""), "{nested}");
let plain = super::emit_kaiv(
&rows,
"data.json",
|n: NodeId| format!("/r/{}", n.0),
|_| ProvenanceList::default(),
8,
)
.unwrap();
assert!(plain.contains(".?q data.json\n"));
assert!(plain.contains("!int?q#r-1\nvalue=7"));
assert!(!plain.contains(".?q-"));
}
fn prov_back(back: &quarb_kaiv::KaivAdapter, q: &str) -> Vec<String> {
match quarb::run(q, back).unwrap() {
quarb::QueryResult::Values(vs) => vs.iter().map(|v| v.to_string()).collect(),
other => panic!("{other:?}"),
}
}
#[test]
fn emit_kaiv_per_field_and_plural() {
use quarb::{NodeId, Origin, Origins, Prov, Provenance, ProvenanceList, Value};
let (t1, _, _) = quarb::temporal::parse_iso("2026-07-17T12:00:00Z").unwrap();
let (t2, _, _) = quarb::temporal::parse_iso("2026-07-18T12:00:00Z").unwrap();
let prov_of = move |n: NodeId| {
ProvenanceList::single(Provenance {
source: Some(format!("https://s.example.com/{}", n.0)),
instant: Some((if n.0 == 2 { t2 } else { t1 }, 0, Some(0))),
dpid: Some(format!("row-{}", n.0)),
..Default::default()
})
};
let at = |ns: &[u64]| {
let mut o = Origins::default();
for &n in ns {
o.insert(Origin {
node: NodeId(n),
command: 1,
});
}
Prov::leaf(o)
};
let rec = Value::Record(vec![
("v".into(), Value::Int(1)),
("n".into(), Value::Int(2)),
("k".into(), Value::Int(3)),
]);
let mut prov = Prov::record(vec![("v".into(), at(&[1])), ("n".into(), at(&[1, 2]))]);
std::sync::Arc::make_mut(&mut prov.parts).push(("k".into(), Prov::default()));
let rows = vec![quarb::Traced {
node: NodeId(3),
topic: Some(rec),
prov,
}];
let out = super::emit_kaiv(&rows, "t", |n| format!("/r/{}", n.0), prov_of, 8).unwrap();
assert!(out.contains(".?src1 https://s.example.com/3\n"), "{out}");
assert!(out.contains(".?src2 https://s.example.com/1\n"), "{out}");
assert!(out.contains(".?src3 https://s.example.com/2\n"), "{out}");
assert!(
out.contains("!int?src2@2026-07-17T12:00:00Z#row-1\nv=1\n"),
"{out}"
);
assert!(
out.contains(
"!int?src2@2026-07-17T12:00:00Z#row-1;src3@2026-07-18T12:00:00Z#row-2\nn=2\n"
),
"{out}"
);
assert!(
out.contains("!int?src1@2026-07-17T12:00:00Z#row-3\nk=3\n"),
"{out}"
);
let back = quarb_kaiv::KaivAdapter::parse_kaiv(&out).unwrap();
let prov = |q: &str| match quarb::run(q, &back).unwrap() {
quarb::QueryResult::Values(vs) => vs.iter().map(|v| v.to_string()).collect::<Vec<_>>(),
other => panic!("{other:?}"),
};
assert_eq!(
prov("/@results/0/n:::provenance"),
[
"?https://s.example.com/1@2026-07-17T12:00:00Z#row-1;https://s.example.com/2@2026-07-18T12:00:00Z#row-2"
]
);
assert_eq!(prov("/@results/0/n:::instant"), ["2026-07-18T12:00:00Z"]);
assert_eq!(prov("/@results/0/n:::@provenance | count"), ["2"]);
assert_eq!(prov("/@results/0/v:::elided"), ["0"]);
let mut o = Origins::default();
for n in 1..=12 {
o.insert(Origin {
node: NodeId(n),
command: 1,
});
}
assert_eq!(o.more, 12 - quarb::ORIGIN_CAP as u32);
let rows = vec![quarb::Traced {
node: NodeId(1),
topic: Some(Value::Int(78)),
prov: Prov::leaf(o),
}];
let out = super::emit_kaiv(&rows, "t", |n| format!("/r/{}", n.0), prov_of, 3).unwrap();
let line = out.lines().find(|l| l.starts_with("!int?")).unwrap();
assert_eq!(line.matches("#row-").count(), 3, "{line}");
assert!(line.ends_with(";+9"), "{line}");
let back = quarb_kaiv::KaivAdapter::parse_kaiv(&out).unwrap();
assert_eq!(prov_back(&back, "/@results/0/value:::elided"), ["9"]);
assert_eq!(
prov_back(&back, "/@results/0/value:::@provenance | count"),
["3"]
);
let same = move |_: NodeId| {
ProvenanceList::single(Provenance {
source: Some("https://s.example.com/x".into()),
instant: None,
dpid: Some("row-x".into()),
..Default::default()
})
};
let rows = vec![quarb::Traced {
node: NodeId(1),
topic: Some(Value::Int(1)),
prov: at(&[1, 2]),
}];
let out = super::emit_kaiv(&rows, "t", |n| format!("/r/{}", n.0), same, 8).unwrap();
assert!(out.contains("!int?src1#row-x\nvalue=1\n"), "{out}");
}
}