use crate::config::ConnectorSpec;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::BTreeMap;
fn default_true() -> bool {
true
}
fn default_snapshot_concurrency() -> usize {
4
}
fn default_discover_interval_secs() -> u64 {
300
}
fn default_max_table_failures() -> u32 {
3
}
fn default_retry_paused_secs() -> u64 {
300
}
fn default_lag_warning_secs() -> u64 {
300
}
fn default_include() -> Vec<String> {
vec!["*".to_string()]
}
#[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,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tables: Option<TablesSpec>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub per_table: BTreeMap<String, TableOverride>,
}
#[derive(Debug, Clone, Serialize, Deserialize, JsonSchema, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct TablesSpec {
#[serde(default = "default_include")]
pub include: Vec<String>,
#[serde(default)]
pub exclude: Vec<String>,
#[serde(default)]
pub new_tables: NewTables,
#[serde(default = "default_discover_interval_secs")]
pub discover_interval_secs: u64,
#[serde(default)]
pub without_primary_key: WithoutPrimaryKey,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub destination: Option<Value>,
#[serde(default)]
pub on_table_error: OnTableError,
#[serde(default = "default_max_table_failures")]
pub max_table_failures: u32,
#[serde(default = "default_retry_paused_secs")]
pub retry_paused_secs: u64,
#[serde(default = "default_lag_warning_secs")]
pub lag_warning_secs: u64,
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum NewTables {
#[default]
Follow,
Ignore,
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum WithoutPrimaryKey {
#[default]
Refuse,
Append,
}
#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, JsonSchema, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum OnTableError {
#[default]
Pause,
Fail,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct TableOverride {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub key: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub write_mode: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub sink: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub snapshot: Option<Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub schema_drift: Option<faucet_core::SchemaDriftSpec>,
}
#[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,
#[serde(default = "default_snapshot_concurrency")]
pub concurrency: usize,
#[serde(default)]
pub shards: usize,
}
#[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());
}
}