use crate::config::{ConnectorSpec, PipelineConfig, PipelineSpec, StateStoreSpec, TransformSpec};
use crate::error::{CliError, CliResult};
use serde_json::{Map, Value};
use std::collections::{BTreeMap, HashMap};
pub fn extract_scope(env: &HashMap<String, String>, prefix: &str) -> CliResult<Value> {
let mut object: Map<String, Value> = Map::new();
let mut json_fields: HashMap<String, String> = HashMap::new();
let mut scalar_fields: HashMap<String, String> = HashMap::new();
for (key, value) in env {
let Some(suffix) = key.strip_prefix(prefix) else {
continue;
};
if suffix.is_empty() {
continue;
}
let lowercase = suffix.to_ascii_lowercase();
if let Some(field) = lowercase.strip_suffix("_json") {
if let Some(scalar_var) = scalar_fields.get(field) {
return Err(CliError::EnvConflict {
field: field.to_owned(),
scalar_var: scalar_var.clone(),
json_var: key.clone(),
});
}
let parsed: Value =
serde_json::from_str(value).map_err(|e| CliError::InvalidEnvJson {
var: key.clone(),
message: e.to_string(),
})?;
object.insert(field.to_owned(), parsed);
json_fields.insert(field.to_owned(), key.clone());
} else {
if let Some(json_var) = json_fields.get(&lowercase) {
return Err(CliError::EnvConflict {
field: lowercase.clone(),
scalar_var: key.clone(),
json_var: json_var.clone(),
});
}
object.insert(lowercase.clone(), coerce_scalar(value));
scalar_fields.insert(lowercase, key.clone());
}
}
Ok(Value::Object(object))
}
fn coerce_scalar(s: &str) -> Value {
serde_json::from_str::<Value>(s).unwrap_or_else(|_| Value::String(s.to_owned()))
}
pub fn build_source(env: &HashMap<String, String>) -> CliResult<ConnectorSpec> {
let kind = env
.get("FAUCET_SOURCE")
.filter(|v| !v.is_empty())
.ok_or_else(|| CliError::MissingEnvSelector {
var: "FAUCET_SOURCE".to_owned(),
})?
.clone();
let prefix = format!(
"FAUCET_SOURCE_{}_",
kind.to_ascii_uppercase().replace('-', "_")
);
let config = extract_scope(env, &prefix)?;
Ok(ConnectorSpec {
kind,
config,
transforms: None,
inherit_transforms: true,
status: None,
tags: Vec::new(),
complete_for: None,
})
}
pub fn build_sink(env: &HashMap<String, String>) -> CliResult<ConnectorSpec> {
let kind = env
.get("FAUCET_SINK")
.filter(|v| !v.is_empty())
.ok_or_else(|| CliError::MissingEnvSelector {
var: "FAUCET_SINK".to_owned(),
})?
.clone();
let prefix = format!(
"FAUCET_SINK_{}_",
kind.to_ascii_uppercase().replace('-', "_")
);
let config = extract_scope(env, &prefix)?;
Ok(ConnectorSpec {
kind,
config,
transforms: None,
inherit_transforms: true,
status: None,
tags: Vec::new(),
complete_for: None,
})
}
pub fn build_state(env: &HashMap<String, String>) -> CliResult<Option<StateStoreSpec>> {
let Some(kind) = env.get("FAUCET_STATE").filter(|v| !v.is_empty()).cloned() else {
return Ok(None);
};
let prefix = format!(
"FAUCET_STATE_{}_",
kind.to_ascii_uppercase().replace('-', "_")
);
let config = extract_scope(env, &prefix)?;
Ok(Some(StateStoreSpec { kind, config }))
}
pub fn build_transforms(env: &HashMap<String, String>) -> CliResult<Vec<TransformSpec>> {
let mut kinds: BTreeMap<u32, String> = BTreeMap::new();
for (key, value) in env {
let Some(rest) = key.strip_prefix("FAUCET_TRANSFORM_") else {
continue;
};
if rest.is_empty() || !rest.chars().all(|c| c.is_ascii_digit()) {
continue;
}
let Ok(idx) = rest.parse::<u32>() else {
continue;
};
kinds.insert(idx, value.clone());
}
if kinds.is_empty() {
return Ok(Vec::new());
}
for (expected, actual) in (1u32..).zip(kinds.keys().copied()) {
if expected != actual {
return Err(CliError::TransformIndexGap { missing: expected });
}
}
let mut out = Vec::with_capacity(kinds.len());
for (idx, kind) in kinds {
let prefix = format!("FAUCET_TRANSFORM_{idx}_");
let config = extract_scope(env, &prefix)?;
out.push(TransformSpec { kind, config });
}
Ok(out)
}
pub fn build_named_sources(
env: &HashMap<String, String>,
) -> CliResult<HashMap<String, ConnectorSpec>> {
build_named_catalog(env, "FAUCET_SOURCES_")
}
pub fn build_named_sinks(
env: &HashMap<String, String>,
) -> CliResult<HashMap<String, ConnectorSpec>> {
build_named_catalog(env, "FAUCET_SINKS_")
}
fn build_named_catalog(
env: &HashMap<String, String>,
prefix: &str,
) -> CliResult<HashMap<String, ConnectorSpec>> {
let mut kinds: HashMap<String, String> = HashMap::new();
for (key, value) in env {
let Some(suffix) = key.strip_prefix(prefix) else {
continue;
};
let Some(name_upper) = suffix.strip_suffix("_TYPE") else {
continue;
};
if name_upper.is_empty() {
continue;
}
kinds.insert(name_upper.to_ascii_lowercase(), value.clone());
}
let scope_prefixes: Vec<String> = kinds
.keys()
.map(|name| format!("{prefix}{}_", name.to_ascii_uppercase()))
.collect();
let mut out: HashMap<String, ConnectorSpec> = HashMap::new();
for (name, kind) in kinds {
let scope_prefix = format!("{prefix}{}_", name.to_ascii_uppercase());
let scoped: HashMap<String, String> = env
.iter()
.filter(|(k, _)| {
k.starts_with(&scope_prefix)
&& !scope_prefixes.iter().any(|other| {
other.len() > scope_prefix.len() && k.starts_with(other.as_str())
})
})
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
let mut config = extract_scope(&scoped, &scope_prefix)?;
if let Value::Object(m) = &mut config {
m.remove("type");
}
out.insert(
name,
ConnectorSpec {
kind,
config,
transforms: None,
inherit_transforms: true,
status: None,
tags: Vec::new(),
complete_for: None,
},
);
}
Ok(out)
}
pub fn build_vars(env: &HashMap<String, String>) -> Option<HashMap<String, Value>> {
let mut out: HashMap<String, Value> = HashMap::new();
for (key, value) in env {
let Some(name_upper) = key.strip_prefix("FAUCET_VARS_") else {
continue;
};
if name_upper.is_empty() {
continue;
}
out.insert(name_upper.to_ascii_lowercase(), coerce_scalar(value));
}
if out.is_empty() { None } else { Some(out) }
}
pub fn build_pipeline_config(env: &HashMap<String, String>) -> CliResult<PipelineConfig> {
let source = match env.get("FAUCET_SOURCE").filter(|v| !v.is_empty()) {
Some(_) => Some(build_source(env)?),
None => None,
};
let sink = match env.get("FAUCET_SINK").filter(|v| !v.is_empty()) {
Some(_) => Some(build_sink(env)?),
None => None,
};
let sources = build_named_sources(env)?;
let sinks = build_named_sinks(env)?;
let state = build_state(env)?;
let transforms = build_transforms(env)?;
let vars = build_vars(env);
let name = env.get("FAUCET_NAME").cloned().filter(|s| !s.is_empty());
if source.is_none() && sources.is_empty() {
return Err(CliError::MissingEnvSelector {
var: "FAUCET_SOURCE (or FAUCET_SOURCES_<NAME>_TYPE)".to_owned(),
});
}
if sink.is_none() && sinks.is_empty() {
return Err(CliError::MissingEnvSelector {
var: "FAUCET_SINK (or FAUCET_SINKS_<NAME>_TYPE)".to_owned(),
});
}
Ok(PipelineConfig {
version: 1,
name,
vars,
params: Default::default(),
auth: None,
pipeline: PipelineSpec {
source,
sink,
sources,
sinks,
transforms,
state,
dlq: None,
#[cfg(feature = "quality")]
quality: None,
#[cfg(feature = "contract")]
contract: None,
#[cfg(feature = "masking")]
masking: None,
schema: None,
nodes: std::collections::HashMap::new(),
edges: Vec::new(),
},
matrix: Vec::new(),
execution: None,
selection: None,
observability: None,
delivery: faucet_core::DeliveryMode::default(),
resilience: None,
sla: None,
shard: None,
replication: None,
backfill: None,
partition: None,
#[cfg(feature = "schedule")]
schedule: None,
#[cfg(feature = "lineage")]
lineage: None,
#[cfg(feature = "catalog")]
catalog: None,
#[cfg(feature = "notify")]
notifications: Vec::new(),
})
}
pub fn from_process_env() -> CliResult<PipelineConfig> {
let env: HashMap<String, String> = std::env::vars().collect();
build_pipeline_config(&env)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn env(pairs: &[(&str, &str)]) -> HashMap<String, String> {
pairs
.iter()
.map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
.collect()
}
#[test]
fn extract_scope_lowercases_field_names() {
let e = env(&[("FAUCET_SOURCE_REST_BASE_URL", "https://x.example")]);
let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
assert_eq!(v, json!({"base_url": "https://x.example"}));
}
#[test]
fn extract_scope_ignores_unrelated_keys() {
let e = env(&[
("FAUCET_SOURCE_REST_BASE_URL", "https://x.example"),
("PATH", "/usr/bin"),
("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
]);
let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
assert_eq!(v, json!({"base_url": "https://x.example"}));
}
#[test]
fn named_templates_do_not_leak_across_prefix_overlapping_names() {
let e = env(&[
("FAUCET_SOURCES_USERS_TYPE", "rest"),
("FAUCET_SOURCES_USERS_BASE_URL", "https://u"),
("FAUCET_SOURCES_USERS_API_TYPE", "rest"),
("FAUCET_SOURCES_USERS_API_BASE_URL", "https://api"),
("FAUCET_SOURCES_USERS_API_TIMEOUT", "30"),
]);
let out = build_named_sources(&e).unwrap();
let users = out.get("users").expect("users template").config.clone();
let users = users.as_object().unwrap();
assert_eq!(
users.get("base_url").and_then(|v| v.as_str()),
Some("https://u")
);
assert!(
!users.contains_key("api_base_url"),
"users must not absorb users_api's vars"
);
assert!(!users.contains_key("api_timeout"));
let api = out
.get("users_api")
.expect("users_api template")
.config
.clone();
let api = api.as_object().unwrap();
assert_eq!(
api.get("base_url").and_then(|v| v.as_str()),
Some("https://api")
);
assert_eq!(api.get("timeout").and_then(|v| v.as_i64()), Some(30));
}
#[test]
fn extract_scope_coerces_numbers_and_bools() {
let e = env(&[
("FAUCET_SOURCE_REST_TIMEOUT_SECS", "30"),
("FAUCET_SOURCE_REST_FOLLOW_REDIRECTS", "true"),
("FAUCET_SOURCE_REST_BASE_URL", "https://x.example"),
]);
let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
assert_eq!(v["timeout_secs"], json!(30));
assert_eq!(v["follow_redirects"], json!(true));
assert_eq!(v["base_url"], json!("https://x.example"));
}
#[test]
fn extract_scope_handles_json_suffix() {
let e = env(&[(
"FAUCET_SOURCE_REST_AUTH_JSON",
r#"{"type":"ApiKey","header":"Authorization","value":"Bearer x"}"#,
)]);
let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
assert_eq!(
v["auth"],
json!({"type": "ApiKey", "header": "Authorization", "value": "Bearer x"})
);
}
#[test]
fn extract_scope_rejects_invalid_json_suffix() {
let e = env(&[("FAUCET_SOURCE_REST_AUTH_JSON", "not-json")]);
let err = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap_err();
match err {
CliError::InvalidEnvJson { var, .. } => {
assert_eq!(var, "FAUCET_SOURCE_REST_AUTH_JSON")
}
other => panic!("expected InvalidEnvJson, got {other:?}"),
}
}
#[test]
fn extract_scope_conflict_scalar_then_json() {
let e = env(&[
("FAUCET_SOURCE_REST_AUTH", "bearer"),
("FAUCET_SOURCE_REST_AUTH_JSON", r#"{"type":"ApiKey"}"#),
]);
let err = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap_err();
match err {
CliError::EnvConflict {
field,
scalar_var,
json_var,
} => {
assert_eq!(field, "auth");
assert_eq!(scalar_var, "FAUCET_SOURCE_REST_AUTH");
assert_eq!(json_var, "FAUCET_SOURCE_REST_AUTH_JSON");
}
other => panic!("expected EnvConflict, got {other:?}"),
}
}
#[test]
fn extract_scope_conflict_detection_is_order_independent() {
for _ in 0..50 {
let e = env(&[
("FAUCET_SOURCE_REST_AUTH", "bearer"),
("FAUCET_SOURCE_REST_AUTH_JSON", r#"{"type":"ApiKey"}"#),
]);
let err = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap_err();
match err {
CliError::EnvConflict {
field,
scalar_var,
json_var,
} => {
assert_eq!(field, "auth");
assert_eq!(scalar_var, "FAUCET_SOURCE_REST_AUTH");
assert_eq!(json_var, "FAUCET_SOURCE_REST_AUTH_JSON");
}
other => panic!("expected EnvConflict, got {other:?}"),
}
}
}
#[test]
fn extract_scope_skips_bare_prefix() {
let e = env(&[("FAUCET_SOURCE_REST_", "ignored")]);
let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
assert_eq!(v, json!({}));
}
#[test]
fn extract_scope_empty_when_no_matches() {
let e = env(&[("PATH", "/usr/bin")]);
let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
assert_eq!(v, json!({}));
}
#[test]
fn build_source_reads_selector_and_scope() {
let e = env(&[
("FAUCET_SOURCE", "rest"),
("FAUCET_SOURCE_REST_BASE_URL", "https://x.example"),
("FAUCET_SOURCE_REST_TIMEOUT_SECS", "30"),
]);
let spec = build_source(&e).unwrap();
assert_eq!(spec.kind, "rest");
assert_eq!(spec.config["base_url"], json!("https://x.example"));
assert_eq!(spec.config["timeout_secs"], json!(30));
}
#[test]
fn build_source_uses_kind_scope_so_other_kinds_dont_leak() {
let e = env(&[
("FAUCET_SOURCE", "csv"),
("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
("FAUCET_SOURCE_REST_BASE_URL", "https://other.example"),
]);
let spec = build_source(&e).unwrap();
assert_eq!(spec.kind, "csv");
assert_eq!(spec.config, json!({"path": "./in.csv"}));
}
#[test]
fn build_source_errors_when_selector_missing() {
let e = env(&[("FAUCET_SOURCE_REST_BASE_URL", "https://x.example")]);
let err = build_source(&e).unwrap_err();
match err {
CliError::MissingEnvSelector { var } => assert_eq!(var, "FAUCET_SOURCE"),
other => panic!("expected MissingEnvSelector, got {other:?}"),
}
}
#[test]
fn build_source_errors_when_selector_empty() {
let e = env(&[("FAUCET_SOURCE", "")]);
let err = build_source(&e).unwrap_err();
match err {
CliError::MissingEnvSelector { var } => assert_eq!(var, "FAUCET_SOURCE"),
other => panic!("expected MissingEnvSelector, got {other:?}"),
}
}
#[test]
fn build_sink_reads_selector_and_scope() {
let e = env(&[
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
]);
let spec = build_sink(&e).unwrap();
assert_eq!(spec.kind, "jsonl");
assert_eq!(spec.config, json!({"path": "./out.jsonl"}));
}
#[test]
fn build_sink_errors_when_selector_missing() {
let e = env(&[("FAUCET_SINK_JSONL_PATH", "./out.jsonl")]);
let err = build_sink(&e).unwrap_err();
match err {
CliError::MissingEnvSelector { var } => assert_eq!(var, "FAUCET_SINK"),
other => panic!("expected MissingEnvSelector, got {other:?}"),
}
}
#[test]
fn build_sink_errors_when_selector_empty() {
let e = env(&[("FAUCET_SINK", "")]);
let err = build_sink(&e).unwrap_err();
match err {
CliError::MissingEnvSelector { var } => assert_eq!(var, "FAUCET_SINK"),
other => panic!("expected MissingEnvSelector, got {other:?}"),
}
}
#[test]
fn build_state_returns_none_when_unset() {
let e = env(&[("FAUCET_SOURCE", "rest")]);
let spec = build_state(&e).unwrap();
assert!(spec.is_none());
}
#[test]
fn build_state_returns_none_when_empty() {
let e = env(&[("FAUCET_STATE", "")]);
let spec = build_state(&e).unwrap();
assert!(spec.is_none());
}
#[test]
fn build_state_reads_file_backend() {
let e = env(&[
("FAUCET_STATE", "file"),
("FAUCET_STATE_FILE_PATH", "./.faucet-state"),
]);
let spec = build_state(&e).unwrap().unwrap();
assert_eq!(spec.kind, "file");
assert_eq!(spec.config, json!({"path": "./.faucet-state"}));
}
#[test]
fn build_state_reads_memory_with_empty_scope() {
let e = env(&[("FAUCET_STATE", "memory")]);
let spec = build_state(&e).unwrap().unwrap();
assert_eq!(spec.kind, "memory");
assert_eq!(spec.config, json!({}));
}
#[test]
fn build_transforms_empty_when_unset() {
let e = env(&[("FAUCET_SOURCE", "rest")]);
let t = build_transforms(&e).unwrap();
assert!(t.is_empty());
}
#[test]
fn build_transforms_single_kind_no_config() {
let e = env(&[("FAUCET_TRANSFORM_1", "snake_case")]);
let t = build_transforms(&e).unwrap();
assert_eq!(t.len(), 1);
assert_eq!(t[0].kind, "snake_case");
assert_eq!(t[0].config, json!({}));
}
#[test]
fn build_transforms_ordered_and_with_config() {
let e = env(&[
("FAUCET_TRANSFORM_1", "snake_case"),
("FAUCET_TRANSFORM_2", "flatten"),
("FAUCET_TRANSFORM_2_SEPARATOR", "__"),
]);
let t = build_transforms(&e).unwrap();
assert_eq!(t.len(), 2);
assert_eq!(t[0].kind, "snake_case");
assert_eq!(t[1].kind, "flatten");
assert_eq!(t[1].config, json!({"separator": "__"}));
}
#[test]
fn build_transforms_handles_double_digit_indices() {
let e = env(&[
("FAUCET_TRANSFORM_1", "snake_case"),
("FAUCET_TRANSFORM_2", "flatten"),
("FAUCET_TRANSFORM_3", "rename_keys"),
]);
let t = build_transforms(&e).unwrap();
assert_eq!(t.len(), 3);
assert_eq!(t[2].kind, "rename_keys");
}
#[test]
fn build_transforms_gap_errors() {
let e = env(&[
("FAUCET_TRANSFORM_1", "snake_case"),
("FAUCET_TRANSFORM_3", "flatten"),
]);
let err = build_transforms(&e).unwrap_err();
match err {
CliError::TransformIndexGap { missing } => assert_eq!(missing, 2),
other => panic!("expected TransformIndexGap, got {other:?}"),
}
}
#[test]
fn build_transforms_must_start_at_one() {
let e = env(&[("FAUCET_TRANSFORM_2", "snake_case")]);
let err = build_transforms(&e).unwrap_err();
match err {
CliError::TransformIndexGap { missing } => assert_eq!(missing, 1),
other => panic!("expected TransformIndexGap, got {other:?}"),
}
}
#[test]
fn build_transforms_ignores_field_vars_when_indexing_kinds() {
let e = env(&[("FAUCET_TRANSFORM_1_SEPARATOR", "__")]);
let t = build_transforms(&e).unwrap();
assert!(t.is_empty());
}
#[test]
fn picks_up_named_source_templates() {
let e = env(&[
("FAUCET_SOURCE", "rest"),
("FAUCET_SOURCE_REST_BASE_URL", "https://default.example"),
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./o.jsonl"),
("FAUCET_SOURCES_USERS_API_TYPE", "rest"),
("FAUCET_SOURCES_USERS_API_BASE_URL", "https://users.example"),
("FAUCET_SOURCES_POSTS_API_TYPE", "rest"),
("FAUCET_SOURCES_POSTS_API_BASE_URL", "https://posts.example"),
("FAUCET_SINKS_ARCHIVE_TYPE", "jsonl"),
("FAUCET_SINKS_ARCHIVE_PATH", "./archive.jsonl"),
]);
let cfg = build_pipeline_config(&e).unwrap();
assert!(cfg.pipeline.source.is_some());
assert_eq!(cfg.pipeline.sources.len(), 2);
assert_eq!(cfg.pipeline.sources["users_api"].kind, "rest");
assert_eq!(
cfg.pipeline.sources["users_api"].config["base_url"],
"https://users.example"
);
assert_eq!(cfg.pipeline.sinks["archive"].kind, "jsonl");
}
#[test]
fn picks_up_vars_block() {
let e = env(&[
("FAUCET_SOURCE", "rest"),
("FAUCET_SOURCE_REST_BASE_URL", "https://x.example"),
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./o.jsonl"),
("FAUCET_VARS_API_BASE", "https://api.example.com"),
("FAUCET_VARS_REGION", "us-east-1"),
]);
let cfg = build_pipeline_config(&e).unwrap();
let vars = cfg.vars.unwrap();
assert_eq!(vars["api_base"], "https://api.example.com");
assert_eq!(vars["region"], "us-east-1");
}
#[test]
fn named_source_only_no_legacy_works() {
let e = env(&[
("FAUCET_SOURCES_USERS_API_TYPE", "rest"),
("FAUCET_SOURCES_USERS_API_BASE_URL", "https://x.example"),
("FAUCET_SINKS_ARCHIVE_TYPE", "jsonl"),
("FAUCET_SINKS_ARCHIVE_PATH", "./o.jsonl"),
]);
let cfg = build_pipeline_config(&e).unwrap();
assert!(cfg.pipeline.source.is_none());
assert!(cfg.pipeline.sink.is_none());
assert_eq!(cfg.pipeline.sources["users_api"].kind, "rest");
assert_eq!(cfg.pipeline.sinks["archive"].kind, "jsonl");
}
#[test]
fn no_source_anywhere_errors() {
let e = env(&[
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./o.jsonl"),
]);
let err = build_pipeline_config(&e).unwrap_err();
assert!(matches!(err, CliError::MissingEnvSelector { .. }));
}
#[test]
fn build_pipeline_config_minimal_csv_to_jsonl() {
let e = env(&[
("FAUCET_SOURCE", "csv"),
("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
]);
let cfg = build_pipeline_config(&e).unwrap();
assert_eq!(cfg.version, 1);
assert_eq!(cfg.pipeline.source.as_ref().unwrap().kind, "csv");
assert_eq!(
cfg.pipeline.source.as_ref().unwrap().config,
json!({"path": "./in.csv"})
);
assert_eq!(cfg.pipeline.sink.as_ref().unwrap().kind, "jsonl");
assert_eq!(
cfg.pipeline.sink.as_ref().unwrap().config,
json!({"path": "./out.jsonl"})
);
assert!(cfg.pipeline.transforms.is_empty());
assert!(cfg.pipeline.state.is_none());
assert!(cfg.name.is_none());
}
#[test]
fn build_pipeline_config_uses_faucet_name_when_set() {
let e = env(&[
("FAUCET_NAME", "github-issues"),
("FAUCET_SOURCE", "csv"),
("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
]);
let cfg = build_pipeline_config(&e).unwrap();
assert_eq!(cfg.name.as_deref(), Some("github-issues"));
}
#[test]
fn build_pipeline_config_treats_empty_name_as_none() {
let e = env(&[
("FAUCET_NAME", ""),
("FAUCET_SOURCE", "csv"),
("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
]);
let cfg = build_pipeline_config(&e).unwrap();
assert!(cfg.name.is_none());
}
#[test]
fn build_pipeline_config_with_state_and_transforms() {
let e = env(&[
("FAUCET_SOURCE", "csv"),
("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
("FAUCET_STATE", "file"),
("FAUCET_STATE_FILE_PATH", "./.faucet-state"),
("FAUCET_TRANSFORM_1", "snake_case"),
("FAUCET_TRANSFORM_2", "flatten"),
("FAUCET_TRANSFORM_2_SEPARATOR", "__"),
]);
let cfg = build_pipeline_config(&e).unwrap();
assert_eq!(cfg.pipeline.transforms.len(), 2);
assert_eq!(cfg.pipeline.state.as_ref().unwrap().kind, "file");
}
#[test]
fn build_pipeline_config_missing_source_errors() {
let e = env(&[
("FAUCET_SINK", "jsonl"),
("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
]);
let err = build_pipeline_config(&e).unwrap_err();
assert!(matches!(err, CliError::MissingEnvSelector { .. }));
}
#[test]
fn build_pipeline_config_missing_sink_errors() {
let e = env(&[
("FAUCET_SOURCE", "csv"),
("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
]);
let err = build_pipeline_config(&e).unwrap_err();
assert!(matches!(err, CliError::MissingEnvSelector { .. }));
}
}