use clap::Parser as _;
use faucet_cli::cli::Cli;
use std::path::Path;
fn on_big_stack<F>(f: impl FnOnce() -> F + Send + 'static)
where
F: std::future::Future<Output = ()>,
{
std::thread::Builder::new()
.stack_size(32 * 1024 * 1024)
.spawn(move || {
tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.unwrap()
.block_on(f())
})
.unwrap()
.join()
.unwrap();
}
async fn run(args: &[&str]) -> Result<(), faucet_cli::error::CliError> {
let mut argv = vec!["faucet"];
argv.extend_from_slice(args);
let cli = Cli::try_parse_from(argv).expect("argv parses");
Box::pin(faucet_cli::run_command(cli)).await
}
fn fixture(dir: &Path) -> (String, String) {
let hub = dir.join("hub");
std::fs::create_dir_all(hub.join("source-templates")).unwrap();
std::fs::create_dir_all(hub.join("sink-templates")).unwrap();
let data = dir.join("data");
std::fs::create_dir_all(&data).unwrap();
std::fs::write(
data.join("orders.csv"),
"id,amount,region\n1,10,eu\n2,20,us\n",
)
.unwrap();
std::fs::write(data.join("customers.csv"), "id,name\n7,Acme\n").unwrap();
std::fs::write(
hub.join("source-templates/shop.yaml"),
format!(
r#"
kind: source-template
name: shop
description: A shop's exports
tags: [demo]
params:
data_dir: {{ type: string, default: "{data}" }}
source:
type: csv
config: {{ path: "${{param.data_dir}}/orders.csv", has_headers: true }}
transforms:
- {{ type: keys_case, config: {{ mode: snake }} }}
streams:
- name: orders
primary_keys: [id]
write: [overwrite, upsert]
- name: customers
source: {{ config: {{ path: "${{param.data_dir}}/customers.csv" }} }}
write: overwrite
"#,
data = data.display()
),
)
.unwrap();
std::fs::write(
hub.join("sink-templates/files.yaml"),
format!(
r#"
kind: sink-template
name: files
description: JSON Lines files
params:
out_dir: {{ type: string, default: "{out}" }}
sink:
type: jsonl
config: {{ append: false }}
per_stream:
path: "${{param.out_dir}}/${{source}}/${{stream}}.jsonl"
write_mode_aliases:
overwrite: append
"#,
out = dir.join("out").display()
),
)
.unwrap();
std::fs::write(
hub.join("sink-templates/plain.yaml"),
r#"
kind: sink-template
name: plain
description: Append-only files
sink:
type: jsonl
config: {}
per_stream:
path: "./plain/${stream}.jsonl"
"#,
)
.unwrap();
(
hub.to_string_lossy().into_owned(),
dir.join("out").to_string_lossy().into_owned(),
)
}
#[tokio::test]
async fn compose_check_list_lint_validate_and_run_a_pairing() {
let dir = tempfile::tempdir().unwrap();
let (hub, out) = fixture(dir.path());
let composed = dir.path().join("composed.yaml");
run(&[
"hub",
"compose",
"--source",
"shop",
"--sink",
"files",
"--hub",
&hub,
"--out",
composed.to_str().unwrap(),
])
.await
.expect("compose");
let text = std::fs::read_to_string(&composed).unwrap();
assert!(
text.contains("name: shop") && text.contains("orders.jsonl"),
"{text}"
);
run(&[
"validate",
composed.to_str().unwrap(),
"--no-env-file",
"--no-secrets",
])
.await
.expect("validate composed file");
run(&[
"hub", "compose", "--source", "shop", "--sink", "files", "--hub", &hub, "--json",
])
.await
.expect("compose --json");
run(&[
"hub", "check", "--source", "shop", "--sink", "files", "--hub", &hub,
])
.await
.expect("check");
run(&[
"hub", "check", "--source", "shop", "--sink", "files", "--hub", &hub, "--json",
])
.await
.expect("check --json");
let err = run(&[
"hub", "check", "--source", "shop", "--sink", "plain", "--hub", &hub,
])
.await
.unwrap_err()
.to_string();
assert!(
err.contains("2 stream(s) of 'shop' have no write mode sink 'plain' supports"),
"{err}"
);
run(&["hub", "list", "--hub", &hub]).await.expect("list");
run(&["hub", "list", "--hub", &hub, "--json"])
.await
.expect("list --json");
for fmt in ["table", "markdown", "json"] {
run(&["hub", "matrix", "--hub", &hub, "--format", fmt])
.await
.expect("matrix");
}
let md = dir.path().join("matrix.md");
run(&[
"hub",
"matrix",
"--hub",
&hub,
"--format",
"markdown",
"--out",
md.to_str().unwrap(),
])
.await
.unwrap();
let md = std::fs::read_to_string(md).unwrap();
assert!(md.contains("| [shop](#shop) | ✓ | — |"), "{md}");
run(&["hub", "lint", "--hub", &hub])
.await
.expect("lint catalog");
let one = format!("{hub}/source-templates/shop.yaml");
run(&["hub", "lint", &one]).await.expect("lint one file");
run(&["hub", "lint", &one, "--json"])
.await
.expect("lint one file --json");
run(&[
"validate",
"--source",
"shop",
"--sink",
"files",
"--hub",
&hub,
"--no-env-file",
"--no-secrets",
])
.await
.expect("validate --source/--sink");
run(&[
"validate",
"--source",
"shop",
"--sink",
"files",
"--hub",
&hub,
"--no-env-file",
"--show-composed",
])
.await
.expect("validate --show-composed");
run(&[
"validate",
"--source",
"shop",
"--sink",
"files",
"--hub",
&hub,
"--no-env-file",
"--no-secrets",
"--json",
])
.await
.expect("validate --json");
run(&[
"run",
"--source",
"shop",
"--sink",
"files",
"--hub",
&hub,
"--no-env-file",
"--quiet",
])
.await
.expect("run --source/--sink");
let orders =
std::fs::read_to_string(Path::new(&out).join("shop/orders.jsonl")).expect("orders written");
assert_eq!(orders.lines().count(), 2, "{orders}");
assert!(orders.contains("\"region\":\"eu\""), "{orders}");
let customers = std::fs::read_to_string(Path::new(&out).join("shop/customers.jsonl"))
.expect("customers written");
assert_eq!(customers.lines().count(), 1);
let other = dir.path().join("other");
run(&[
"run",
"--source",
"shop",
"--sink",
"files",
"--hub",
&hub,
"--no-env-file",
"--quiet",
"--param",
&format!("out_dir={}", other.display()),
])
.await
.expect("run with --param");
assert!(other.join("shop/orders.jsonl").is_file());
}
#[tokio::test]
async fn hub_templates_handed_to_run_or_validate_directly_are_redirected() {
let dir = tempfile::tempdir().unwrap();
let (hub, _) = fixture(dir.path());
let tpl = format!("{hub}/source-templates/shop.yaml");
let err = run(&["run", &tpl, "--no-env-file"])
.await
.unwrap_err()
.to_string();
assert!(
err.contains("is a hub source-template")
&& err.contains("--source <source-template> --sink <sink-template>"),
"{err}"
);
let err = run(&["validate", &tpl, "--no-env-file", "--no-secrets"])
.await
.unwrap_err()
.to_string();
assert!(err.contains("is a hub source-template"), "{err}");
let err = run(&[
"run",
"--source",
"nope",
"--sink",
"files",
"--hub",
&hub,
"--no-env-file",
])
.await
.unwrap_err()
.to_string();
assert!(
err.contains("no hub template 'nope'") && err.contains("known: shop"),
"{err}"
);
let err = run(&[
"run",
"--source",
"shop",
"--sink",
"plain",
"--hub",
&hub,
"--no-env-file",
])
.await
.unwrap_err()
.to_string();
assert!(
err.contains("cannot compose with sink-template 'plain'"),
"{err}"
);
let plain = dir.path().join("plain.yaml");
std::fs::write(&plain, "version: 1\n").unwrap();
let err = run(&["hub", "lint", plain.to_str().unwrap()])
.await
.unwrap_err()
.to_string();
assert!(err.contains("not a hub template"), "{err}");
let bad = format!("{hub}/source-templates/bad.yaml");
std::fs::write(
&bad,
"kind: source-template\nname: bad\nsource: {type: csv, config: {path: x.csv, api_key: literal-key}}\nstreams: [{name: t}]\n",
)
.unwrap();
let err = run(&["hub", "lint", &bad]).await.unwrap_err().to_string();
assert!(err.contains("1 template(s) with findings"), "{err}");
let err = run(&["hub", "lint", "--hub", &hub, "--json"])
.await
.unwrap_err()
.to_string();
assert!(err.contains("with findings"), "{err}");
}
#[tokio::test]
async fn schema_targets_print() {
run(&["schema", "source-template"])
.await
.expect("schema source-template");
run(&["schema", "sink-template"])
.await
.expect("schema sink-template");
}
#[test]
fn hub_flags_are_mutually_required_and_exclusive_with_a_config_path() {
assert!(Cli::try_parse_from(["faucet", "run", "--source", "a"]).is_err());
assert!(Cli::try_parse_from(["faucet", "run", "--sink", "b"]).is_err());
assert!(
Cli::try_parse_from(["faucet", "run", "cfg.yaml", "--source", "a", "--sink", "b"]).is_err()
);
assert!(
Cli::try_parse_from([
"faucet",
"run",
"--from-env",
"--source",
"a",
"--sink",
"b"
])
.is_err()
);
assert!(Cli::try_parse_from(["faucet", "run", "--source", "a", "--sink", "b"]).is_ok());
assert!(
Cli::try_parse_from([
"faucet", "validate", "cfg.yaml", "--source", "a", "--sink", "b"
])
.is_err()
);
assert!(Cli::try_parse_from(["faucet", "validate", "--source", "a", "--sink", "b"]).is_ok());
}
#[cfg(all(feature = "source-csv", feature = "sink-sqlite"))]
#[test]
fn shipped_example_pairing_runs_twice_against_a_fresh_database() {
on_big_stack(|| async {
let repo = Path::new(env!("CARGO_MANIFEST_DIR")).parent().unwrap();
let hub = repo.join("hub");
let data = hub.join("examples/data");
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("fresh.db");
let csv_rows = |name: &str| {
std::fs::read_to_string(data.join(name))
.unwrap()
.lines()
.skip(1)
.filter(|l| !l.trim().is_empty())
.count() as i64
};
for pass in 1..=2 {
run(&[
"run",
"--source",
"faucet-hq/example-csv",
"--sink",
"faucet-hq/sqlite",
"--hub",
hub.to_str().unwrap(),
"--no-env-file",
"--quiet",
"--param",
&format!("data_dir={}", data.display()),
"--param",
&format!("sqlite_path={}", db.display()),
])
.await
.unwrap_or_else(|e| panic!("pass {pass}: {e}"));
let pool = sqlx::SqlitePool::connect(&format!("sqlite:{}", db.display()))
.await
.unwrap();
for (table, file) in [("orders", "orders.csv"), ("customers", "customers.csv")] {
let n: i64 = sqlx::query_scalar(&format!("SELECT COUNT(*) FROM {table}"))
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(n, csv_rows(file), "pass {pass}: {table}");
}
pool.close().await;
}
});
}
#[cfg(all(feature = "source-csv", feature = "sink-sqlite"))]
#[test]
fn a_deployment_overlay_reaches_the_composed_run() {
on_big_stack(|| async {
let repo = Path::new(env!("CARGO_MANIFEST_DIR")).parent().unwrap();
let hub = repo.join("hub");
let data = hub.join("examples/data");
let dir = tempfile::tempdir().unwrap();
let state = dir.path().join("state");
let overlay = dir.path().join("ops.yaml");
std::fs::write(
&overlay,
format!(
"kind: deployment\nname: ops\nstate: {{ type: file, config: {{ path: \"{}\" }} }}\nstreams:\n orders: {{ sla: {{ min_rows_per_run: 1 }} }}\n",
state.display()
),
)
.unwrap();
let args = |verb: &'static str| {
vec![
verb.to_string(),
"--source".into(),
"faucet-hq/example-csv".into(),
"--sink".into(),
"faucet-hq/sqlite".into(),
"--overlay".into(),
overlay.display().to_string(),
"--hub".into(),
hub.display().to_string(),
"--no-env-file".into(),
"--param".into(),
format!("data_dir={}", data.display()),
"--param".into(),
format!("sqlite_path={}", dir.path().join("o.db").display()),
]
};
let validate = args("validate");
run(&validate.iter().map(String::as_str).collect::<Vec<_>>())
.await
.expect("validate --overlay");
let mut run_args = args("run");
run_args.push("--quiet".into());
run(&run_args.iter().map(String::as_str).collect::<Vec<_>>())
.await
.expect("run --overlay");
let written: Vec<String> = std::fs::read_dir(&state)
.expect("the overlay's state store was used")
.map(|e| e.unwrap().file_name().to_string_lossy().into_owned())
.collect();
assert!(
written
.iter()
.any(|f| f.contains("orders") && f.contains("__sla__")),
"per-stream SLA history in the overlay's store: {written:?}"
);
assert!(
!written
.iter()
.any(|f| f.contains("customers") && f.contains("__sla__")),
"only the stream the overlay gave an SLA: {written:?}"
);
let err = run(&["run", overlay.to_str().unwrap(), "--no-env-file"])
.await
.unwrap_err()
.to_string();
assert!(
err.contains("deployment overlay") && err.contains("--overlay"),
"{err}"
);
});
}
#[tokio::test]
async fn a_source_and_a_sink_compose_across_two_hubs() {
let dir = tempfile::tempdir().unwrap();
let (private, _out) = fixture(dir.path());
let public = dir.path().join("public-hub");
std::fs::create_dir_all(&public).unwrap();
std::fs::rename(
Path::new(&private).join("sink-templates"),
public.join("sink-templates"),
)
.unwrap();
std::fs::create_dir_all(Path::new(&private).join("sink-templates")).unwrap();
std::fs::create_dir_all(public.join("source-templates")).unwrap();
let public = public.to_string_lossy().into_owned();
let composed = dir.path().join("composed.yaml");
let read = |p: &Path| std::fs::read_to_string(p).unwrap();
run(&[
"hub",
"compose",
"--source",
"shop",
"--source-hub",
&private,
"--sink",
"files",
"--hub",
&public,
"--out",
composed.to_str().unwrap(),
])
.await
.expect("compose across hubs");
let text = read(&composed);
assert!(
text.contains(&format!("# source: shop from {private}")),
"{text}"
);
assert!(
text.contains(&format!("# sink: files from {public}")),
"{text}"
);
let both = format!("{private},{public}");
run(&[
"hub", "check", "--source", "shop", "--sink", "files", "--hub", &both,
])
.await
.expect("an ordered hub list");
run(&[
"validate",
"--source",
"shop",
"--sink",
"files",
"--hub",
&private,
"--hub",
&public,
"--no-env-file",
"--no-secrets",
])
.await
.expect("validate across hubs");
run(&[
"hub",
"compose",
"--source",
"shop",
"--hub",
&private,
"--sink",
&format!("{public}:files"),
"--json",
])
.await
.expect("a qualified sink locator");
std::fs::write(
Path::new(&private).join("sink-templates/files.yaml"),
"kind: sink-template\nname: files\ndescription: Private files\nsink:\n type: jsonl\n config: {append: false}\nper_stream:\n path: \"./private-copy/${stream}.jsonl\"\nwrite_mode_aliases:\n overwrite: append\n",
)
.unwrap();
run(&[
"hub",
"compose",
"--source",
"shop",
"--sink",
"files",
"--hub",
&both,
"--out",
composed.to_str().unwrap(),
])
.await
.unwrap();
assert!(
read(&composed).contains("private-copy"),
"the first hub wins"
);
let err = run(&[
"hub", "check", "--source", "shop", "--sink", "nope", "--hub", &both,
])
.await
.unwrap_err()
.to_string();
assert!(err.contains("in any of the 2 hubs searched"), "{err}");
assert!(err.contains(&private) && err.contains(&public), "{err}");
}
#[cfg(feature = "hub-remote")]
#[test]
fn a_per_owner_github_token_takes_precedence() {
use faucet_cli::hub::remote::{github_token, owner_token_var};
let var = owner_token_var("zz-hub-cli-696/private-hub");
assert_eq!(var, "FAUCET_GITHUB_TOKEN_ZZ_HUB_CLI_696");
unsafe { std::env::set_var(&var, "owner-token") };
assert_eq!(
github_token("zz-hub-cli-696/private-hub").as_deref(),
Some("owner-token")
);
let loc = faucet_cli::hub::HubLocation::parse("github:zz-hub-cli-696/private-hub").unwrap();
assert!(faucet_cli::hub::remote::GithubHub::new(&loc, "http://127.0.0.1:1").is_ok());
unsafe { std::env::remove_var(&var) };
}