Skip to main content

crabka_connect/config/
def.rs

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
14/// Incoming connector configuration as a JSON object.
15pub type RawConfig = Map<String, Value>;
16
17/// Supported logical connector configuration kinds.
18#[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    /// Return a human-readable type expectation for diagnostics.
34    #[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/// One connector configuration field definition.
51#[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/// ConfigDef-style connector configuration schema.
62#[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    /// Create an empty configuration definition for a connector.
97    #[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    /// Return the connector definition name.
106    #[must_use]
107    pub fn name(&self) -> &str {
108        &self.name
109    }
110
111    /// Iterate over declared config keys in stable name order.
112    pub fn keys(&self) -> impl Iterator<Item = &ConfigKey> {
113        self.keys.values()
114    }
115
116    /// Add a required configuration key.
117    #[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    /// Add an optional configuration key.
131    #[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    /// Add an optional configuration key with a default value.
145    #[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    /// Add a required secret configuration key.
168    #[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    /// Validate raw configuration and resolve secret references.
182    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    /// Validate raw configuration and resolve secret references with options.
192    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}