1use crate::config::{ConnectorSpec, PipelineConfig, PipelineSpec, StateStoreSpec, TransformSpec};
10use crate::error::{CliError, CliResult};
11use serde_json::{Map, Value};
12use std::collections::{BTreeMap, HashMap};
13
14pub fn extract_scope(env: &HashMap<String, String>, prefix: &str) -> CliResult<Value> {
22 let mut object: Map<String, Value> = Map::new();
23 let mut json_fields: HashMap<String, String> = HashMap::new();
25 let mut scalar_fields: HashMap<String, String> = HashMap::new();
26
27 for (key, value) in env {
28 let Some(suffix) = key.strip_prefix(prefix) else {
29 continue;
30 };
31 if suffix.is_empty() {
32 continue;
34 }
35 let lowercase = suffix.to_ascii_lowercase();
36 if let Some(field) = lowercase.strip_suffix("_json") {
37 if let Some(scalar_var) = scalar_fields.get(field) {
38 return Err(CliError::EnvConflict {
39 field: field.to_owned(),
40 scalar_var: scalar_var.clone(),
41 json_var: key.clone(),
42 });
43 }
44 let parsed: Value =
45 serde_json::from_str(value).map_err(|e| CliError::InvalidEnvJson {
46 var: key.clone(),
47 message: e.to_string(),
48 })?;
49 object.insert(field.to_owned(), parsed);
50 json_fields.insert(field.to_owned(), key.clone());
51 } else {
52 if let Some(json_var) = json_fields.get(&lowercase) {
53 return Err(CliError::EnvConflict {
54 field: lowercase.clone(),
55 scalar_var: key.clone(),
56 json_var: json_var.clone(),
57 });
58 }
59 object.insert(lowercase.clone(), coerce_scalar(value));
60 scalar_fields.insert(lowercase, key.clone());
61 }
62 }
63 Ok(Value::Object(object))
64}
65
66fn coerce_scalar(s: &str) -> Value {
80 serde_json::from_str::<Value>(s).unwrap_or_else(|_| Value::String(s.to_owned()))
81}
82
83pub fn build_source(env: &HashMap<String, String>) -> CliResult<ConnectorSpec> {
85 let kind = env
86 .get("FAUCET_SOURCE")
87 .filter(|v| !v.is_empty())
88 .ok_or_else(|| CliError::MissingEnvSelector {
89 var: "FAUCET_SOURCE".to_owned(),
90 })?
91 .clone();
92 let prefix = format!(
93 "FAUCET_SOURCE_{}_",
94 kind.to_ascii_uppercase().replace('-', "_")
95 );
96 let config = extract_scope(env, &prefix)?;
97 Ok(ConnectorSpec {
98 kind,
99 config,
100 transforms: None,
101 inherit_transforms: true,
102 status: None,
103 tags: Vec::new(),
104 complete_for: None,
105 })
106}
107
108pub fn build_sink(env: &HashMap<String, String>) -> CliResult<ConnectorSpec> {
110 let kind = env
111 .get("FAUCET_SINK")
112 .filter(|v| !v.is_empty())
113 .ok_or_else(|| CliError::MissingEnvSelector {
114 var: "FAUCET_SINK".to_owned(),
115 })?
116 .clone();
117 let prefix = format!(
118 "FAUCET_SINK_{}_",
119 kind.to_ascii_uppercase().replace('-', "_")
120 );
121 let config = extract_scope(env, &prefix)?;
122 Ok(ConnectorSpec {
123 kind,
124 config,
125 transforms: None,
126 inherit_transforms: true,
127 status: None,
128 tags: Vec::new(),
129 complete_for: None,
130 })
131}
132
133pub fn build_state(env: &HashMap<String, String>) -> CliResult<Option<StateStoreSpec>> {
136 let Some(kind) = env.get("FAUCET_STATE").filter(|v| !v.is_empty()).cloned() else {
137 return Ok(None);
138 };
139 let prefix = format!(
140 "FAUCET_STATE_{}_",
141 kind.to_ascii_uppercase().replace('-', "_")
142 );
143 let config = extract_scope(env, &prefix)?;
144 Ok(Some(StateStoreSpec { kind, config }))
145}
146
147pub fn build_transforms(env: &HashMap<String, String>) -> CliResult<Vec<TransformSpec>> {
152 let mut kinds: BTreeMap<u32, String> = BTreeMap::new();
155 for (key, value) in env {
156 let Some(rest) = key.strip_prefix("FAUCET_TRANSFORM_") else {
157 continue;
158 };
159 if rest.is_empty() || !rest.chars().all(|c| c.is_ascii_digit()) {
160 continue;
161 }
162 let Ok(idx) = rest.parse::<u32>() else {
163 continue;
164 };
165 kinds.insert(idx, value.clone());
166 }
167 if kinds.is_empty() {
168 return Ok(Vec::new());
169 }
170 for (expected, actual) in (1u32..).zip(kinds.keys().copied()) {
172 if expected != actual {
173 return Err(CliError::TransformIndexGap { missing: expected });
174 }
175 }
176 let mut out = Vec::with_capacity(kinds.len());
178 for (idx, kind) in kinds {
179 let prefix = format!("FAUCET_TRANSFORM_{idx}_");
180 let config = extract_scope(env, &prefix)?;
181 out.push(TransformSpec { kind, config });
182 }
183 Ok(out)
184}
185
186pub fn build_named_sources(
189 env: &HashMap<String, String>,
190) -> CliResult<HashMap<String, ConnectorSpec>> {
191 build_named_catalog(env, "FAUCET_SOURCES_")
192}
193
194pub fn build_named_sinks(
196 env: &HashMap<String, String>,
197) -> CliResult<HashMap<String, ConnectorSpec>> {
198 build_named_catalog(env, "FAUCET_SINKS_")
199}
200
201fn build_named_catalog(
202 env: &HashMap<String, String>,
203 prefix: &str,
204) -> CliResult<HashMap<String, ConnectorSpec>> {
205 let mut kinds: HashMap<String, String> = HashMap::new();
208 for (key, value) in env {
209 let Some(suffix) = key.strip_prefix(prefix) else {
210 continue;
211 };
212 let Some(name_upper) = suffix.strip_suffix("_TYPE") else {
213 continue;
214 };
215 if name_upper.is_empty() {
216 continue;
217 }
218 kinds.insert(name_upper.to_ascii_lowercase(), value.clone());
219 }
220 let scope_prefixes: Vec<String> = kinds
226 .keys()
227 .map(|name| format!("{prefix}{}_", name.to_ascii_uppercase()))
228 .collect();
229
230 let mut out: HashMap<String, ConnectorSpec> = HashMap::new();
232 for (name, kind) in kinds {
233 let scope_prefix = format!("{prefix}{}_", name.to_ascii_uppercase());
234 let scoped: HashMap<String, String> = env
238 .iter()
239 .filter(|(k, _)| {
240 k.starts_with(&scope_prefix)
241 && !scope_prefixes.iter().any(|other| {
242 other.len() > scope_prefix.len() && k.starts_with(other.as_str())
243 })
244 })
245 .map(|(k, v)| (k.clone(), v.clone()))
246 .collect();
247 let mut config = extract_scope(&scoped, &scope_prefix)?;
248 if let Value::Object(m) = &mut config {
250 m.remove("type");
251 }
252 out.insert(
253 name,
254 ConnectorSpec {
255 kind,
256 config,
257 transforms: None,
258 inherit_transforms: true,
259 status: None,
260 tags: Vec::new(),
261 complete_for: None,
262 },
263 );
264 }
265 Ok(out)
266}
267
268pub fn build_vars(env: &HashMap<String, String>) -> Option<HashMap<String, Value>> {
271 let mut out: HashMap<String, Value> = HashMap::new();
272 for (key, value) in env {
273 let Some(name_upper) = key.strip_prefix("FAUCET_VARS_") else {
274 continue;
275 };
276 if name_upper.is_empty() {
277 continue;
278 }
279 out.insert(name_upper.to_ascii_lowercase(), coerce_scalar(value));
280 }
281 if out.is_empty() { None } else { Some(out) }
282}
283
284pub fn build_pipeline_config(env: &HashMap<String, String>) -> CliResult<PipelineConfig> {
286 let source = match env.get("FAUCET_SOURCE").filter(|v| !v.is_empty()) {
287 Some(_) => Some(build_source(env)?),
288 None => None,
289 };
290 let sink = match env.get("FAUCET_SINK").filter(|v| !v.is_empty()) {
291 Some(_) => Some(build_sink(env)?),
292 None => None,
293 };
294 let sources = build_named_sources(env)?;
295 let sinks = build_named_sinks(env)?;
296 let state = build_state(env)?;
297 let transforms = build_transforms(env)?;
298 let vars = build_vars(env);
299 let name = env.get("FAUCET_NAME").cloned().filter(|s| !s.is_empty());
300
301 if source.is_none() && sources.is_empty() {
303 return Err(CliError::MissingEnvSelector {
304 var: "FAUCET_SOURCE (or FAUCET_SOURCES_<NAME>_TYPE)".to_owned(),
305 });
306 }
307 if sink.is_none() && sinks.is_empty() {
308 return Err(CliError::MissingEnvSelector {
309 var: "FAUCET_SINK (or FAUCET_SINKS_<NAME>_TYPE)".to_owned(),
310 });
311 }
312 Ok(PipelineConfig {
313 version: 1,
314 name,
315 vars,
316 params: Default::default(),
319 auth: None,
322 pipeline: PipelineSpec {
323 source,
324 sink,
325 sources,
326 sinks,
327 transforms,
328 state,
329 dlq: None,
330 #[cfg(feature = "quality")]
331 quality: None,
332 #[cfg(feature = "contract")]
333 contract: None,
334 #[cfg(feature = "masking")]
335 masking: None,
336 schema: None,
337 nodes: std::collections::HashMap::new(),
339 edges: Vec::new(),
340 },
341 matrix: Vec::new(),
342 execution: None,
343 selection: None,
344 observability: None,
345 delivery: faucet_core::DeliveryMode::default(),
346 resilience: None,
347 sla: None,
349 reconcile: None,
350 shard: None,
351 replication: None,
352 backfill: None,
353 metadata_columns: None,
354 partition: None,
355 #[cfg(feature = "schedule")]
356 schedule: None,
357 #[cfg(feature = "lineage")]
358 lineage: None,
359 #[cfg(feature = "catalog")]
361 catalog: None,
362 #[cfg(feature = "notify")]
363 notifications: Vec::new(),
364 })
365}
366
367pub fn from_process_env() -> CliResult<PipelineConfig> {
369 let env: HashMap<String, String> = std::env::vars().collect();
370 build_pipeline_config(&env)
371}
372
373#[cfg(test)]
374mod tests {
375 use super::*;
376 use serde_json::json;
377
378 fn env(pairs: &[(&str, &str)]) -> HashMap<String, String> {
379 pairs
380 .iter()
381 .map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
382 .collect()
383 }
384
385 #[test]
386 fn extract_scope_lowercases_field_names() {
387 let e = env(&[("FAUCET_SOURCE_REST_BASE_URL", "https://x.example")]);
388 let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
389 assert_eq!(v, json!({"base_url": "https://x.example"}));
390 }
391
392 #[test]
393 fn extract_scope_ignores_unrelated_keys() {
394 let e = env(&[
395 ("FAUCET_SOURCE_REST_BASE_URL", "https://x.example"),
396 ("PATH", "/usr/bin"),
397 ("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
398 ]);
399 let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
400 assert_eq!(v, json!({"base_url": "https://x.example"}));
401 }
402
403 #[test]
404 fn named_templates_do_not_leak_across_prefix_overlapping_names() {
405 let e = env(&[
408 ("FAUCET_SOURCES_USERS_TYPE", "rest"),
409 ("FAUCET_SOURCES_USERS_BASE_URL", "https://u"),
410 ("FAUCET_SOURCES_USERS_API_TYPE", "rest"),
411 ("FAUCET_SOURCES_USERS_API_BASE_URL", "https://api"),
412 ("FAUCET_SOURCES_USERS_API_TIMEOUT", "30"),
413 ]);
414 let out = build_named_sources(&e).unwrap();
415
416 let users = out.get("users").expect("users template").config.clone();
417 let users = users.as_object().unwrap();
418 assert_eq!(
419 users.get("base_url").and_then(|v| v.as_str()),
420 Some("https://u")
421 );
422 assert!(
423 !users.contains_key("api_base_url"),
424 "users must not absorb users_api's vars"
425 );
426 assert!(!users.contains_key("api_timeout"));
427
428 let api = out
429 .get("users_api")
430 .expect("users_api template")
431 .config
432 .clone();
433 let api = api.as_object().unwrap();
434 assert_eq!(
435 api.get("base_url").and_then(|v| v.as_str()),
436 Some("https://api")
437 );
438 assert_eq!(api.get("timeout").and_then(|v| v.as_i64()), Some(30));
439 }
440
441 #[test]
442 fn extract_scope_coerces_numbers_and_bools() {
443 let e = env(&[
444 ("FAUCET_SOURCE_REST_TIMEOUT_SECS", "30"),
445 ("FAUCET_SOURCE_REST_FOLLOW_REDIRECTS", "true"),
446 ("FAUCET_SOURCE_REST_BASE_URL", "https://x.example"),
447 ]);
448 let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
449 assert_eq!(v["timeout_secs"], json!(30));
450 assert_eq!(v["follow_redirects"], json!(true));
451 assert_eq!(v["base_url"], json!("https://x.example"));
452 }
453
454 #[test]
455 fn extract_scope_handles_json_suffix() {
456 let e = env(&[(
457 "FAUCET_SOURCE_REST_AUTH_JSON",
458 r#"{"type":"ApiKey","header":"Authorization","value":"Bearer x"}"#,
459 )]);
460 let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
461 assert_eq!(
462 v["auth"],
463 json!({"type": "ApiKey", "header": "Authorization", "value": "Bearer x"})
464 );
465 }
466
467 #[test]
468 fn extract_scope_rejects_invalid_json_suffix() {
469 let e = env(&[("FAUCET_SOURCE_REST_AUTH_JSON", "not-json")]);
470 let err = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap_err();
471 match err {
472 CliError::InvalidEnvJson { var, .. } => {
473 assert_eq!(var, "FAUCET_SOURCE_REST_AUTH_JSON")
474 }
475 other => panic!("expected InvalidEnvJson, got {other:?}"),
476 }
477 }
478
479 #[test]
480 fn extract_scope_conflict_scalar_then_json() {
481 let e = env(&[
482 ("FAUCET_SOURCE_REST_AUTH", "bearer"),
483 ("FAUCET_SOURCE_REST_AUTH_JSON", r#"{"type":"ApiKey"}"#),
484 ]);
485 let err = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap_err();
486 match err {
487 CliError::EnvConflict {
488 field,
489 scalar_var,
490 json_var,
491 } => {
492 assert_eq!(field, "auth");
493 assert_eq!(scalar_var, "FAUCET_SOURCE_REST_AUTH");
494 assert_eq!(json_var, "FAUCET_SOURCE_REST_AUTH_JSON");
495 }
496 other => panic!("expected EnvConflict, got {other:?}"),
497 }
498 }
499
500 #[test]
501 fn extract_scope_conflict_detection_is_order_independent() {
502 for _ in 0..50 {
506 let e = env(&[
507 ("FAUCET_SOURCE_REST_AUTH", "bearer"),
508 ("FAUCET_SOURCE_REST_AUTH_JSON", r#"{"type":"ApiKey"}"#),
509 ]);
510 let err = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap_err();
511 match err {
512 CliError::EnvConflict {
513 field,
514 scalar_var,
515 json_var,
516 } => {
517 assert_eq!(field, "auth");
518 assert_eq!(scalar_var, "FAUCET_SOURCE_REST_AUTH");
519 assert_eq!(json_var, "FAUCET_SOURCE_REST_AUTH_JSON");
520 }
521 other => panic!("expected EnvConflict, got {other:?}"),
522 }
523 }
524 }
525
526 #[test]
527 fn extract_scope_skips_bare_prefix() {
528 let e = env(&[("FAUCET_SOURCE_REST_", "ignored")]);
529 let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
530 assert_eq!(v, json!({}));
531 }
532
533 #[test]
534 fn extract_scope_empty_when_no_matches() {
535 let e = env(&[("PATH", "/usr/bin")]);
536 let v = extract_scope(&e, "FAUCET_SOURCE_REST_").unwrap();
537 assert_eq!(v, json!({}));
538 }
539
540 #[test]
541 fn build_source_reads_selector_and_scope() {
542 let e = env(&[
543 ("FAUCET_SOURCE", "rest"),
544 ("FAUCET_SOURCE_REST_BASE_URL", "https://x.example"),
545 ("FAUCET_SOURCE_REST_TIMEOUT_SECS", "30"),
546 ]);
547 let spec = build_source(&e).unwrap();
548 assert_eq!(spec.kind, "rest");
549 assert_eq!(spec.config["base_url"], json!("https://x.example"));
550 assert_eq!(spec.config["timeout_secs"], json!(30));
551 }
552
553 #[test]
554 fn build_source_uses_kind_scope_so_other_kinds_dont_leak() {
555 let e = env(&[
556 ("FAUCET_SOURCE", "csv"),
557 ("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
558 ("FAUCET_SOURCE_REST_BASE_URL", "https://other.example"),
559 ]);
560 let spec = build_source(&e).unwrap();
561 assert_eq!(spec.kind, "csv");
562 assert_eq!(spec.config, json!({"path": "./in.csv"}));
563 }
564
565 #[test]
566 fn build_source_errors_when_selector_missing() {
567 let e = env(&[("FAUCET_SOURCE_REST_BASE_URL", "https://x.example")]);
568 let err = build_source(&e).unwrap_err();
569 match err {
570 CliError::MissingEnvSelector { var } => assert_eq!(var, "FAUCET_SOURCE"),
571 other => panic!("expected MissingEnvSelector, got {other:?}"),
572 }
573 }
574
575 #[test]
576 fn build_source_errors_when_selector_empty() {
577 let e = env(&[("FAUCET_SOURCE", "")]);
578 let err = build_source(&e).unwrap_err();
579 match err {
580 CliError::MissingEnvSelector { var } => assert_eq!(var, "FAUCET_SOURCE"),
581 other => panic!("expected MissingEnvSelector, got {other:?}"),
582 }
583 }
584
585 #[test]
586 fn build_sink_reads_selector_and_scope() {
587 let e = env(&[
588 ("FAUCET_SINK", "jsonl"),
589 ("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
590 ]);
591 let spec = build_sink(&e).unwrap();
592 assert_eq!(spec.kind, "jsonl");
593 assert_eq!(spec.config, json!({"path": "./out.jsonl"}));
594 }
595
596 #[test]
597 fn build_sink_errors_when_selector_missing() {
598 let e = env(&[("FAUCET_SINK_JSONL_PATH", "./out.jsonl")]);
599 let err = build_sink(&e).unwrap_err();
600 match err {
601 CliError::MissingEnvSelector { var } => assert_eq!(var, "FAUCET_SINK"),
602 other => panic!("expected MissingEnvSelector, got {other:?}"),
603 }
604 }
605
606 #[test]
607 fn build_sink_errors_when_selector_empty() {
608 let e = env(&[("FAUCET_SINK", "")]);
609 let err = build_sink(&e).unwrap_err();
610 match err {
611 CliError::MissingEnvSelector { var } => assert_eq!(var, "FAUCET_SINK"),
612 other => panic!("expected MissingEnvSelector, got {other:?}"),
613 }
614 }
615
616 #[test]
617 fn build_state_returns_none_when_unset() {
618 let e = env(&[("FAUCET_SOURCE", "rest")]);
619 let spec = build_state(&e).unwrap();
620 assert!(spec.is_none());
621 }
622
623 #[test]
624 fn build_state_returns_none_when_empty() {
625 let e = env(&[("FAUCET_STATE", "")]);
626 let spec = build_state(&e).unwrap();
627 assert!(spec.is_none());
628 }
629
630 #[test]
631 fn build_state_reads_file_backend() {
632 let e = env(&[
633 ("FAUCET_STATE", "file"),
634 ("FAUCET_STATE_FILE_PATH", "./.faucet-state"),
635 ]);
636 let spec = build_state(&e).unwrap().unwrap();
637 assert_eq!(spec.kind, "file");
638 assert_eq!(spec.config, json!({"path": "./.faucet-state"}));
639 }
640
641 #[test]
642 fn build_state_reads_memory_with_empty_scope() {
643 let e = env(&[("FAUCET_STATE", "memory")]);
644 let spec = build_state(&e).unwrap().unwrap();
645 assert_eq!(spec.kind, "memory");
646 assert_eq!(spec.config, json!({}));
647 }
648
649 #[test]
650 fn build_transforms_empty_when_unset() {
651 let e = env(&[("FAUCET_SOURCE", "rest")]);
652 let t = build_transforms(&e).unwrap();
653 assert!(t.is_empty());
654 }
655
656 #[test]
657 fn build_transforms_single_kind_no_config() {
658 let e = env(&[("FAUCET_TRANSFORM_1", "snake_case")]);
659 let t = build_transforms(&e).unwrap();
660 assert_eq!(t.len(), 1);
661 assert_eq!(t[0].kind, "snake_case");
662 assert_eq!(t[0].config, json!({}));
663 }
664
665 #[test]
666 fn build_transforms_ordered_and_with_config() {
667 let e = env(&[
668 ("FAUCET_TRANSFORM_1", "snake_case"),
669 ("FAUCET_TRANSFORM_2", "flatten"),
670 ("FAUCET_TRANSFORM_2_SEPARATOR", "__"),
671 ]);
672 let t = build_transforms(&e).unwrap();
673 assert_eq!(t.len(), 2);
674 assert_eq!(t[0].kind, "snake_case");
675 assert_eq!(t[1].kind, "flatten");
676 assert_eq!(t[1].config, json!({"separator": "__"}));
677 }
678
679 #[test]
680 fn build_transforms_handles_double_digit_indices() {
681 let e = env(&[
682 ("FAUCET_TRANSFORM_1", "snake_case"),
683 ("FAUCET_TRANSFORM_2", "flatten"),
684 ("FAUCET_TRANSFORM_3", "rename_keys"),
685 ]);
686 let t = build_transforms(&e).unwrap();
687 assert_eq!(t.len(), 3);
688 assert_eq!(t[2].kind, "rename_keys");
689 }
690
691 #[test]
692 fn build_transforms_gap_errors() {
693 let e = env(&[
694 ("FAUCET_TRANSFORM_1", "snake_case"),
695 ("FAUCET_TRANSFORM_3", "flatten"),
696 ]);
697 let err = build_transforms(&e).unwrap_err();
698 match err {
699 CliError::TransformIndexGap { missing } => assert_eq!(missing, 2),
700 other => panic!("expected TransformIndexGap, got {other:?}"),
701 }
702 }
703
704 #[test]
705 fn build_transforms_must_start_at_one() {
706 let e = env(&[("FAUCET_TRANSFORM_2", "snake_case")]);
707 let err = build_transforms(&e).unwrap_err();
708 match err {
709 CliError::TransformIndexGap { missing } => assert_eq!(missing, 1),
710 other => panic!("expected TransformIndexGap, got {other:?}"),
711 }
712 }
713
714 #[test]
715 fn build_transforms_ignores_field_vars_when_indexing_kinds() {
716 let e = env(&[("FAUCET_TRANSFORM_1_SEPARATOR", "__")]);
718 let t = build_transforms(&e).unwrap();
721 assert!(t.is_empty());
722 }
723
724 #[test]
725 fn picks_up_named_source_templates() {
726 let e = env(&[
727 ("FAUCET_SOURCE", "rest"),
729 ("FAUCET_SOURCE_REST_BASE_URL", "https://default.example"),
730 ("FAUCET_SINK", "jsonl"),
731 ("FAUCET_SINK_JSONL_PATH", "./o.jsonl"),
732 ("FAUCET_SOURCES_USERS_API_TYPE", "rest"),
734 ("FAUCET_SOURCES_USERS_API_BASE_URL", "https://users.example"),
735 ("FAUCET_SOURCES_POSTS_API_TYPE", "rest"),
736 ("FAUCET_SOURCES_POSTS_API_BASE_URL", "https://posts.example"),
737 ("FAUCET_SINKS_ARCHIVE_TYPE", "jsonl"),
738 ("FAUCET_SINKS_ARCHIVE_PATH", "./archive.jsonl"),
739 ]);
740 let cfg = build_pipeline_config(&e).unwrap();
741 assert!(cfg.pipeline.source.is_some());
742 assert_eq!(cfg.pipeline.sources.len(), 2);
743 assert_eq!(cfg.pipeline.sources["users_api"].kind, "rest");
744 assert_eq!(
745 cfg.pipeline.sources["users_api"].config["base_url"],
746 "https://users.example"
747 );
748 assert_eq!(cfg.pipeline.sinks["archive"].kind, "jsonl");
749 }
750
751 #[test]
752 fn picks_up_vars_block() {
753 let e = env(&[
754 ("FAUCET_SOURCE", "rest"),
755 ("FAUCET_SOURCE_REST_BASE_URL", "https://x.example"),
756 ("FAUCET_SINK", "jsonl"),
757 ("FAUCET_SINK_JSONL_PATH", "./o.jsonl"),
758 ("FAUCET_VARS_API_BASE", "https://api.example.com"),
759 ("FAUCET_VARS_REGION", "us-east-1"),
760 ]);
761 let cfg = build_pipeline_config(&e).unwrap();
762 let vars = cfg.vars.unwrap();
763 assert_eq!(vars["api_base"], "https://api.example.com");
764 assert_eq!(vars["region"], "us-east-1");
765 }
766
767 #[test]
768 fn named_source_only_no_legacy_works() {
769 let e = env(&[
773 ("FAUCET_SOURCES_USERS_API_TYPE", "rest"),
774 ("FAUCET_SOURCES_USERS_API_BASE_URL", "https://x.example"),
775 ("FAUCET_SINKS_ARCHIVE_TYPE", "jsonl"),
776 ("FAUCET_SINKS_ARCHIVE_PATH", "./o.jsonl"),
777 ]);
778 let cfg = build_pipeline_config(&e).unwrap();
779 assert!(cfg.pipeline.source.is_none());
780 assert!(cfg.pipeline.sink.is_none());
781 assert_eq!(cfg.pipeline.sources["users_api"].kind, "rest");
782 assert_eq!(cfg.pipeline.sinks["archive"].kind, "jsonl");
783 }
784
785 #[test]
786 fn no_source_anywhere_errors() {
787 let e = env(&[
789 ("FAUCET_SINK", "jsonl"),
790 ("FAUCET_SINK_JSONL_PATH", "./o.jsonl"),
791 ]);
792 let err = build_pipeline_config(&e).unwrap_err();
793 assert!(matches!(err, CliError::MissingEnvSelector { .. }));
794 }
795
796 #[test]
797 fn build_pipeline_config_minimal_csv_to_jsonl() {
798 let e = env(&[
799 ("FAUCET_SOURCE", "csv"),
800 ("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
801 ("FAUCET_SINK", "jsonl"),
802 ("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
803 ]);
804 let cfg = build_pipeline_config(&e).unwrap();
805 assert_eq!(cfg.version, 1);
806 assert_eq!(cfg.pipeline.source.as_ref().unwrap().kind, "csv");
807 assert_eq!(
808 cfg.pipeline.source.as_ref().unwrap().config,
809 json!({"path": "./in.csv"})
810 );
811 assert_eq!(cfg.pipeline.sink.as_ref().unwrap().kind, "jsonl");
812 assert_eq!(
813 cfg.pipeline.sink.as_ref().unwrap().config,
814 json!({"path": "./out.jsonl"})
815 );
816 assert!(cfg.pipeline.transforms.is_empty());
817 assert!(cfg.pipeline.state.is_none());
818 assert!(cfg.name.is_none());
819 }
820
821 #[test]
822 fn build_pipeline_config_uses_faucet_name_when_set() {
823 let e = env(&[
824 ("FAUCET_NAME", "github-issues"),
825 ("FAUCET_SOURCE", "csv"),
826 ("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
827 ("FAUCET_SINK", "jsonl"),
828 ("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
829 ]);
830 let cfg = build_pipeline_config(&e).unwrap();
831 assert_eq!(cfg.name.as_deref(), Some("github-issues"));
832 }
833
834 #[test]
835 fn build_pipeline_config_treats_empty_name_as_none() {
836 let e = env(&[
837 ("FAUCET_NAME", ""),
838 ("FAUCET_SOURCE", "csv"),
839 ("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
840 ("FAUCET_SINK", "jsonl"),
841 ("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
842 ]);
843 let cfg = build_pipeline_config(&e).unwrap();
844 assert!(cfg.name.is_none());
845 }
846
847 #[test]
848 fn build_pipeline_config_with_state_and_transforms() {
849 let e = env(&[
850 ("FAUCET_SOURCE", "csv"),
851 ("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
852 ("FAUCET_SINK", "jsonl"),
853 ("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
854 ("FAUCET_STATE", "file"),
855 ("FAUCET_STATE_FILE_PATH", "./.faucet-state"),
856 ("FAUCET_TRANSFORM_1", "snake_case"),
857 ("FAUCET_TRANSFORM_2", "flatten"),
858 ("FAUCET_TRANSFORM_2_SEPARATOR", "__"),
859 ]);
860 let cfg = build_pipeline_config(&e).unwrap();
861 assert_eq!(cfg.pipeline.transforms.len(), 2);
862 assert_eq!(cfg.pipeline.state.as_ref().unwrap().kind, "file");
863 }
864
865 #[test]
866 fn build_pipeline_config_missing_source_errors() {
867 let e = env(&[
868 ("FAUCET_SINK", "jsonl"),
869 ("FAUCET_SINK_JSONL_PATH", "./out.jsonl"),
870 ]);
871 let err = build_pipeline_config(&e).unwrap_err();
872 assert!(matches!(err, CliError::MissingEnvSelector { .. }));
873 }
874
875 #[test]
876 fn build_pipeline_config_missing_sink_errors() {
877 let e = env(&[
878 ("FAUCET_SOURCE", "csv"),
879 ("FAUCET_SOURCE_CSV_PATH", "./in.csv"),
880 ]);
881 let err = build_pipeline_config(&e).unwrap_err();
882 assert!(matches!(err, CliError::MissingEnvSelector { .. }));
883 }
884}