1use std::{
2 collections::{BTreeMap, BTreeSet},
3 fmt,
4};
5
6use serde_json::{Map, Value};
7
8use super::{
9 error::{ConfigError, ConfigResult},
10 resolved::ResolvedConfig,
11 secret::{ResolveOptions, SecretRef, SecretResolver, SecretString},
12};
13
14pub type RawConfig = Map<String, Value>;
16
17#[derive(Debug, Clone, Copy, Eq, PartialEq)]
19pub enum ConfigKind {
20 String,
21 Bool,
22 Integer,
23 UnsignedInteger,
24 Float,
25 DurationMillis,
26 DurationMs,
27 StringList,
28 Json,
29 Secret,
30}
31
32impl ConfigKind {
33 #[must_use]
35 pub fn expected(self) -> &'static str {
36 match self {
37 Self::String => "string",
38 Self::Bool => "bool",
39 Self::Integer => "integer",
40 Self::UnsignedInteger => "unsigned integer",
41 Self::Float => "float",
42 Self::DurationMillis | Self::DurationMs => "duration milliseconds",
43 Self::StringList => "string list",
44 Self::Json => "json value",
45 Self::Secret => "secret reference",
46 }
47 }
48}
49
50#[non_exhaustive]
52#[derive(Clone, Eq, PartialEq)]
53pub struct ConfigKey {
54 pub name: String,
55 pub kind: ConfigKind,
56 pub required: bool,
57 pub default: Option<Value>,
58 pub description: Option<String>,
59}
60
61#[derive(Clone, Default, Eq, PartialEq)]
63pub struct ConfigDef {
64 name: String,
65 keys: BTreeMap<String, ConfigKey>,
66}
67
68impl fmt::Debug for ConfigKey {
69 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
70 let redacted_default = if self.kind == ConfigKind::Secret && self.default.is_some() {
71 Some(Value::String("<redacted>".to_owned()))
72 } else {
73 self.default.clone()
74 };
75
76 f.debug_struct("ConfigKey")
77 .field("name", &self.name)
78 .field("kind", &self.kind)
79 .field("required", &self.required)
80 .field("default", &redacted_default)
81 .field("description", &self.description)
82 .finish()
83 }
84}
85
86impl fmt::Debug for ConfigDef {
87 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
88 f.debug_struct("ConfigDef")
89 .field("name", &self.name)
90 .field("keys", &self.keys)
91 .finish()
92 }
93}
94
95impl ConfigDef {
96 #[must_use]
98 pub fn new(name: impl Into<String>) -> Self {
99 Self {
100 name: name.into(),
101 keys: BTreeMap::new(),
102 }
103 }
104
105 #[must_use]
107 pub fn name(&self) -> &str {
108 &self.name
109 }
110
111 pub fn keys(&self) -> impl Iterator<Item = &ConfigKey> {
113 self.keys.values()
114 }
115
116 #[must_use]
118 pub fn required(mut self, name: impl Into<String>, kind: ConfigKind) -> Self {
119 let name = name.into();
120 self.define_key(ConfigKey {
121 name,
122 kind,
123 required: true,
124 default: None,
125 description: None,
126 });
127 self
128 }
129
130 #[must_use]
132 pub fn optional(mut self, name: impl Into<String>, kind: ConfigKind) -> Self {
133 let name = name.into();
134 self.define_key(ConfigKey {
135 name,
136 kind,
137 required: false,
138 default: None,
139 description: None,
140 });
141 self
142 }
143
144 #[must_use]
146 pub fn default(
147 mut self,
148 name: impl Into<String>,
149 kind: ConfigKind,
150 default: impl Into<Value>,
151 ) -> Self {
152 let name = name.into();
153 assert!(
154 kind != ConfigKind::Secret,
155 "secret connector config key `{name}` cannot have a default"
156 );
157 self.define_key(ConfigKey {
158 name,
159 kind,
160 required: false,
161 default: Some(default.into()),
162 description: None,
163 });
164 self
165 }
166
167 #[must_use]
169 pub fn secret(self, name: impl Into<String>) -> Self {
170 self.required(name, ConfigKind::Secret)
171 }
172
173 fn define_key(&mut self, key: ConfigKey) {
174 let name = key.name.clone();
175 assert!(
176 self.keys.insert(name.clone(), key).is_none(),
177 "duplicate connector config key `{name}`"
178 );
179 }
180
181 pub async fn resolve(
183 &self,
184 raw: RawConfig,
185 resolver: &dyn SecretResolver,
186 ) -> ConfigResult<ResolvedConfig> {
187 self.resolve_with_options(raw, resolver, ResolveOptions::default())
188 .await
189 }
190
191 pub async fn resolve_with_options(
193 &self,
194 raw: RawConfig,
195 resolver: &dyn SecretResolver,
196 options: ResolveOptions,
197 ) -> ConfigResult<ResolvedConfig> {
198 self.reject_unknown_keys(&raw)?;
199 self.validate_defaults(options)?;
200
201 let mut resolved_config = ResolvedConfig::default();
202 for key in self.keys.values() {
203 let (value, is_default) = if let Some(value) = raw.get(&key.name).cloned() {
204 (Some(value), false)
205 } else {
206 (key.default.clone(), true)
207 };
208 let Some(value) = value else {
209 if key.required {
210 return Err(ConfigError::MissingRequired {
211 key: key.name.clone(),
212 });
213 }
214 continue;
215 };
216
217 if key.kind == ConfigKind::Secret {
218 let secret = resolve_secret_value(&key.name, value, resolver, options).await?;
219 resolved_config.insert_secret(key.name.clone(), secret);
220 } else if let Err(err) = validate_kind(&key.name, key.kind, &value) {
221 if is_default {
222 return Err(ConfigError::InvalidDefault {
223 key: key.name.clone(),
224 reason: err.to_string(),
225 });
226 }
227 return Err(err);
228 } else {
229 resolved_config.insert_plain(key.name.clone(), value);
230 }
231 }
232
233 Ok(resolved_config)
234 }
235
236 fn reject_unknown_keys(&self, raw: &RawConfig) -> ConfigResult<()> {
237 let known = self.keys.keys().collect::<BTreeSet<_>>();
238 for key in raw.keys() {
239 if !known.contains(key) {
240 return Err(ConfigError::UnknownKey { key: key.clone() });
241 }
242 }
243 Ok(())
244 }
245
246 fn validate_defaults(&self, options: ResolveOptions) -> ConfigResult<()> {
247 for key in self.keys.values() {
248 let Some(default) = &key.default else {
249 continue;
250 };
251 let result = if key.kind == ConfigKind::Secret {
252 validate_secret_default(default, options)
253 } else {
254 validate_kind(&key.name, key.kind, default).map_err(|err| err.to_string())
255 };
256 if let Err(err) = result {
257 return Err(ConfigError::InvalidDefault {
258 key: key.name.clone(),
259 reason: err,
260 });
261 }
262 }
263 Ok(())
264 }
265}
266
267fn validate_kind(key: &str, kind: ConfigKind, value: &Value) -> ConfigResult<()> {
268 let valid = match kind {
269 ConfigKind::String => value.is_string(),
270 ConfigKind::Bool => value.is_boolean(),
271 ConfigKind::Integer => value.as_i64().is_some(),
272 ConfigKind::DurationMillis | ConfigKind::DurationMs => {
273 value.as_i64().is_some_and(|millis| millis >= 0)
274 }
275 ConfigKind::UnsignedInteger => value.as_u64().is_some(),
276 ConfigKind::Float => value.as_f64().is_some(),
277 ConfigKind::StringList => value
278 .as_array()
279 .is_some_and(|items| items.iter().all(Value::is_string)),
280 ConfigKind::Json => true,
281 ConfigKind::Secret => unreachable!("secret kind validated separately"),
282 };
283
284 if valid {
285 Ok(())
286 } else {
287 Err(ConfigError::WrongType {
288 key: key.into(),
289 expected: kind.expected(),
290 })
291 }
292}
293
294fn validate_secret_default(value: &Value, options: ResolveOptions) -> Result<(), String> {
295 if value.is_string() {
296 if options.allow_literal_secrets {
297 return Ok(());
298 }
299 return Err("literal secret strings are disabled".into());
300 }
301
302 serde_json::from_value::<SecretRef>(value.clone())
303 .map(|_| ())
304 .map_err(|source| source.to_string())
305}
306
307async fn resolve_secret_value(
308 key: &str,
309 value: Value,
310 resolver: &dyn SecretResolver,
311 options: ResolveOptions,
312) -> ConfigResult<SecretString> {
313 if let Some(literal) = value.as_str() {
314 if options.allow_literal_secrets {
315 return Ok(SecretString::new(literal));
316 }
317 return Err(ConfigError::InvalidSecretRef {
318 key: key.into(),
319 reason: "literal secret strings are disabled".into(),
320 });
321 }
322
323 let secret_ref: SecretRef =
324 serde_json::from_value(value).map_err(|source| ConfigError::InvalidSecretRef {
325 key: key.into(),
326 reason: source.to_string(),
327 })?;
328
329 resolver
330 .resolve(&secret_ref)
331 .await
332 .map_err(|source| ConfigError::SecretResolution {
333 key: key.into(),
334 source: Box::new(source),
335 })
336}
337
338#[cfg(test)]
339mod tests {
340 use assert2::check;
341 use async_trait::async_trait;
342 use serde_json::{Value, json};
343
344 use super::*;
345 use crate::config::{ConfigError, EnvSecretResolver, ResolveOptions, SecretResolutionError};
346
347 fn raw(entries: impl IntoIterator<Item = (&'static str, Value)>) -> RawConfig {
348 entries
349 .into_iter()
350 .map(|(key, value)| (key.to_string(), value))
351 .collect()
352 }
353
354 #[tokio::test]
355 async fn resolve_applies_defaults_and_validates_required_fields() {
356 let def = ConfigDef::new("demo")
357 .required("database_url", ConfigKind::String)
358 .default("schema", ConfigKind::String, "public");
359 let raw = raw([("database_url", json!("postgres://localhost/app"))]);
360
361 let resolved = def.resolve(raw, &EnvSecretResolver).await.unwrap();
362
363 assert_eq!(
364 resolved.get_string("database_url").unwrap(),
365 "postgres://localhost/app"
366 );
367 assert_eq!(resolved.get_string("schema").unwrap(), "public");
368 assert!(!resolved.contains_key("missing_optional"));
369 }
370
371 #[tokio::test]
372 async fn resolve_rejects_unknown_keys() {
373 let def = ConfigDef::new("demo").required("database_url", ConfigKind::String);
374 let raw = raw([
375 ("database_url", json!("postgres://localhost/app")),
376 ("extra", json!(true)),
377 ]);
378
379 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
380
381 assert!(matches!(err, ConfigError::UnknownKey { key } if key == "extra"));
382 }
383
384 #[tokio::test]
385 async fn resolve_rejects_missing_required_keys() {
386 let def = ConfigDef::new("demo").required("database_url", ConfigKind::String);
387
388 let err = def
389 .resolve(RawConfig::new(), &EnvSecretResolver)
390 .await
391 .unwrap_err();
392
393 assert!(matches!(err, ConfigError::MissingRequired { key } if key == "database_url"));
394 }
395
396 #[tokio::test]
397 async fn resolve_rejects_wrong_types() {
398 let def = ConfigDef::new("demo").required("database_url", ConfigKind::String);
399 let raw = raw([("database_url", json!(42))]);
400
401 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
402
403 assert!(
404 matches!(err, ConfigError::WrongType { key, expected: "string" } if key == "database_url")
405 );
406 }
407
408 #[tokio::test]
409 async fn defaults_are_validated() {
410 let def = ConfigDef::new("demo").default("topics", ConfigKind::StringList, json!(["a", 7]));
411
412 let err = def
413 .resolve(RawConfig::new(), &EnvSecretResolver)
414 .await
415 .unwrap_err();
416
417 assert!(matches!(err, ConfigError::InvalidDefault { key, .. } if key == "topics"));
418 }
419
420 #[tokio::test]
421 async fn defaults_are_validated_even_when_raw_value_is_supplied() {
422 let def = ConfigDef::new("demo").default("topics", ConfigKind::StringList, json!(["a", 7]));
423 let raw = raw([("topics", json!(["user"]))]);
424
425 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
426
427 assert!(matches!(err, ConfigError::InvalidDefault { key, .. } if key == "topics"));
428 }
429
430 #[tokio::test]
431 async fn integer_kind_rejects_values_outside_i64_range() {
432 let def = ConfigDef::new("demo").required("limit", ConfigKind::Integer);
433 let raw = raw([("limit", json!(i64::MAX as u64 + 1))]);
434
435 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
436
437 assert!(
438 matches!(err, ConfigError::WrongType { key, expected: "integer" } if key == "limit")
439 );
440 }
441
442 #[tokio::test]
443 async fn unsigned_integer_kind_accepts_u64_max() {
444 let def = ConfigDef::new("demo").required("limit", ConfigKind::UnsignedInteger);
445 let raw = raw([("limit", json!(u64::MAX))]);
446
447 let resolved = def.resolve(raw, &EnvSecretResolver).await.unwrap();
448
449 assert_eq!(resolved.get_u64("limit").unwrap(), u64::MAX);
450 }
451
452 #[tokio::test]
453 async fn unsigned_integer_kind_rejects_negative_values() {
454 let def = ConfigDef::new("demo").required("limit", ConfigKind::UnsignedInteger);
455 let raw = raw([("limit", json!(-1))]);
456
457 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
458
459 assert!(
460 matches!(err, ConfigError::WrongType { key, expected: "unsigned integer" } if key == "limit")
461 );
462 }
463
464 #[tokio::test]
465 async fn duration_kinds_reject_negative_milliseconds() {
466 for kind in [ConfigKind::DurationMillis, ConfigKind::DurationMs] {
467 let def = ConfigDef::new("demo").required("timeout", kind);
468 let raw = raw([("timeout", json!(-1))]);
469
470 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
471
472 assert!(
473 matches!(err, ConfigError::WrongType { key, expected: "duration milliseconds" } if key == "timeout")
474 );
475 }
476 }
477
478 #[tokio::test]
479 async fn typed_getters_return_resolved_values() {
480 let def = ConfigDef::new("demo")
481 .required("name", ConfigKind::String)
482 .required("enabled", ConfigKind::Bool)
483 .required("limit", ConfigKind::Integer)
484 .required("unsigned_limit", ConfigKind::UnsignedInteger)
485 .required("timeout_ms", ConfigKind::DurationMillis)
486 .required("ratio", ConfigKind::Float)
487 .required("topics", ConfigKind::StringList)
488 .required("metadata", ConfigKind::Json)
489 .secret("password");
490 let raw = raw([
491 ("name", json!("source-a")),
492 ("enabled", json!(true)),
493 ("limit", json!(42)),
494 ("unsigned_limit", json!(u64::MAX)),
495 ("timeout_ms", json!(2500)),
496 ("ratio", json!(0.75)),
497 ("topics", json!(["alpha", "beta"])),
498 ("metadata", json!({"mode": "snapshot"})),
499 ("password", json!("literal-secret")),
500 ]);
501
502 let resolved = def
503 .resolve_with_options(
504 raw,
505 &EnvSecretResolver,
506 ResolveOptions {
507 allow_literal_secrets: true,
508 },
509 )
510 .await
511 .unwrap();
512
513 assert_eq!(resolved.get_string("name").unwrap(), "source-a");
514 assert!(resolved.get_bool("enabled").unwrap());
515 assert_eq!(resolved.get_i64("limit").unwrap(), 42);
516 assert_eq!(resolved.get_u64("unsigned_limit").unwrap(), u64::MAX);
517 assert_eq!(resolved.get_u64("timeout_ms").unwrap(), 2500);
518 assert!((resolved.get_f64("ratio").unwrap() - 0.75).abs() < f64::EPSILON);
519 assert_eq!(
520 resolved.get_string_list("topics").unwrap(),
521 vec!["alpha".to_string(), "beta".to_string()]
522 );
523 assert_eq!(
524 resolved.get_json("metadata").unwrap(),
525 json!({"mode": "snapshot"})
526 );
527 assert_eq!(
528 resolved.get_secret("password").unwrap().expose_secret(),
529 "literal-secret"
530 );
531 }
532
533 #[tokio::test]
534 async fn duration_ms_spelling_remains_supported() {
535 let def = ConfigDef::new("demo").required("timeout_ms", ConfigKind::DurationMs);
536 let raw = raw([("timeout_ms", json!(2500))]);
537
538 let resolved = def.resolve(raw, &EnvSecretResolver).await.unwrap();
539
540 assert_eq!(resolved.get_u64("timeout_ms").unwrap(), 2500);
541 assert_eq!(
542 ConfigKind::DurationMs.expected(),
543 ConfigKind::DurationMillis.expected()
544 );
545 }
546
547 #[test]
548 fn config_def_reports_its_name() {
549 let def = ConfigDef::new("postgres-source");
550
551 assert_eq!(def.name(), "postgres-source");
552 }
553
554 #[test]
555 fn config_def_debug_includes_name_and_redacts_secret_defaults() {
556 let key = ConfigKey {
557 name: "password".to_string(),
558 kind: ConfigKind::Secret,
559 required: false,
560 default: Some(json!("literal-secret")),
561 description: Some("database password".to_string()),
562 };
563 let def = ConfigDef {
564 name: "demo".to_string(),
565 keys: BTreeMap::from_iter([("password".to_string(), key)]),
566 };
567
568 let debug = format!("{def:?}");
569
570 check!(debug.contains("ConfigDef"));
571 check!(debug.contains("demo"));
572 check!(debug.contains("password"));
573 check!(debug.contains("<redacted>"));
574 check!(!debug.contains("literal-secret"));
575 }
576
577 #[test]
578 #[should_panic(expected = "duplicate connector config key `database_url`")]
579 fn duplicate_config_def_keys_panic_at_definition_time() {
580 let _ = ConfigDef::new("demo")
581 .required("database_url", ConfigKind::String)
582 .optional("database_url", ConfigKind::String);
583 }
584
585 #[test]
586 #[should_panic(expected = "secret connector config key `password` cannot have a default")]
587 fn secret_defaults_panic_at_definition_time() {
588 let _ = ConfigDef::new("demo").default("password", ConfigKind::Secret, "literal-secret");
589 }
590
591 #[test]
592 fn secret_defaults_are_redacted_in_debug_defensively() {
593 let key = ConfigKey {
594 name: "password".to_string(),
595 kind: ConfigKind::Secret,
596 required: false,
597 default: Some(json!("literal-secret")),
598 description: None,
599 };
600
601 let debug = format!("{key:?}");
602
603 assert!(!debug.contains("literal-secret"));
604 assert!(debug.contains("<redacted>"));
605 }
606
607 #[test]
608 fn non_secret_defaults_are_not_redacted_in_debug() {
609 let key = ConfigKey {
610 name: "schema".to_string(),
611 kind: ConfigKind::String,
612 required: false,
613 default: Some(json!("public")),
614 description: None,
615 };
616
617 let debug = format!("{key:?}");
618
619 assert!(debug.contains("public"));
620 assert!(!debug.contains("<redacted>"));
621 }
622
623 #[test]
624 fn secret_default_validation_rejects_literals_without_opt_in() {
625 let err = validate_secret_default(&json!("literal-secret"), ResolveOptions::default())
626 .expect_err("literal secret defaults require explicit opt in");
627
628 assert_eq!(err, "literal secret strings are disabled");
629 }
630
631 #[tokio::test]
632 async fn secret_literals_are_rejected_by_default() {
633 let def = ConfigDef::new("demo").secret("password");
634 let raw = raw([("password", json!("secret"))]);
635
636 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
637
638 assert!(matches!(err, ConfigError::InvalidSecretRef { key, .. } if key == "password"));
639 }
640
641 #[tokio::test]
642 async fn secret_literals_can_be_allowed_for_local_use_and_debug_is_redacted() {
643 let def = ConfigDef::new("demo").secret("password");
644 let raw = raw([("password", json!("secret"))]);
645
646 let resolved = def
647 .resolve_with_options(
648 raw,
649 &EnvSecretResolver,
650 ResolveOptions {
651 allow_literal_secrets: true,
652 },
653 )
654 .await
655 .unwrap();
656
657 assert_eq!(
658 resolved.get_secret("password").unwrap().expose_secret(),
659 "secret"
660 );
661 assert!(!format!("{resolved:?}").contains("secret"));
662 assert!(format!("{resolved:?}").contains("[REDACTED]"));
663 }
664
665 #[tokio::test]
666 async fn structured_secret_refs_resolve_through_provider() {
667 struct RecordingResolver;
668
669 #[async_trait]
670 impl SecretResolver for RecordingResolver {
671 async fn resolve(
672 &self,
673 secret_ref: &SecretRef,
674 ) -> Result<SecretString, SecretResolutionError> {
675 assert_eq!(
676 secret_ref,
677 &SecretRef::Env {
678 name: "POSTGRES_PASSWORD".into()
679 }
680 );
681 Ok(SecretString::new("resolved-password"))
682 }
683 }
684
685 let def = ConfigDef::new("demo").secret("password");
686 let raw = raw([(
687 "password",
688 json!({"from": "env", "name": "POSTGRES_PASSWORD"}),
689 )]);
690
691 let resolved = def.resolve(raw, &RecordingResolver).await.unwrap();
692
693 assert_eq!(
694 resolved.get_secret("password").unwrap().expose_secret(),
695 "resolved-password"
696 );
697 }
698
699 #[tokio::test]
700 async fn env_resolver_failure_reports_config_field_key() {
701 let def = ConfigDef::new("demo").secret("password");
702 let raw = raw([(
703 "password",
704 json!({"from": "env", "name": "CRABKA_CONNECT_TEST_MISSING_PASSWORD"}),
705 )]);
706
707 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
708
709 assert!(matches!(err, ConfigError::SecretResolution { ref key, .. } if key == "password"));
710 assert!(
711 !matches!(err, ConfigError::SecretResolution { ref key, .. } if key == "CRABKA_CONNECT_TEST_MISSING_PASSWORD")
712 );
713 }
714
715 #[tokio::test]
716 async fn secret_ref_unknown_fields_report_invalid_secret_ref_for_config_field() {
717 let def = ConfigDef::new("demo").secret("password");
718 let raw = raw([(
719 "password",
720 json!({"from": "env", "name": "POSTGRES_PASSWORD", "extra": true}),
721 )]);
722
723 let err = def.resolve(raw, &EnvSecretResolver).await.unwrap_err();
724
725 assert!(matches!(err, ConfigError::InvalidSecretRef { ref key, .. } if key == "password"));
726 assert!(err.to_string().contains("unknown field"));
727 }
728}