use crate::config::ConnectorSpec;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
fn default_true() -> bool {
true
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct ReplicationSpec {
pub mode: ReplicationMode,
pub snapshot: SnapshotSpec,
#[serde(default = "default_true")]
pub continuous: bool,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum ReplicationMode {
SnapshotThenCdc,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct SnapshotSpec {
pub source: ConnectorSpec,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_minimal_replication_block() {
let yaml = r#"
mode: snapshot_then_cdc
snapshot:
source:
type: postgres
config: { connection_url: "postgres://x", query: "SELECT * FROM orders" }
"#;
let spec: ReplicationSpec = serde_yaml::from_str(yaml).unwrap();
assert_eq!(spec.mode, ReplicationMode::SnapshotThenCdc);
assert_eq!(spec.snapshot.source.kind, "postgres");
assert!(spec.continuous, "continuous defaults to true");
}
#[test]
fn continuous_false_parses() {
let yaml = r#"
mode: snapshot_then_cdc
continuous: false
snapshot:
source: { type: postgres, config: {} }
"#;
let spec: ReplicationSpec = serde_yaml::from_str(yaml).unwrap();
assert!(!spec.continuous);
}
#[test]
fn rejects_unknown_field() {
let yaml = r#"
mode: snapshot_then_cdc
snapshot: { source: { type: postgres, config: {} } }
bogus: true
"#;
assert!(serde_yaml::from_str::<ReplicationSpec>(yaml).is_err());
}
}