use crate::cli::DiscoverArgs;
use crate::config::{ConnectorSpec, PipelineConfig, PipelineSpec};
use crate::error::{CliError, CliResult};
use faucet_core::DatasetDescriptor;
use serde_json::{Value, json};
pub async fn run(args: DiscoverArgs) -> CliResult<()> {
let cwd = std::env::current_dir()?;
let env_path =
crate::env_loader::resolve_env_file(args.env_file.as_deref(), args.no_env_file, &cwd)?;
crate::env_loader::load_env_file_if_present(env_path.as_deref())?;
let path = match args.config {
Some(p) => p,
None => crate::env_loader::discover_config_path(&cwd).ok_or(CliError::NoConfigOrFromEnv)?,
};
let cfg = PipelineConfig::from_path_async(&path, args.profile.as_deref()).await?;
let raw_composed = crate::compose::compose(&path, args.profile.as_deref())?;
let template = args.source.as_deref().unwrap_or("default");
let spec = select_source_template(&cfg.pipeline, template)?;
let auth = crate::auth_catalog::build_auth_catalog(cfg.auth.as_ref())?;
let source =
crate::registry::build_source(&spec.kind, spec.config.clone(), &auth, None).await?;
if !source.supports_discover() {
return Err(CliError::Config(format!(
"source '{}' does not support dataset discovery — discovery is available for \
catalog-backed sources (postgres, mysql, mssql, sqlite, mongodb, elasticsearch, \
bigquery, snowflake, spanner, s3, gcs)",
spec.kind
)));
}
let datasets = source.discover().await.map_err(|e| {
CliError::Config(format!(
"discovery against source '{}' failed: {e}",
spec.kind
))
})?;
let total = datasets.len();
let datasets = filter_datasets(datasets, &args.include, &args.exclude);
if datasets.is_empty() {
return Err(CliError::Config(format!(
"discovery found no datasets{} (source '{}' reported {total} before filtering)",
if args.include.is_empty() && args.exclude.is_empty() {
""
} else {
" matching the --include/--exclude filters"
},
spec.kind
)));
}
if args.json {
let out = json!({ "source": spec.kind, "datasets": datasets });
println!(
"{}",
serde_json::to_string_pretty(&out)
.map_err(|e| CliError::Internal(format!("json render: {e}")))?
);
return Ok(());
}
let doc = render_discovered_config(&raw_composed, template, &datasets)?;
if let Err(e) =
PipelineConfig::from_text(&doc, &path).and_then(|c| crate::expand::expand(&c).map(|_| ()))
{
tracing::warn!(error = %e, "generated config failed offline re-validation — review it before running");
}
match args.output {
Some(out) => {
if out.exists() && !args.force {
return Err(CliError::Config(format!(
"output file {} already exists — pass --force to overwrite",
out.display()
)));
}
std::fs::write(&out, &doc)?;
eprintln!(
"wrote {} ({} dataset{} from '{}' source)",
out.display(),
datasets.len(),
if datasets.len() == 1 { "" } else { "s" },
spec.kind
);
}
None => print!("{doc}"),
}
Ok(())
}
fn select_source_template<'a>(
pipeline: &'a PipelineSpec,
name: &str,
) -> CliResult<&'a ConnectorSpec> {
if let Some(spec) = pipeline.sources.get(name) {
return Ok(spec);
}
if name == "default"
&& let Some(spec) = pipeline.source.as_ref()
{
return Ok(spec);
}
let mut available: Vec<&str> = pipeline.sources.keys().map(String::as_str).collect();
if pipeline.source.is_some() {
available.push("default");
}
available.sort_unstable();
Err(CliError::Config(format!(
"no source template named '{name}' — available: {}",
if available.is_empty() {
"none (the config has no `pipeline.source` or `pipeline.sources`)".to_string()
} else {
available.join(", ")
}
)))
}
fn glob_match(pattern: &str, name: &str) -> bool {
let segments: Vec<&str> = pattern.split('*').collect();
if segments.len() == 1 {
return pattern == name;
}
let mut rest = name;
for (i, seg) in segments.iter().enumerate() {
if seg.is_empty() {
continue;
}
if i == 0 {
match rest.strip_prefix(seg) {
Some(r) => rest = r,
None => return false,
}
} else if i == segments.len() - 1 {
return rest.ends_with(seg) && rest.len() >= seg.len();
} else {
match rest.find(seg) {
Some(pos) => rest = &rest[pos + seg.len()..],
None => return false,
}
}
}
true
}
fn filter_datasets(
datasets: Vec<DatasetDescriptor>,
include: &[String],
exclude: &[String],
) -> Vec<DatasetDescriptor> {
datasets
.into_iter()
.filter(|d| {
let included = include.is_empty() || include.iter().any(|p| glob_match(p, &d.name));
let excluded = exclude.iter().any(|p| glob_match(p, &d.name));
included && !excluded
})
.collect()
}
fn sanitize_row_id(name: &str) -> String {
let mut out = String::with_capacity(name.len());
let mut last_was_sep = false;
for c in name.chars() {
if c.is_ascii_alphanumeric() || c == '_' || c == '-' {
out.push(c);
last_was_sep = false;
} else if !last_was_sep && !out.is_empty() {
out.push('_');
last_was_sep = true;
}
}
let trimmed = out.trim_matches('_').to_string();
if trimmed.is_empty() {
"dataset".to_string()
} else {
trimmed
}
}
fn unique_row_ids(datasets: &[DatasetDescriptor]) -> Vec<String> {
let mut seen = std::collections::HashMap::<String, usize>::new();
datasets
.iter()
.map(|d| {
let base = sanitize_row_id(&d.name);
let n = seen.entry(base.clone()).or_insert(0);
*n += 1;
if *n == 1 { base } else { format!("{base}-{n}") }
})
.collect()
}
fn schema_summary(schema: &Value, max_cols: usize) -> Option<String> {
let props = schema.get("properties")?.as_object()?;
if props.is_empty() {
return None;
}
let mut parts: Vec<String> = Vec::new();
for (name, fragment) in props.iter().take(max_cols) {
let (ty, nullable) = match fragment.get("type") {
Some(Value::String(t)) => (t.clone(), false),
Some(Value::Array(a)) => {
let base = a
.iter()
.filter_map(|v| v.as_str())
.find(|t| *t != "null")
.unwrap_or("any");
(base.to_string(), a.iter().any(|v| v == "null"))
}
_ => ("any".to_string(), false),
};
parts.push(format!("{name} {ty}{}", if nullable { "?" } else { "" }));
}
if props.len() > max_cols {
parts.push("…".to_string());
}
Some(parts.join(", "))
}
fn yaml_seq_item(value: &Value) -> CliResult<String> {
let body = serde_yaml::to_string(value)
.map_err(|e| CliError::Internal(format!("yaml render: {e}")))?;
let mut out = String::new();
for (i, line) in body.trim_end().lines().enumerate() {
if i == 0 {
out.push_str(" - ");
} else {
out.push_str(" ");
}
out.push_str(line);
out.push('\n');
}
Ok(out)
}
fn render_discovered_config(
raw_composed: &str,
template: &str,
datasets: &[DatasetDescriptor],
) -> CliResult<String> {
let mut root: serde_yaml::Value = serde_yaml::from_str(raw_composed)
.map_err(|e| CliError::Config(format!("could not re-parse composed config: {e}")))?;
let map = root
.as_mapping_mut()
.ok_or_else(|| CliError::Config("composed config is not a mapping".into()))?;
let replaced_matrix = map
.remove(serde_yaml::Value::String("matrix".into()))
.is_some();
let head = serde_yaml::to_string(&root)
.map_err(|e| CliError::Internal(format!("yaml render: {e}")))?;
let row_ids = unique_row_ids(datasets);
let mut doc = String::new();
doc.push_str(&head);
if !head.ends_with('\n') {
doc.push('\n');
}
doc.push('\n');
doc.push_str(&format!(
"# Generated by `faucet discover` — one row per discovered dataset ({}).\n",
datasets.len()
));
if replaced_matrix {
doc.push_str("# NOTE: the input config's `matrix:` block was replaced.\n");
}
doc.push_str("matrix:\n");
for (d, id) in datasets.iter().zip(&row_ids) {
let est = d
.estimated_rows
.map(|n| format!(", ~{n} rows"))
.unwrap_or_default();
doc.push_str(&format!(" # {} ({}{})\n", d.name, d.kind, est));
if let Some(summary) = d.schema.as_ref().and_then(|s| schema_summary(s, 12)) {
doc.push_str(&format!(" # columns: {summary}\n"));
}
let mut row = json!({ "id": id, "source": { "config": d.config_patch } });
if template != "default" {
row["source"]["ref"] = json!(template);
}
doc.push_str(&yaml_seq_item(&row)?);
}
Ok(doc)
}
#[cfg(test)]
mod tests {
use super::*;
fn ds(name: &str, kind: &str, patch: Value) -> DatasetDescriptor {
DatasetDescriptor::new(name, kind, patch)
}
#[test]
fn glob_exact_and_wildcards() {
assert!(glob_match("orders", "orders"));
assert!(!glob_match("orders", "orders2"));
assert!(glob_match("public.*", "public.orders"));
assert!(!glob_match("public.*", "sales.orders"));
assert!(glob_match("*.orders", "public.orders"));
assert!(glob_match("*order*", "public.orders"));
assert!(glob_match("*", "anything"));
assert!(glob_match("a*b*c", "aXXbYYc"));
assert!(!glob_match("a*b*c", "aXXcYYb"));
assert!(!glob_match("ab*ab", "ab"));
}
#[test]
fn filter_include_exclude_composition() {
let all = vec![
ds("public.orders", "table", json!({})),
ds("public.users", "table", json!({})),
ds("audit.log", "table", json!({})),
];
let got = filter_datasets(all.clone(), &["public.*".into()], &["*.users".into()]);
let names: Vec<&str> = got.iter().map(|d| d.name.as_str()).collect();
assert_eq!(names, vec!["public.orders"]);
let got = filter_datasets(all, &[], &["audit.*".into()]);
assert_eq!(got.len(), 2);
}
#[test]
fn row_ids_sanitized_and_unique() {
assert_eq!(sanitize_row_id("public.orders"), "public_orders");
assert_eq!(sanitize_row_id("raw/orders/"), "raw_orders");
assert_eq!(sanitize_row_id("...."), "dataset");
let ids = unique_row_ids(&[
ds("a.b", "table", json!({})),
ds("a/b", "table", json!({})),
ds("a.b", "table", json!({})),
]);
assert_eq!(ids, vec!["a_b", "a_b-2", "a_b-3"]);
}
#[test]
fn schema_summary_marks_nullable_and_caps() {
let schema = json!({
"type": "object",
"properties": {
"id": {"type": "integer"},
"note": {"type": ["string", "null"]},
}
});
let s = schema_summary(&schema, 12).unwrap();
assert!(s.contains("id integer"), "{s}");
assert!(s.contains("note string?"), "{s}");
let mut props = serde_json::Map::new();
for i in 0..15 {
props.insert(format!("c{i:02}"), json!({"type": "integer"}));
}
let big = json!({"type": "object", "properties": props});
let s = schema_summary(&big, 12).unwrap();
assert!(s.ends_with("…"), "capped: {s}");
}
#[test]
fn schema_summary_empty_is_none() {
assert!(schema_summary(&json!({"type": "object", "properties": {}}), 12).is_none());
assert!(schema_summary(&json!({"type": "string"}), 12).is_none());
}
const RAW: &str = r#"
version: 1
name: conn
pipeline:
source:
type: postgres
config:
connection_url: ${env:DATABASE_URL}
query: SELECT 1
sink:
type: jsonl
config:
path: ./out.jsonl
"#;
#[test]
fn render_emits_matrix_rows_with_comments() {
let datasets = vec![
ds(
"public.orders",
"table",
json!({"query": "SELECT * FROM \"public\".\"orders\""}),
)
.with_schema(json!({
"type": "object",
"properties": {"id": {"type": "integer"}}
}))
.with_estimated_rows(120),
ds(
"sales.leads",
"table",
json!({"query": "SELECT * FROM \"sales\".\"leads\""}),
),
];
let doc = render_discovered_config(RAW, "default", &datasets).unwrap();
assert!(doc.contains("${env:DATABASE_URL}"), "{doc}");
assert!(doc.contains("# public.orders (table, ~120 rows)"), "{doc}");
assert!(doc.contains("# columns: id integer"), "{doc}");
assert!(doc.contains("- id: public_orders"), "{doc}");
assert!(doc.contains("SELECT * FROM \"public\".\"orders\""), "{doc}");
let cfg = crate::config::parse_with_extension(&doc, "yaml").unwrap();
let nodes = crate::expand::expand(&cfg).unwrap();
assert_eq!(nodes.len(), 2);
assert_eq!(nodes[0].id, "public_orders");
assert_eq!(
nodes[0].source.config["query"], "SELECT * FROM \"public\".\"orders\"",
"row patch deep-merged over the connection config"
);
assert_eq!(
nodes[0].source.config["connection_url"], "${env:DATABASE_URL}",
"connection settings inherited"
);
}
#[test]
fn render_replaces_existing_matrix_and_notes_it() {
let raw = format!("{RAW}matrix:\n - id: old\n");
let datasets = vec![ds("t", "table", json!({"query": "SELECT * FROM t"}))];
let doc = render_discovered_config(&raw, "default", &datasets).unwrap();
assert!(doc.contains("was replaced"), "{doc}");
assert!(!doc.contains("id: old"), "{doc}");
}
#[test]
fn render_named_template_sets_source_ref() {
let raw = r#"
version: 1
pipeline:
sources:
warehouse:
type: postgres
config: { connection_url: "postgres://x", query: "SELECT 1" }
sinks:
default:
type: jsonl
config: { path: ./out.jsonl }
"#;
let datasets = vec![ds("t", "table", json!({"query": "SELECT * FROM t"}))];
let doc = render_discovered_config(raw, "warehouse", &datasets).unwrap();
assert!(doc.contains("ref: warehouse"), "{doc}");
let cfg = crate::config::parse_with_extension(&doc, "yaml").unwrap();
let nodes = crate::expand::expand(&cfg).unwrap();
assert_eq!(nodes.len(), 1);
assert_eq!(nodes[0].source.kind, "postgres");
}
#[test]
fn select_template_falls_back_to_singular_default() {
let cfg = crate::config::parse_with_extension(RAW, "yaml").unwrap();
let spec = select_source_template(&cfg.pipeline, "default").unwrap();
assert_eq!(spec.kind, "postgres");
let err = select_source_template(&cfg.pipeline, "nope").unwrap_err();
assert!(err.to_string().contains("available: default"), "{err}");
}
}
#[cfg(all(test, feature = "source-sqlite", feature = "sink-jsonl"))]
mod run_tests {
use super::run;
use crate::cli::DiscoverArgs;
fn args(config: std::path::PathBuf) -> DiscoverArgs {
DiscoverArgs {
config: Some(config),
source: None,
include: vec![],
exclude: vec![],
output: None,
force: false,
json: false,
env_file: None,
no_env_file: true,
profile: None,
}
}
fn write_config(dir: &std::path::Path, db: &str) -> std::path::PathBuf {
let cfg = dir.join("conn.yaml");
std::fs::write(
&cfg,
format!(
"version: 1\nname: conn\npipeline:\n source:\n type: sqlite\n config:\n database_url: \"sqlite://{db}\"\n query: SELECT 1\n sink:\n type: jsonl\n config: {{ path: ./out.jsonl }}\n"
),
)
.unwrap();
cfg
}
async fn seed_db(dir: &std::path::Path) -> String {
let db = dir.join("cat.db").display().to_string();
let pool = sqlx::SqlitePool::connect(&format!("sqlite://{db}?mode=rwc"))
.await
.expect("create db");
sqlx::query("CREATE TABLE orders (id INTEGER PRIMARY KEY, note TEXT)")
.execute(&pool)
.await
.unwrap();
sqlx::query("CREATE TABLE users (id INTEGER NOT NULL, active BOOLEAN)")
.execute(&pool)
.await
.unwrap();
pool.close().await;
db
}
#[tokio::test]
async fn run_errors_on_unsupported_source_kind() {
let dir = tempfile::tempdir().unwrap();
let cfg = dir.path().join("conn.yaml");
std::fs::write(
&cfg,
"version: 1\npipeline:\n source: { type: csv, config: { path: ./in.csv } }\n sink: { type: jsonl, config: { path: ./o.jsonl } }\n",
)
.unwrap();
let err = run(args(cfg)).await.unwrap_err();
assert!(
err.to_string()
.contains("does not support dataset discovery"),
"{err}"
);
}
#[tokio::test]
async fn run_writes_generated_config_that_expands() {
let dir = tempfile::tempdir().unwrap();
let db = seed_db(dir.path()).await;
let cfg = write_config(dir.path(), &db);
let out = dir.path().join("generated.yaml");
let mut a = args(cfg.clone());
a.output = Some(out.clone());
run(a).await.expect("discover runs");
let text = std::fs::read_to_string(&out).unwrap();
let generated = crate::config::parse_with_extension(&text, "yaml").unwrap();
let nodes = crate::expand::expand(&generated).unwrap();
assert_eq!(nodes.len(), 2, "one row per table: {text}");
assert!(text.contains("orders"), "{text}");
assert!(text.contains("users"), "{text}");
let mut a = args(cfg.clone());
a.output = Some(out.clone());
let err = run(a).await.unwrap_err();
assert!(err.to_string().contains("--force"), "{err}");
let mut a = args(cfg.clone());
a.output = Some(dir.path().join("orders-only.yaml"));
a.include = vec!["orders".into()];
run(a).await.expect("filtered discover runs");
let text = std::fs::read_to_string(dir.path().join("orders-only.yaml")).unwrap();
assert!(
text.contains("orders") && !text.contains("- id: users"),
"{text}"
);
let mut a = args(cfg);
a.json = true;
run(a).await.expect("json mode runs");
}
#[tokio::test]
async fn run_errors_when_filters_match_nothing() {
let dir = tempfile::tempdir().unwrap();
let db = seed_db(dir.path()).await;
let cfg = write_config(dir.path(), &db);
let mut a = args(cfg);
a.include = vec!["nothing_matches_*".into()];
let err = run(a).await.unwrap_err();
assert!(err.to_string().contains("no datasets"), "{err}");
}
}