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