laminar_connectors/postgres/cdc/
mod.rs1pub mod changelog;
4pub mod config;
5pub mod decoder;
6pub mod lsn;
7pub mod metrics;
8pub mod postgres_io;
9pub mod schema;
10pub mod source;
11pub mod types;
12
13pub use crate::postgres::SslMode;
15pub use config::PostgresCdcConfig;
16pub use lsn::Lsn;
17pub use source::PostgresCdcSource;
18
19use std::sync::Arc;
20
21use crate::config::{ConfigKeySpec, ConnectorInfo};
22use crate::registry::ConnectorRegistry;
23
24pub fn register_postgres_cdc_source(
30 registry: &ConnectorRegistry,
31) -> Result<(), crate::error::ConnectorError> {
32 let info = ConnectorInfo {
33 name: "postgres-cdc".to_string(),
34 display_name: "PostgreSQL CDC Source".to_string(),
35 version: env!("CARGO_PKG_VERSION").to_string(),
36 is_source: true,
37 is_sink: false,
38 config_keys: postgres_cdc_config_keys(),
39 };
40
41 registry.register_source(
42 "postgres-cdc",
43 info,
44 Arc::new(|registry: Option<&Arc<prometheus::Registry>>| {
45 Ok(Box::new(PostgresCdcSource::new(
46 PostgresCdcConfig::default(),
47 registry.map(Arc::as_ref),
48 )))
49 }),
50 )?;
51
52 let pg_info = ConnectorInfo {
54 name: "postgres".to_string(),
55 display_name: "PostgreSQL Lookup Source".to_string(),
56 version: env!("CARGO_PKG_VERSION").to_string(),
57 is_source: true,
58 is_sink: false,
59 config_keys: postgres_lookup_config_keys(),
60 };
61 registry.register_table_source(
62 "postgres",
63 pg_info.clone(),
64 Arc::new(|config, declared_schema| {
65 Ok(Box::new(
66 crate::postgres::reference::PostgresReferenceTableSource::new(
67 config.clone(),
68 declared_schema,
69 ),
70 ))
71 }),
72 )?;
73
74 registry.register_lookup_source("postgres", pg_info, Arc::new(PostgresLookupFactory))
76}
77
78fn postgres_lookup_config_keys() -> Vec<ConfigKeySpec> {
79 vec![
80 ConfigKeySpec::required("table", "Qualified PostgreSQL table name"),
81 ConfigKeySpec::optional(
82 "connection",
83 "libpq connection string (alternative to individual connection properties)",
84 "",
85 ),
86 ConfigKeySpec::optional("connection_string", "Alias for connection", ""),
87 ConfigKeySpec::optional("host", "PostgreSQL host", "localhost"),
88 ConfigKeySpec::optional("port", "PostgreSQL port", "5432"),
89 ConfigKeySpec::optional("database", "Database name (alternative to dbname)", ""),
90 ConfigKeySpec::optional("dbname", "Alias for database", ""),
91 ConfigKeySpec::optional("user", "Database user (alternative to username)", ""),
92 ConfigKeySpec::optional("username", "Alias for user", ""),
93 ConfigKeySpec::optional("password", "Database password", ""),
94 ConfigKeySpec::optional("options", "PostgreSQL command-line options", ""),
95 ConfigKeySpec::optional("pool_size", "On-demand lookup connection-pool size", "4"),
96 ConfigKeySpec::optional(
97 "ssl.mode",
98 "Connection security: verify-full or explicit disable",
99 "verify-full",
100 ),
101 ConfigKeySpec::optional(
102 "ssl.ca.cert.path",
103 "PEM file with trusted CA certificates; defaults to webpki roots",
104 "",
105 ),
106 ]
107}
108
109struct PostgresLookupFactory;
110
111#[async_trait::async_trait]
112impl crate::registry::LookupSourceFactory for PostgresLookupFactory {
113 async fn build(
114 &self,
115 config: crate::config::ConnectorConfig,
116 _declared_schema: Option<arrow_schema::SchemaRef>,
117 ) -> Result<Arc<dyn laminar_core::lookup::source::LookupSourceDyn>, crate::error::ConnectorError>
118 {
119 use crate::postgres::lookup::{PostgresLookupSource, PostgresLookupSourceConfig};
120
121 let pk_columns: Vec<String> = config
122 .get("_primary_key_columns")
123 .unwrap_or("")
124 .split(',')
125 .map(|s| s.trim().to_string())
126 .filter(|s| !s.is_empty())
127 .collect();
128 if pk_columns.is_empty() {
129 return Err(crate::error::ConnectorError::ConfigurationError(
130 "postgres lookup source requires primary key columns".into(),
131 ));
132 }
133
134 let table = config
135 .get("table")
136 .ok_or_else(|| {
137 crate::error::ConnectorError::ConfigurationError(
138 "postgres lookup source requires a 'table' property".into(),
139 )
140 })?
141 .to_string();
142
143 let pool_size = if let Some(s) = config.get("pool_size") {
144 s.parse::<usize>().map_err(|e| {
145 crate::error::ConnectorError::ConfigurationError(format!(
146 "invalid 'pool_size' value '{s}': {e}"
147 ))
148 })?
149 } else {
150 4
151 };
152
153 let lookup_config = PostgresLookupSourceConfig {
154 properties: config.properties().clone(),
155 table,
156 primary_key_columns: pk_columns,
157 pool_size,
158 };
159
160 let source = PostgresLookupSource::open(lookup_config).await?;
161 Ok(Arc::new(source) as Arc<dyn laminar_core::lookup::source::LookupSourceDyn>)
162 }
163}
164
165fn postgres_cdc_config_keys() -> Vec<ConfigKeySpec> {
166 vec![
167 ConfigKeySpec::required("host", "PostgreSQL host address"),
168 ConfigKeySpec::required("database", "Database name"),
169 ConfigKeySpec::required("slot.name", "Logical replication slot name"),
170 ConfigKeySpec::required("publication", "Publication name"),
171 ConfigKeySpec::optional("port", "PostgreSQL port", "5432"),
172 ConfigKeySpec::optional("username", "Connection username", "postgres"),
173 ConfigKeySpec::optional("password", "Connection password", ""),
174 ConfigKeySpec::optional(
175 "ssl.mode",
176 "Connection security: verify-full or explicit disable",
177 "verify-full",
178 ),
179 ConfigKeySpec::optional(
180 "ssl.ca.cert.path",
181 "PEM file with trusted CA certificates; defaults to webpki roots",
182 "",
183 ),
184 ConfigKeySpec::optional(
185 "table.include",
186 "Comma-separated schema-qualified tables to include",
187 "",
188 ),
189 ConfigKeySpec::optional(
190 "table.exclude",
191 "Comma-separated schema-qualified tables to exclude",
192 "",
193 ),
194 ConfigKeySpec::optional(
195 "max.buffered.bytes",
196 "Total connector-owned payload budget in bytes",
197 "268435456",
198 ),
199 ]
200}
201
202#[cfg(test)]
203mod tests {
204 use super::*;
205
206 #[test]
207 fn test_register_postgres_cdc_source() {
208 let registry = ConnectorRegistry::new();
209 register_postgres_cdc_source(®istry).unwrap();
210
211 let info = registry.source_info("postgres-cdc");
212 assert!(info.is_some());
213 let info = info.unwrap();
214 assert_eq!(info.name, "postgres-cdc");
215 assert!(info.is_source);
216 assert!(!info.is_sink);
217 assert!(!info.config_keys.is_empty());
218 assert!(!registry
219 .list_table_sources()
220 .contains(&"postgres-cdc".to_string()));
221 let error = registry
222 .create_table_source(
223 &crate::config::ConnectorConfig::new("postgres-cdc"),
224 Arc::new(arrow_schema::Schema::empty()),
225 )
226 .err()
227 .expect("CDC polling gaps cannot define snapshot completion");
228 assert!(error.to_string().contains("snapshot-capable table source"));
229 }
230
231 #[test]
232 fn test_config_keys() {
233 let keys = postgres_cdc_config_keys();
234 let required: Vec<&str> = keys
235 .iter()
236 .filter(|k| k.required)
237 .map(|k| k.key.as_str())
238 .collect();
239 assert!(required.contains(&"host"));
240 assert!(required.contains(&"database"));
241 assert!(required.contains(&"slot.name"));
242 assert!(required.contains(&"publication"));
243 assert!(!required.contains(&"ssl.mode"));
244 assert!(keys.iter().any(|key| key.key == "ssl.ca.cert.path"));
245 for removed in [
246 "snapshot.mode",
247 "max.poll.records",
248 "wal.sender.timeout.ms",
249 "poll.timeout.ms",
250 "keepalive.interval.ms",
251 "backpressure.high.watermark",
252 "max.buffered.events",
253 "start.lsn",
254 "ssl.client.cert.path",
255 "ssl.client.key.path",
256 "ssl.sni.hostname",
257 ] {
258 assert!(keys.iter().all(|key| key.key != removed));
259 }
260 assert!(keys.iter().any(|key| key.key == "max.buffered.bytes"));
261 }
262
263 #[test]
264 fn test_lookup_config_keys_match_snapshot_and_on_demand_paths() {
265 let keys = postgres_lookup_config_keys();
266 let required: Vec<&str> = keys
267 .iter()
268 .filter(|key| key.required)
269 .map(|key| key.key.as_str())
270 .collect();
271 assert_eq!(required, ["table"]);
272
273 for supported in [
274 "connection",
275 "connection_string",
276 "host",
277 "port",
278 "database",
279 "dbname",
280 "user",
281 "username",
282 "password",
283 "options",
284 "pool_size",
285 "ssl.mode",
286 "ssl.ca.cert.path",
287 ] {
288 assert!(
289 keys.iter().any(|key| key.key == supported),
290 "PostgreSQL lookup descriptor omits {supported}"
291 );
292 }
293
294 for internal_or_rejected in ["_primary_key_columns", "sslmode", "sslrootcert"] {
295 assert!(keys.iter().all(|key| key.key != internal_or_rejected));
296 }
297 }
298}