Skip to main content

faucet_source_postgres_cdc/
config.rs

1//! Configuration for `PostgresCdcSource`.
2
3use faucet_core::{DEFAULT_BATCH_SIZE, FaucetError};
4use schemars::JsonSchema;
5use serde::{Deserialize, Serialize};
6use std::time::Duration;
7
8fn default_true() -> bool {
9    true
10}
11fn default_proto_version() -> u32 {
12    1
13}
14fn default_idle_timeout() -> Duration {
15    Duration::from_secs(30)
16}
17fn default_status_update_interval() -> Duration {
18    Duration::from_secs(10)
19}
20fn default_tcp_keepalive() -> Duration {
21    Duration::from_secs(60)
22}
23fn default_batch_size() -> usize {
24    DEFAULT_BATCH_SIZE
25}
26fn default_slot_acquire_retries() -> u32 {
27    10
28}
29
30/// Configuration for [`PostgresCdcSource`](crate::PostgresCdcSource).
31#[derive(Clone, Serialize, Deserialize, JsonSchema)]
32#[serde(deny_unknown_fields)]
33pub struct PostgresCdcSourceConfig {
34    /// Connection URL pointing at the database whose WAL we want to read.
35    /// The crate internally upgrades the connection to `replication=database`
36    /// — callers do **not** need to add it themselves.
37    pub connection_url: String,
38
39    /// Logical replication slot name. Must match the Postgres naming rules:
40    /// 1–63 chars, lowercase letters / digits / underscores only.
41    pub slot_name: String,
42
43    /// Publication name on the server. Must already exist (faucet does not
44    /// create publications — they're a DBA-level concern that determines
45    /// which tables are replicated).
46    pub publication_name: String,
47
48    /// If the slot does not exist, create it as a logical/`pgoutput` slot
49    /// at connection time. Default: `true`.
50    #[serde(default = "default_true")]
51    pub create_slot_if_missing: bool,
52
53    /// Whether a newly-created slot is `permanent` (survives disconnect) or
54    /// `temporary` (auto-dropped when the replication connection closes).
55    ///
56    /// Default `permanent` (back-compatible). **A permanent slot pins WAL on
57    /// the server until it is consumed or dropped** — an abandoned permanent
58    /// slot fills `pg_wal` and can take the whole instance down. Use
59    /// `temporary` for ephemeral / test runs (note: a temporary slot resets on
60    /// reconnect, so bookmark-based resume across runs requires a permanent
61    /// slot). Drop an unused permanent slot explicitly with
62    /// [`PostgresCdcSource::drop_slot`](crate::PostgresCdcSource::drop_slot).
63    #[serde(default)]
64    pub slot_type: SlotType,
65
66    /// Number of times to retry acquiring the replication slot when the server
67    /// reports it is still **active** (held by a not-yet-released prior
68    /// connection). On a rapid restart — a scheduler or `serve` re-running the
69    /// pipeline before the previous backend has dropped the slot — both the
70    /// pre-stream `pg_replication_slot_advance` and `START_REPLICATION` fail
71    /// with *"replication slot … is active for PID …"*. Each retry waits an
72    /// exponentially increasing backoff (250 ms, doubling, capped at 4 s).
73    /// `0` disables retries (fail fast). Defaults to 10.
74    #[serde(default = "default_slot_acquire_retries")]
75    pub slot_acquire_retries: u32,
76
77    /// TLS settings for the replication connection. Default `disable`
78    /// (plaintext) for back-compatibility, but credentials and all WAL data
79    /// then travel unencrypted — set `require`/`verify_ca`/`verify_full` in
80    /// production.
81    #[serde(default)]
82    pub tls: CdcTls,
83
84    /// Optional starting LSN override (e.g. `"0/16A4F88"`). Ignored when a
85    /// state-store-managed bookmark is present (that bookmark wins).
86    /// When neither is set, replication starts from the slot's
87    /// `confirmed_flush_lsn`.
88    #[serde(default)]
89    pub start_lsn: Option<String>,
90
91    /// pgoutput protocol version. Only `1` is fully exercised in v1; `2` is
92    /// accepted but streaming-transaction messages (S/E/c/A) are not yet
93    /// decoded. Default: `1`.
94    #[serde(default = "default_proto_version")]
95    pub proto_version: u32,
96
97    /// Maximum time to wait for new replication messages before returning
98    /// the current batch. Default: 30 s.
99    #[serde(
100        default = "default_idle_timeout",
101        with = "faucet_core::config::duration_secs"
102    )]
103    #[schemars(with = "u64")]
104    pub idle_timeout: Duration,
105
106    /// Optional cap on the number of change events drained per fetch call.
107    /// Acts as a safety bound — `idle_timeout` is the primary terminator.
108    ///
109    /// **Note:** the cap is checked **after each COMMIT**, never mid-
110    /// transaction. A single transaction larger than `max_messages` will
111    /// still be emitted atomically (the fetch returns only after that
112    /// transaction's COMMIT and may produce more records than `max_messages`).
113    /// To bound the memory a *single* in-progress transaction can consume,
114    /// use [`max_staged_records`](Self::max_staged_records) instead.
115    #[serde(default)]
116    pub max_messages: Option<usize>,
117
118    /// Maximum number of change records buffered in memory for a *single*
119    /// in-progress transaction before it is aborted.
120    ///
121    /// Logical replication requires a transaction to be buffered until its
122    /// COMMIT so it can be emitted atomically (partial transactions must
123    /// never leak downstream). A single bulk `UPDATE`/`DELETE`/`COPY` of
124    /// millions of rows therefore buffers every decoded row as a
125    /// `serde_json::Value` in RAM, which can OOM the process. This bound is a
126    /// safety valve: when an in-progress transaction's staged record count
127    /// exceeds it, the source aborts with a typed
128    /// [`FaucetError::Source`] rather than
129    /// being OOM-killed.
130    ///
131    /// `None` (the default) means unbounded — atomic delivery of arbitrarily
132    /// large transactions at the cost of unbounded memory. Set a value sized
133    /// to your available memory if you replicate tables subject to large
134    /// bulk writes.
135    #[serde(default)]
136    pub max_staged_records: Option<usize>,
137
138    /// Interval at which Standby Status Update keepalives are sent to the
139    /// server. Must be shorter than `idle_timeout` and well under the
140    /// server's `wal_sender_timeout` (default 60 s). Default: 10 s.
141    #[serde(
142        default = "default_status_update_interval",
143        with = "faucet_core::config::duration_secs"
144    )]
145    #[schemars(with = "u64")]
146    pub status_update_interval: Duration,
147
148    /// TCP keepalive for the replication connection. Default: 60 s.
149    #[serde(
150        default = "default_tcp_keepalive",
151        with = "faucet_core::config::duration_secs"
152    )]
153    #[schemars(with = "u64")]
154    pub tcp_keepalive: Duration,
155
156    /// Advisory page size for
157    /// [`Source::stream_pages`](faucet_core::Source::stream_pages). The CDC
158    /// source emits **one `StreamPage` per committed transaction** so the
159    /// pipeline gets per-transaction durability via its per-page bookmark
160    /// persist. Because transactions are atomic units they are never split
161    /// across pages — a single transaction whose record count exceeds
162    /// `batch_size` still emits as one page. Defaults to
163    /// [`DEFAULT_BATCH_SIZE`].
164    ///
165    /// `batch_size = 0` is the "no batching" sentinel: every committed
166    /// transaction during the run window is accumulated into a single page
167    /// that is emitted at the end with `bookmark = max(commit_lsn)`. This
168    /// negates per-transaction durability and is only useful for tests or
169    /// initial-snapshot style runs.
170    #[serde(default = "default_batch_size")]
171    pub batch_size: usize,
172}
173
174/// Lifetime of a newly-created replication slot.
175#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
176#[serde(rename_all = "snake_case")]
177pub enum SlotType {
178    /// Survives disconnect; pins WAL until consumed or dropped. Default.
179    #[default]
180    Permanent,
181    /// Auto-dropped by the server when the replication connection closes.
182    Temporary,
183}
184
185/// TLS configuration for the CDC replication connection.
186#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
187#[serde(tag = "mode", rename_all = "snake_case")]
188pub enum CdcTls {
189    /// No TLS — plaintext (default, back-compatible).
190    #[default]
191    Disable,
192    /// Require TLS but do not verify the server certificate.
193    Require,
194    /// Require TLS and verify the certificate chain against `ca_path` (or the
195    /// system roots when `None`).
196    VerifyCa {
197        #[serde(default, skip_serializing_if = "Option::is_none")]
198        ca_path: Option<String>,
199    },
200    /// Require TLS and verify both the certificate chain and the hostname.
201    VerifyFull {
202        #[serde(default, skip_serializing_if = "Option::is_none")]
203        ca_path: Option<String>,
204    },
205}
206
207impl std::fmt::Debug for PostgresCdcSourceConfig {
208    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
209        f.debug_struct("PostgresCdcSourceConfig")
210            .field("connection_url", &"***")
211            .field("slot_name", &self.slot_name)
212            .field("publication_name", &self.publication_name)
213            .field("create_slot_if_missing", &self.create_slot_if_missing)
214            .field("slot_type", &self.slot_type)
215            .field("tls", &self.tls)
216            .field("start_lsn", &self.start_lsn)
217            .field("proto_version", &self.proto_version)
218            .field("idle_timeout", &self.idle_timeout)
219            .field("max_messages", &self.max_messages)
220            .field("max_staged_records", &self.max_staged_records)
221            .field("status_update_interval", &self.status_update_interval)
222            .field("tcp_keepalive", &self.tcp_keepalive)
223            .field("batch_size", &self.batch_size)
224            .field("slot_acquire_retries", &self.slot_acquire_retries)
225            .finish()
226    }
227}
228
229impl PostgresCdcSourceConfig {
230    /// Override the advisory per-page record count emitted by
231    /// [`Source::stream_pages`](faucet_core::Source::stream_pages).
232    ///
233    /// Pass `0` to disable per-transaction emission — every transaction in
234    /// the run window will be accumulated into a single trailing page with
235    /// `bookmark = max(commit_lsn)`. Transactions are never split regardless
236    /// of `batch_size`.
237    pub fn with_batch_size(mut self, batch_size: usize) -> Self {
238        self.batch_size = batch_size;
239        self
240    }
241
242    /// Validate fail-fast invariants. Called from `PostgresCdcSource::new`.
243    pub fn validate(&self) -> Result<(), FaucetError> {
244        if self.connection_url.trim().is_empty() {
245            return Err(FaucetError::Config(
246                "postgres-cdc: connection_url must not be empty".into(),
247            ));
248        }
249        validate_slot_name(&self.slot_name)?;
250        if self.publication_name.is_empty() {
251            return Err(FaucetError::Config(
252                "postgres-cdc: publication_name must not be empty".into(),
253            ));
254        }
255        if self.proto_version != 1 {
256            return Err(FaucetError::Config(format!(
257                "postgres-cdc: proto_version must be 1 (v2 streaming-transaction \
258                 support is not yet available via pgwire-replication), got {}",
259                self.proto_version
260            )));
261        }
262        if self.idle_timeout.is_zero() {
263            return Err(FaucetError::Config(
264                "postgres-cdc: idle_timeout must be > 0".into(),
265            ));
266        }
267        if self.status_update_interval >= self.idle_timeout {
268            return Err(FaucetError::Config(format!(
269                "postgres-cdc: status_update_interval ({}s) must be \
270                 strictly less than idle_timeout ({}s)",
271                self.status_update_interval.as_secs(),
272                self.idle_timeout.as_secs()
273            )));
274        }
275        Ok(())
276    }
277}
278
279fn validate_slot_name(name: &str) -> Result<(), FaucetError> {
280    if name.is_empty() {
281        return Err(FaucetError::Config(
282            "postgres-cdc: slot_name must not be empty".into(),
283        ));
284    }
285    if name.len() > 63 {
286        return Err(FaucetError::Config(format!(
287            "postgres-cdc: slot_name '{name}' exceeds Postgres' 63-char limit"
288        )));
289    }
290    if !name
291        .chars()
292        .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '_')
293    {
294        return Err(FaucetError::Config(format!(
295            "postgres-cdc: slot_name '{name}' must contain only \
296             [a-z0-9_]"
297        )));
298    }
299    Ok(())
300}
301
302#[cfg(test)]
303mod tests {
304    use super::*;
305
306    fn minimal() -> PostgresCdcSourceConfig {
307        PostgresCdcSourceConfig {
308            connection_url: "postgres://u:p@localhost/db".into(),
309            slot_name: "faucet_slot".into(),
310            publication_name: "faucet_pub".into(),
311            create_slot_if_missing: true,
312            slot_type: SlotType::Permanent,
313            tls: CdcTls::Disable,
314            start_lsn: None,
315            proto_version: 1,
316            idle_timeout: std::time::Duration::from_secs(30),
317            max_messages: None,
318            max_staged_records: None,
319            status_update_interval: std::time::Duration::from_secs(10),
320            tcp_keepalive: std::time::Duration::from_secs(60),
321            batch_size: DEFAULT_BATCH_SIZE,
322            slot_acquire_retries: default_slot_acquire_retries(),
323        }
324    }
325
326    #[test]
327    fn defaults_via_serde() {
328        let value: PostgresCdcSourceConfig = serde_json::from_value(serde_json::json!({
329            "connection_url": "postgres://u:p@localhost/db",
330            "slot_name": "faucet_slot",
331            "publication_name": "faucet_pub",
332        }))
333        .unwrap();
334        assert!(value.create_slot_if_missing);
335        assert_eq!(value.proto_version, 1);
336        assert_eq!(value.idle_timeout.as_secs(), 30);
337        assert_eq!(value.status_update_interval.as_secs(), 10);
338        assert_eq!(value.tcp_keepalive.as_secs(), 60);
339        assert!(value.start_lsn.is_none());
340        assert!(value.max_messages.is_none());
341        assert_eq!(value.batch_size, DEFAULT_BATCH_SIZE);
342    }
343
344    #[test]
345    fn batch_size_defaults_to_default_batch_size() {
346        let c = minimal();
347        assert_eq!(c.batch_size, DEFAULT_BATCH_SIZE);
348    }
349
350    #[test]
351    fn with_batch_size_overrides_default() {
352        let c = minimal().with_batch_size(64);
353        assert_eq!(c.batch_size, 64);
354    }
355
356    #[test]
357    fn batch_size_zero_is_accepted_as_no_batching_sentinel() {
358        let c = minimal().with_batch_size(0);
359        assert_eq!(c.batch_size, 0);
360        assert!(faucet_core::validate_batch_size(c.batch_size).is_ok());
361    }
362
363    #[test]
364    fn batch_size_above_max_is_rejected_by_validate_batch_size() {
365        let c = minimal().with_batch_size(faucet_core::MAX_BATCH_SIZE + 1);
366        assert!(faucet_core::validate_batch_size(c.batch_size).is_err());
367    }
368
369    #[test]
370    fn batch_size_deserializes_from_json() {
371        let v: PostgresCdcSourceConfig = serde_json::from_value(serde_json::json!({
372            "connection_url": "postgres://u:p@localhost/db",
373            "slot_name": "faucet_slot",
374            "publication_name": "faucet_pub",
375            "batch_size": 256,
376        }))
377        .unwrap();
378        assert_eq!(v.batch_size, 256);
379    }
380
381    #[test]
382    fn rejects_empty_slot_name() {
383        let mut c = minimal();
384        c.slot_name = String::new();
385        assert!(c.validate().is_err());
386    }
387
388    #[test]
389    fn rejects_invalid_slot_name_chars() {
390        let mut c = minimal();
391        c.slot_name = "Faucet-Slot".into(); // uppercase + dash both disallowed
392        assert!(c.validate().is_err());
393    }
394
395    #[test]
396    fn rejects_slot_name_over_63_chars() {
397        let mut c = minimal();
398        c.slot_name = "a".repeat(64);
399        assert!(c.validate().is_err());
400    }
401
402    #[test]
403    fn rejects_empty_publication_name() {
404        let mut c = minimal();
405        c.publication_name = String::new();
406        assert!(c.validate().is_err());
407    }
408
409    #[test]
410    fn rejects_zero_idle_timeout() {
411        let mut c = minimal();
412        c.idle_timeout = std::time::Duration::from_secs(0);
413        assert!(c.validate().is_err());
414    }
415
416    #[test]
417    fn rejects_status_update_interval_longer_than_idle_timeout() {
418        // Keepalives must fire before idle_timeout would terminate the loop.
419        let mut c = minimal();
420        c.status_update_interval = std::time::Duration::from_secs(60);
421        c.idle_timeout = std::time::Duration::from_secs(30);
422        assert!(c.validate().is_err());
423    }
424
425    #[test]
426    fn rejects_invalid_proto_version() {
427        // 0, 2, and 3 are all rejected — only 1 is supported.
428        let mut c = minimal();
429        c.proto_version = 0;
430        assert!(c.validate().is_err());
431        c.proto_version = 2;
432        assert!(c.validate().is_err());
433        c.proto_version = 3;
434        assert!(c.validate().is_err());
435    }
436
437    #[test]
438    fn accepts_proto_version_one() {
439        let mut c = minimal();
440        c.proto_version = 1;
441        assert!(c.validate().is_ok());
442    }
443
444    #[test]
445    fn rejects_empty_connection_url() {
446        let mut c = minimal();
447        c.connection_url = String::new();
448        assert!(c.validate().is_err());
449    }
450
451    #[test]
452    fn rejects_whitespace_connection_url() {
453        let mut c = minimal();
454        c.connection_url = "   ".into();
455        assert!(c.validate().is_err());
456    }
457
458    #[test]
459    fn debug_redacts_connection_url() {
460        let cfg = minimal();
461        let dbg = format!("{cfg:?}");
462        assert!(dbg.contains("connection_url: \"***\""));
463        assert!(!dbg.contains("u:p@localhost"));
464    }
465
466    #[test]
467    fn schema_for_config_includes_required_fields() {
468        let schema = schemars::schema_for!(PostgresCdcSourceConfig);
469        let json = serde_json::to_value(&schema).unwrap();
470        let required = json["required"].as_array().expect("required array");
471        let names: Vec<_> = required.iter().filter_map(|v| v.as_str()).collect();
472        assert!(names.contains(&"connection_url"));
473        assert!(names.contains(&"slot_name"));
474        assert!(names.contains(&"publication_name"));
475    }
476}