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