faucet_source_postgres_cdc/
config.rs1use 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#[derive(Clone, Serialize, Deserialize, JsonSchema)]
32#[serde(deny_unknown_fields)]
33pub struct PostgresCdcSourceConfig {
34 pub connection_url: String,
38
39 pub slot_name: String,
42
43 pub publication_name: String,
47
48 #[serde(default = "default_true")]
51 pub create_slot_if_missing: bool,
52
53 #[serde(default)]
64 pub slot_type: SlotType,
65
66 #[serde(default = "default_slot_acquire_retries")]
75 pub slot_acquire_retries: u32,
76
77 #[serde(default)]
82 pub tls: CdcTls,
83
84 #[serde(default)]
89 pub start_lsn: Option<String>,
90
91 #[serde(default = "default_proto_version")]
95 pub proto_version: u32,
96
97 #[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 #[serde(default)]
116 pub max_messages: Option<usize>,
117
118 #[serde(default)]
136 pub max_staged_records: Option<usize>,
137
138 #[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 #[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 #[serde(default = "default_batch_size")]
171 pub batch_size: usize,
172}
173
174#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
176#[serde(rename_all = "snake_case")]
177pub enum SlotType {
178 #[default]
180 Permanent,
181 Temporary,
183}
184
185#[derive(Clone, Debug, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
187#[serde(tag = "mode", rename_all = "snake_case")]
188pub enum CdcTls {
189 #[default]
191 Disable,
192 Require,
194 VerifyCa {
197 #[serde(default, skip_serializing_if = "Option::is_none")]
198 ca_path: Option<String>,
199 },
200 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 pub fn with_batch_size(mut self, batch_size: usize) -> Self {
238 self.batch_size = batch_size;
239 self
240 }
241
242 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(); 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 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 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}