1use crate::config::PostgresCdcSourceConfig;
4use crate::pgoutput::decoder::decode_message;
5use crate::pgoutput::messages::{
6 Delete, Insert, Message, Relation, Truncate, TupleCell, TupleData, Update,
7};
8use crate::pgoutput::registry::RelationRegistry;
9use crate::pgoutput::values::text_to_json;
10use crate::replication::{
11 self, ReplicationEvent, ReplicationParams, postgres_clock_to_unix_ms, recv, send_status_update,
12};
13use crate::state::{Bookmark, format_lsn, parse_lsn, state_key};
14use async_trait::async_trait;
15use faucet_core::{FaucetError, Source, Stream, StreamPage};
16use serde_json::{Map, Value, json};
17use std::collections::HashMap;
18use std::pin::Pin;
19use std::time::{Duration, Instant};
20use tokio::sync::Mutex;
21
22pub const UNCHANGED_TOAST_FIELD: &str = "__unchanged_toast__";
32
33pub struct PostgresCdcSource {
34 config: PostgresCdcSourceConfig,
35 state_key_value: String,
36 pending_bookmark: Mutex<Option<Bookmark>>,
40 confirmed_lsn: Mutex<u64>,
50 emitted_lsn: std::sync::atomic::AtomicU64,
54}
55
56fn slot_lag_bytes(current_wal: u64, slot_confirmed: Option<u64>, local: u64) -> u64 {
60 current_wal.saturating_sub(slot_confirmed.unwrap_or(0).max(local))
61}
62
63impl PostgresCdcSource {
64 pub async fn new(config: PostgresCdcSourceConfig) -> Result<Self, FaucetError> {
65 config.validate()?;
66 let key = state_key(&config.slot_name);
67 let initial_lsn = match config.start_lsn.as_deref() {
68 Some(s) => parse_lsn(s)?,
69 None => 0,
70 };
71 Ok(Self {
72 config,
73 state_key_value: key,
74 pending_bookmark: Mutex::new(None),
75 confirmed_lsn: Mutex::new(initial_lsn),
76 emitted_lsn: std::sync::atomic::AtomicU64::new(0),
77 })
78 }
79
80 pub async fn drop_slot(&self) -> Result<(), FaucetError> {
85 replication::drop_slot(
86 &self.config.connection_url,
87 &self.config.slot_name,
88 &self.config.tls,
89 )
90 .await
91 }
92}
93
94#[async_trait]
95impl Source for PostgresCdcSource {
96 async fn fetch_with_context(
97 &self,
98 ctx: &HashMap<String, Value>,
99 ) -> Result<Vec<Value>, FaucetError> {
100 let (records, _bookmark) = self.fetch_with_context_incremental(ctx).await?;
101 Ok(records)
102 }
103
104 async fn fetch_with_context_incremental(
116 &self,
117 ctx: &HashMap<String, Value>,
118 ) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
119 use futures::StreamExt;
120 let mut pages = self.stream_pages_with_batch_size(ctx, 0);
121 let mut all: Vec<Value> = Vec::new();
122 let mut bookmark: Option<Value> = None;
123 while let Some(page) = pages.next().await {
124 let page = page?;
125 all.extend(page.records);
126 if page.bookmark.is_some() {
127 bookmark = page.bookmark;
128 }
129 }
130 Ok((all, bookmark))
131 }
132
133 fn stream_pages<'a>(
156 &'a self,
157 ctx: &'a HashMap<String, Value>,
158 _batch_size: usize,
159 ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>> {
160 self.stream_pages_with_batch_size(ctx, self.config.batch_size)
161 }
162
163 fn config_schema(&self) -> Value {
164 let schema = schemars::schema_for!(PostgresCdcSourceConfig);
165 serde_json::to_value(&schema).unwrap_or(Value::Null)
166 }
167
168 fn state_key(&self) -> Option<String> {
169 Some(self.state_key_value.clone())
170 }
171
172 async fn apply_start_bookmark(&self, bookmark: Value) -> Result<(), FaucetError> {
173 let parsed = Bookmark::from_value(bookmark)?;
174 *self.confirmed_lsn.lock().await = parsed.as_u64()?;
177 *self.pending_bookmark.lock().await = Some(parsed);
178 Ok(())
179 }
180
181 async fn capture_resume_position(&self) -> Result<Option<Value>, FaucetError> {
182 if matches!(self.config.slot_type, crate::config::SlotType::Temporary) {
185 return Err(FaucetError::Config(
186 "postgres-cdc: replication snapshot handoff requires a permanent slot \
187 (slot_type: permanent) so WAL is retained across the snapshot — a \
188 temporary slot is dropped when the capture connection closes"
189 .into(),
190 ));
191 }
192 let lsn = replication::ensure_slot_and_current_lsn(
193 &self.config.connection_url,
194 &self.config.slot_name,
195 self.config.create_slot_if_missing,
196 self.config.slot_type,
197 &self.config.tls,
198 )
199 .await?;
200 Ok(Some(crate::state::Bookmark::from_u64(lsn).to_value()?))
201 }
202
203 async fn lag(&self) -> Result<Option<faucet_core::SourceLag>, FaucetError> {
204 let local = (*self.confirmed_lsn.lock().await)
205 .max(self.emitted_lsn.load(std::sync::atomic::Ordering::Relaxed));
206 let positions = replication::slot_positions(
207 &self.config.connection_url,
208 &self.config.slot_name,
209 &self.config.tls,
210 )
211 .await?;
212 Ok(positions.map(|(current, confirmed)| {
213 faucet_core::SourceLag::bytes(slot_lag_bytes(current, confirmed, local))
214 }))
215 }
216
217 fn supports_exactly_once(&self) -> bool {
218 true
222 }
223
224 fn connector_name(&self) -> &'static str {
225 "postgres-cdc"
226 }
227
228 fn record_table(&self, record: &Value) -> Option<String> {
229 schema_table(record)
230 }
231
232 fn position_le(&self, a: &Value, b: &Value) -> Option<bool> {
233 let lsn = |v: &Value| Bookmark::from_value(v.clone()).ok()?.as_u64().ok();
234 Some(lsn(a)? <= lsn(b)?)
235 }
236
237 fn dataset_uri(&self) -> String {
238 format!(
239 "{}?publication={}",
240 faucet_core::redact_uri_credentials(&self.config.connection_url),
241 self.config.publication_name
242 )
243 }
244
245 async fn check(
261 &self,
262 ctx: &faucet_core::check::CheckContext,
263 ) -> Result<faucet_core::check::CheckReport, FaucetError> {
264 use faucet_core::check::{CheckReport, Probe};
265 use sqlx::ConnectOptions as _;
266 use sqlx::postgres::{PgConnectOptions, PgConnection};
267
268 let start = std::time::Instant::now();
269
270 let opts: PgConnectOptions = match self.config.connection_url.parse() {
272 Ok(o) => o,
273 Err(e) => {
274 return Ok(CheckReport::single(Probe::fail_hint(
275 "auth",
276 start.elapsed(),
277 format!("invalid connection URL: {e}"),
278 "connection_url must be a valid postgres:// URL",
279 )));
280 }
281 };
282
283 let probe = async {
286 let mut conn: PgConnection = opts.connect().await.map_err(|e| {
287 Probe::fail_hint(
288 "auth",
289 start.elapsed(),
290 format!("could not connect: {e}"),
291 "verify the host is reachable and credentials are valid",
292 )
293 })?;
294
295 let row: Option<(String,)> = sqlx::query_as(
296 "SELECT slot_name::text FROM pg_replication_slots WHERE slot_name = $1",
297 )
298 .bind(&self.config.slot_name)
299 .fetch_optional(&mut conn)
300 .await
301 .map_err(|e| {
302 Probe::fail(
303 "slot",
304 start.elapsed(),
305 format!("could not query pg_replication_slots: {e}"),
306 )
307 })?;
308
309 Ok::<Probe, Probe>(match row {
310 Some(_) => Probe::pass("slot", start.elapsed()),
311 None => Probe::skip(
312 "slot",
313 format!(
314 "replication slot {} does not exist yet (faucet run can create it)",
315 self.config.slot_name
316 ),
317 ),
318 })
319 };
320
321 let probe = match tokio::time::timeout(ctx.timeout, probe).await {
322 Ok(Ok(p)) | Ok(Err(p)) => p,
323 Err(_elapsed) => Probe::fail_hint(
324 "auth",
325 start.elapsed(),
326 "connection timed out",
327 "the database did not respond within the check timeout",
328 ),
329 };
330 Ok(CheckReport::single(probe))
331 }
332}
333
334impl PostgresCdcSource {
335 fn stream_pages_with_batch_size<'a>(
346 &'a self,
347 _ctx: &'a HashMap<String, Value>,
348 batch_size: usize,
349 ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>> {
350 let max_messages = self.config.max_messages.unwrap_or(usize::MAX);
351 let idle_timeout = self.config.idle_timeout;
352 let per_transaction = batch_size != 0;
353
354 Box::pin(async_stream::try_stream! {
355 let pending = {
357 let mut g = self.pending_bookmark.lock().await;
358 g.take()
359 };
360 let start_lsn = if let Some(b) = pending.as_ref() {
361 let lsn = b.as_u64()?;
362 *self.confirmed_lsn.lock().await = lsn;
363 Some(lsn)
364 } else {
365 self.config
366 .start_lsn
367 .as_deref()
368 .map(parse_lsn)
369 .transpose()?
370 };
371
372 let params = ReplicationParams {
374 connection_url: &self.config.connection_url,
375 slot_name: &self.config.slot_name,
376 publication_name: &self.config.publication_name,
377 proto_version: self.config.proto_version,
378 create_slot_if_missing: self.config.create_slot_if_missing,
379 start_lsn,
380 status_update_interval: self.config.status_update_interval,
381 tcp_keepalive: self.config.tcp_keepalive,
382 slot_type: self.config.slot_type,
383 tls: &self.config.tls,
384 };
385 let client = replication::connect(¶ms).await?;
386 replication::ensure_slot(
387 &client,
388 &self.config.connection_url,
389 &self.config.slot_name,
390 self.config.create_slot_if_missing,
391 self.config.slot_type,
392 &self.config.tls,
393 )
394 .await?;
395
396 if let Some(lsn) = start_lsn {
407 replication::retry_on_slot_active(self.config.slot_acquire_retries, || {
412 replication::advance_slot(
413 &self.config.connection_url,
414 &self.config.slot_name,
415 lsn,
416 &self.config.tls,
417 )
418 })
419 .await?;
420 }
421
422 let mut duplex = replication::retry_on_slot_active(
423 self.config.slot_acquire_retries,
424 || replication::start_replication(&client, ¶ms),
425 )
426 .await?;
427
428 let initial_confirmed = *self.confirmed_lsn.lock().await;
431 send_status_update(&mut duplex, initial_confirmed, false).await?;
432
433 let mut registry = RelationRegistry::new();
440 let mut state = TxnState {
441 max_staged_records: self.config.max_staged_records,
442 ..TxnState::default()
443 };
444 let mut agg_records: Vec<Value> = Vec::new();
445 let mut total_records: usize = 0;
446 let mut last_message_at = Instant::now();
447
448 loop {
449 let idle_deadline = last_message_at + idle_timeout;
450 let budget = idle_deadline
451 .checked_duration_since(Instant::now())
452 .unwrap_or(Duration::ZERO);
453
454 let mut stop = false;
458 let mut just_committed: Option<(u64, Vec<Value>)> = None;
459 let mut fatal: Option<FaucetError> = None;
460 let mut unexpected_end = false;
461 tokio::select! {
462 biased;
463 _ = tokio::signal::ctrl_c() => {
464 tracing::info!("postgres-cdc: ctrl_c received, stopping cleanly");
465 stop = true;
466 }
467 ev = tokio::time::timeout(budget, recv(&mut duplex)) => {
468 match ev {
469 Ok(Ok(Some(event))) => {
470 last_message_at = Instant::now();
474 let was_in_txn = state.in_txn;
475 let pre_commit_count = state.last_committed;
476 let mut committed_records: Vec<Value> = Vec::new();
477 if let Err(e) = handle_event(
478 event,
479 &mut registry,
480 &mut state,
481 &mut committed_records,
482 ) {
483 fatal = Some(e);
484 } else if was_in_txn
485 && !state.in_txn
486 && state.last_committed != pre_commit_count
487 {
488 let lsn = state.last_committed
492 .expect("last_committed set on commit");
493 total_records += committed_records.len();
494 just_committed = Some((lsn, committed_records));
495 }
496 }
497 Ok(Ok(None)) => {
498 unexpected_end = true;
499 }
500 Ok(Err(e)) => {
501 fatal = Some(e);
502 }
503 Err(_timeout) => {
504 tracing::debug!(
505 "postgres-cdc: idle_timeout reached, stopping"
506 );
507 stop = true;
508 }
509 }
510 }
511 }
512
513 if let Some(e) = fatal {
514 Err(e)?;
515 }
516 if unexpected_end {
517 Err(FaucetError::Source(
518 "postgres-cdc: replication stream ended unexpectedly".into(),
519 ))?;
520 }
521 if let Some((lsn, drained)) = just_committed {
522 if per_transaction {
533 self.emitted_lsn
534 .fetch_max(lsn, std::sync::atomic::Ordering::Relaxed);
535 let bookmark = Some(Bookmark::from_u64(lsn).to_value()?);
536 yield StreamPage {
537 records: drained,
538 bookmark,
539 };
540 } else {
541 agg_records.extend(drained);
542 }
543 if total_records >= max_messages {
544 stop = true;
545 }
546 }
547
548 if stop {
549 break;
550 }
551 }
552
553 if !per_transaction
559 && let Some(lsn) = state.last_committed
560 {
561 self.emitted_lsn
562 .fetch_max(lsn, std::sync::atomic::Ordering::Relaxed);
563 let bookmark = Some(Bookmark::from_u64(lsn).to_value()?);
564 yield StreamPage {
565 records: agg_records,
566 bookmark,
567 };
568 }
569
570 tracing::info!(
571 records = total_records,
572 batch_size,
573 "postgres-cdc: stream complete",
574 );
575 })
576 }
577}
578
579struct TupleRow {
582 values: Map<String, Value>,
583 unchanged_toast: Vec<String>,
584}
585
586#[derive(Default)]
588struct TxnState {
589 staged: Vec<Value>,
594 last_committed: Option<u64>,
596 in_progress_ts: i64,
599 in_progress_lsn: u64,
601 in_txn: bool,
603 max_staged_records: Option<usize>,
608}
609
610impl TxnState {
611 fn push_staged(&mut self, record: Value) -> Result<(), FaucetError> {
614 if let Some(max) = self.max_staged_records
615 && self.staged.len() >= max
616 {
617 return Err(FaucetError::Source(format!(
618 "postgres-cdc: in-progress transaction exceeded max_staged_records ({max}); \
619 aborting to avoid unbounded memory growth. Raise max_staged_records or \
620 reduce the size of the source transaction."
621 )));
622 }
623 self.staged.push(record);
624 Ok(())
625 }
626}
627
628fn handle_event(
629 event: ReplicationEvent,
630 registry: &mut RelationRegistry,
631 state: &mut TxnState,
632 out: &mut Vec<Value>,
633) -> Result<(), FaucetError> {
634 match event {
635 ReplicationEvent::Begin {
636 final_lsn,
637 commit_time_micros,
638 xid: _,
639 } => {
640 if state.in_txn {
641 return Err(FaucetError::Source(format!(
646 "postgres-cdc: BEGIN received while a previous transaction was still \
647 in progress ({} records staged) — replication stream desync",
648 state.staged.len()
649 )));
650 }
651 state.in_txn = true;
652 state.in_progress_lsn = final_lsn.as_u64();
653 state.in_progress_ts = commit_time_micros;
654 state.staged.clear();
655 }
656 ReplicationEvent::Commit {
657 lsn: _,
658 commit_time_micros: _,
659 end_lsn,
660 } => {
661 if !state.in_txn {
662 return Err(FaucetError::Source(
663 "postgres-cdc: COMMIT without BEGIN".into(),
664 ));
665 }
666 out.append(&mut state.staged);
684 state.last_committed = Some(end_lsn.as_u64());
685 state.in_txn = false;
686 }
687 ReplicationEvent::XLogData { data, .. } => {
688 let msg = decode_message(&data)?;
689 handle_pgoutput(msg, registry, state)?;
690 }
691 ReplicationEvent::Message { .. } => {
692 }
695 other => {
699 return Err(FaucetError::Source(format!(
700 "postgres-cdc: unhandled ReplicationEvent variant {other:?} — refusing to \
701 continue rather than risk silently dropping change data"
702 )));
703 }
704 }
705 Ok(())
706}
707
708fn handle_pgoutput(
709 msg: Message,
710 registry: &mut RelationRegistry,
711 state: &mut TxnState,
712) -> Result<(), FaucetError> {
713 match msg {
714 Message::Relation(r) => registry.insert(r),
715 Message::Origin | Message::Type => {} Message::Insert(i) => stage_insert(state, registry, i)?,
717 Message::Update(u) => stage_update(state, registry, u)?,
718 Message::Delete(d) => stage_delete(state, registry, d)?,
719 Message::Truncate(t) => stage_truncate(state, registry, t)?,
720 Message::Begin(_) | Message::Commit(_) => {
725 tracing::warn!(
726 "postgres-cdc: pgoutput Begin/Commit reached pgoutput decoder; \
727 pgwire-replication should have intercepted it"
728 );
729 }
730 }
731 Ok(())
732}
733
734fn stage_insert(
735 state: &mut TxnState,
736 registry: &RelationRegistry,
737 i: Insert,
738) -> Result<(), FaucetError> {
739 let rel = registry.get(i.relation_oid)?;
740 let after = tuple_to_object(rel, &i.new)?;
741 let r = record(rel, "insert", state, None, Some(after));
742 state.push_staged(r)
743}
744
745fn stage_update(
746 state: &mut TxnState,
747 registry: &RelationRegistry,
748 u: Update,
749) -> Result<(), FaucetError> {
750 let rel = registry.get(u.relation_oid)?;
751 let before = match &u.old {
752 Some(t) => Some(tuple_to_object(rel, t)?),
753 None => None,
754 };
755 let after = tuple_to_object(rel, &u.new)?;
756 let r = record(rel, "update", state, before, Some(after));
757 state.push_staged(r)
758}
759
760fn stage_delete(
761 state: &mut TxnState,
762 registry: &RelationRegistry,
763 d: Delete,
764) -> Result<(), FaucetError> {
765 let rel = registry.get(d.relation_oid)?;
766 let before = Some(tuple_to_object(rel, &d.old)?);
767 let r = record(rel, "delete", state, before, None);
768 state.push_staged(r)
769}
770
771fn stage_truncate(
772 state: &mut TxnState,
773 registry: &RelationRegistry,
774 t: Truncate,
775) -> Result<(), FaucetError> {
776 for oid in &t.relation_oids {
777 let rel = registry.get(*oid)?;
778 let r = record(rel, "truncate", state, None, None);
779 state.push_staged(r)?;
780 }
781 Ok(())
782}
783
784fn record(
785 rel: &Relation,
786 op: &str,
787 state: &TxnState,
788 before: Option<TupleRow>,
789 after: Option<TupleRow>,
790) -> Value {
791 fn to_value(row: TupleRow) -> Value {
792 let mut o = row.values;
793 if !row.unchanged_toast.is_empty() {
794 o.insert(UNCHANGED_TOAST_FIELD.into(), json!(row.unchanged_toast));
795 }
796 Value::Object(o)
797 }
798 let mut obj = Map::new();
799 obj.insert("op".into(), json!(op));
800 obj.insert("schema".into(), json!(rel.namespace));
801 obj.insert("table".into(), json!(rel.name));
802 obj.insert("lsn".into(), json!(format_lsn(state.in_progress_lsn)));
803 obj.insert(
804 "ts_ms".into(),
805 json!(postgres_clock_to_unix_ms(state.in_progress_ts)),
806 );
807 obj.insert("before".into(), before.map(to_value).unwrap_or(Value::Null));
808 obj.insert("after".into(), after.map(to_value).unwrap_or(Value::Null));
809 Value::Object(obj)
810}
811
812fn tuple_to_object(rel: &Relation, tup: &TupleData) -> Result<TupleRow, FaucetError> {
814 if tup.cells.len() != rel.columns.len() {
815 return Err(FaucetError::Source(format!(
816 "postgres-cdc: tuple has {} cells but relation {}.{} has {} columns",
817 tup.cells.len(),
818 rel.namespace,
819 rel.name,
820 rel.columns.len()
821 )));
822 }
823 let mut values = Map::with_capacity(rel.columns.len());
824 let mut unchanged_toast = Vec::new();
825 for (col, cell) in rel.columns.iter().zip(&tup.cells) {
826 match cell {
827 TupleCell::Null => {
828 values.insert(col.name.clone(), Value::Null);
829 }
830 TupleCell::UnchangedToast => {
831 unchanged_toast.push(col.name.clone());
832 }
833 TupleCell::Text(s) => {
834 values.insert(col.name.clone(), text_to_json(col.type_oid, s)?);
835 }
836 }
837 }
838 Ok(TupleRow {
839 values,
840 unchanged_toast,
841 })
842}
843
844fn schema_table(record: &Value) -> Option<String> {
847 let schema = record.get("schema")?.as_str()?;
848 let table = record.get("table")?.as_str()?;
849 Some(format!("{schema}.{table}"))
850}
851
852#[cfg(test)]
853mod tests {
854
855 #[tokio::test]
856 async fn routes_by_schema_table_and_orders_by_lsn() {
857 let src = PostgresCdcSource::new(
858 serde_json::from_value(json!({
859 "connection_url": "postgres://u:p@localhost/db",
860 "slot_name": "s",
861 "publication_name": "p"
862 }))
863 .unwrap(),
864 )
865 .await
866 .unwrap();
867 assert_eq!(
868 src.record_table(&json!({"schema": "public", "table": "orders"})),
869 Some("public.orders".into())
870 );
871 assert_eq!(src.record_table(&json!({"table": "orders"})), None);
872 let a = json!({"last_lsn": "0/16A4F88"});
873 let b = json!({"last_lsn": "1/0"});
874 assert_eq!(src.position_le(&a, &b), Some(true));
875 assert_eq!(src.position_le(&b, &a), Some(false));
876 assert_eq!(src.position_le(&a, &a), Some(true));
877 assert_eq!(src.position_le(&a, &json!({"last_lsn": "bad"})), None);
878 assert_eq!(src.position_min(&[b.clone(), a.clone()]), Some(a));
879 }
880
881 #[test]
882 fn slot_lag_bytes_measures_from_the_furthest_known_position() {
883 assert_eq!(slot_lag_bytes(1000, Some(400), 0), 600);
884 assert_eq!(slot_lag_bytes(1000, Some(400), 900), 100);
885 assert_eq!(slot_lag_bytes(1000, None, 250), 750);
886 assert_eq!(slot_lag_bytes(1000, Some(1200), 0), 0);
887 }
888
889 use super::*;
890 use crate::pgoutput::messages::{ColumnDesc, ReplicaIdentity};
891 use crate::replication::ReplicationEvent;
892 use pgwire_replication::Lsn;
893
894 fn rel_users() -> Relation {
895 Relation {
896 oid: 16384,
897 namespace: "public".into(),
898 name: "users".into(),
899 replica_identity: ReplicaIdentity::Default,
900 columns: vec![
901 ColumnDesc {
902 flags: 1,
903 name: "id".into(),
904 type_oid: 23,
905 type_modifier: -1,
906 },
907 ColumnDesc {
908 flags: 0,
909 name: "name".into(),
910 type_oid: 25,
911 type_modifier: -1,
912 },
913 ],
914 }
915 }
916
917 fn xlogdata(payload: Vec<u8>) -> ReplicationEvent {
918 ReplicationEvent::XLogData {
919 wal_start: Lsn::from_u64(0),
920 wal_end: Lsn::from_u64(0x16A_4F88),
921 server_time_micros: 0,
922 data: bytes::Bytes::from(payload),
923 }
924 }
925
926 fn insert_payload(relation_oid: u32, cells: &[(&str, &str)]) -> Vec<u8> {
927 let mut buf: Vec<u8> = Vec::new();
928 buf.push(b'I');
929 buf.extend_from_slice(&relation_oid.to_be_bytes());
930 buf.push(b'N');
931 buf.extend_from_slice(&(cells.len() as u16).to_be_bytes());
932 for (_, val) in cells {
933 text_cell(&mut buf, val);
934 }
935 buf
936 }
937
938 fn update_full_payload(
940 relation_oid: u32,
941 old_cells: &[(&str, &str)],
942 new_cells: &[(&str, &str)],
943 ) -> Vec<u8> {
944 let mut buf: Vec<u8> = Vec::new();
945 buf.push(b'U');
946 buf.extend_from_slice(&relation_oid.to_be_bytes());
947 buf.push(b'O');
948 buf.extend_from_slice(&(old_cells.len() as u16).to_be_bytes());
949 for (_, val) in old_cells {
950 text_cell(&mut buf, val);
951 }
952 buf.push(b'N');
953 buf.extend_from_slice(&(new_cells.len() as u16).to_be_bytes());
954 for (_, val) in new_cells {
955 text_cell(&mut buf, val);
956 }
957 buf
958 }
959
960 fn delete_full_payload(relation_oid: u32, old_cells: &[(&str, &str)]) -> Vec<u8> {
962 let mut buf: Vec<u8> = Vec::new();
963 buf.push(b'D');
964 buf.extend_from_slice(&relation_oid.to_be_bytes());
965 buf.push(b'O');
966 buf.extend_from_slice(&(old_cells.len() as u16).to_be_bytes());
967 for (_, val) in old_cells {
968 text_cell(&mut buf, val);
969 }
970 buf
971 }
972
973 fn truncate_payload(relation_oids: &[u32]) -> Vec<u8> {
975 let mut buf: Vec<u8> = Vec::new();
976 buf.push(b'T');
977 buf.extend_from_slice(&(relation_oids.len() as u32).to_be_bytes());
978 buf.push(0u8); for oid in relation_oids {
980 buf.extend_from_slice(&oid.to_be_bytes());
981 }
982 buf
983 }
984
985 fn text_cell(buf: &mut Vec<u8>, val: &str) {
986 buf.push(b't');
987 buf.extend_from_slice(&(val.len() as u32).to_be_bytes());
988 buf.extend_from_slice(val.as_bytes());
989 }
990
991 fn begin_event(final_lsn: u64) -> ReplicationEvent {
992 ReplicationEvent::Begin {
993 final_lsn: Lsn::from_u64(final_lsn),
994 xid: 1,
995 commit_time_micros: 0,
996 }
997 }
998
999 fn commit_event(lsn: u64) -> ReplicationEvent {
1000 ReplicationEvent::Commit {
1001 lsn: Lsn::from_u64(lsn),
1002 end_lsn: Lsn::from_u64(lsn + 0x10),
1003 commit_time_micros: 0,
1004 }
1005 }
1006
1007 #[test]
1008 fn full_transaction_promotes_to_output_on_commit() {
1009 let mut registry = RelationRegistry::new();
1010 registry.insert(rel_users());
1011 let mut state = TxnState::default();
1012 let mut out = vec![];
1013
1014 handle_event(begin_event(0x16A_4F88), &mut registry, &mut state, &mut out).unwrap();
1015 assert!(out.is_empty());
1016
1017 handle_event(
1018 xlogdata(insert_payload(16384, &[("id", "1"), ("name", "alice")])),
1019 &mut registry,
1020 &mut state,
1021 &mut out,
1022 )
1023 .unwrap();
1024 assert!(out.is_empty(), "records stay staged until COMMIT");
1025
1026 handle_event(
1027 commit_event(0x16A_4F88),
1028 &mut registry,
1029 &mut state,
1030 &mut out,
1031 )
1032 .unwrap();
1033
1034 assert_eq!(out.len(), 1);
1035 assert_eq!(out[0]["op"], "insert");
1036 assert_eq!(out[0]["schema"], "public");
1037 assert_eq!(out[0]["table"], "users");
1038 assert_eq!(out[0]["lsn"], "0/16A4F88");
1039 assert_eq!(out[0]["after"]["id"], 1);
1040 assert_eq!(out[0]["after"]["name"], "alice");
1041 assert_eq!(out[0]["before"], Value::Null);
1042
1043 assert_eq!(state.last_committed, Some(0x16A_4F88 + 0x10));
1048 }
1049
1050 #[test]
1051 fn staging_beyond_max_staged_records_aborts() {
1052 let mut registry = RelationRegistry::new();
1056 registry.insert(rel_users());
1057 let mut state = TxnState {
1058 max_staged_records: Some(2),
1059 ..TxnState::default()
1060 };
1061 let mut out = vec![];
1062
1063 handle_event(begin_event(0x16A_4F88), &mut registry, &mut state, &mut out).unwrap();
1064 for id in ["1", "2"] {
1066 handle_event(
1067 xlogdata(insert_payload(16384, &[("id", id), ("name", "x")])),
1068 &mut registry,
1069 &mut state,
1070 &mut out,
1071 )
1072 .unwrap();
1073 }
1074 let err = handle_event(
1076 xlogdata(insert_payload(16384, &[("id", "3"), ("name", "x")])),
1077 &mut registry,
1078 &mut state,
1079 &mut out,
1080 )
1081 .unwrap_err();
1082 assert!(
1083 format!("{err}").contains("max_staged_records"),
1084 "error must name the cap: {err}"
1085 );
1086 assert!(matches!(err, FaucetError::Source(_)));
1087 }
1088
1089 #[test]
1090 fn no_cap_allows_large_transactions() {
1091 let mut registry = RelationRegistry::new();
1094 registry.insert(rel_users());
1095 let mut state = TxnState::default();
1096 let mut out = vec![];
1097
1098 handle_event(begin_event(0x16A_4F88), &mut registry, &mut state, &mut out).unwrap();
1099 for id in 0..50 {
1100 handle_event(
1101 xlogdata(insert_payload(
1102 16384,
1103 &[("id", &id.to_string()), ("name", "x")],
1104 )),
1105 &mut registry,
1106 &mut state,
1107 &mut out,
1108 )
1109 .unwrap();
1110 }
1111 handle_event(
1112 commit_event(0x16A_4F88),
1113 &mut registry,
1114 &mut state,
1115 &mut out,
1116 )
1117 .unwrap();
1118 assert_eq!(out.len(), 50);
1119 }
1120
1121 #[test]
1122 fn commit_without_begin_errors() {
1123 let mut registry = RelationRegistry::new();
1124 let mut state = TxnState::default();
1125 let mut out = vec![];
1126
1127 let err = handle_event(
1128 ReplicationEvent::Commit {
1129 lsn: Lsn::from_u64(1),
1130 end_lsn: Lsn::from_u64(2),
1131 commit_time_micros: 0,
1132 },
1133 &mut registry,
1134 &mut state,
1135 &mut out,
1136 )
1137 .unwrap_err();
1138 assert!(format!("{err}").contains("COMMIT without BEGIN"));
1139 }
1140
1141 #[test]
1142 fn double_begin_errors() {
1143 let mut registry = RelationRegistry::new();
1146 registry.insert(rel_users());
1147 let mut state = TxnState::default();
1148 let mut out = vec![];
1149
1150 handle_event(begin_event(0x100), &mut registry, &mut state, &mut out).unwrap();
1151 handle_event(
1152 xlogdata(insert_payload(16384, &[("id", "1"), ("name", "alice")])),
1153 &mut registry,
1154 &mut state,
1155 &mut out,
1156 )
1157 .unwrap();
1158
1159 let err =
1160 handle_event(begin_event(0x200), &mut registry, &mut state, &mut out).unwrap_err();
1161 assert!(format!("{err}").contains("desync"), "{err}");
1162 }
1163
1164 #[test]
1165 fn unknown_relation_in_insert_errors() {
1166 let mut registry = RelationRegistry::new();
1167 let mut state = TxnState::default();
1168 let mut out = vec![];
1169
1170 handle_event(begin_event(1), &mut registry, &mut state, &mut out).unwrap();
1171 let err = handle_event(
1173 xlogdata(insert_payload(99999, &[("id", "1"), ("name", "alice")])),
1174 &mut registry,
1175 &mut state,
1176 &mut out,
1177 )
1178 .unwrap_err();
1179 assert!(format!("{err}").contains("99999"));
1180 }
1181
1182 #[test]
1183 fn update_with_replica_identity_full_emits_before_and_after() {
1184 let mut registry = RelationRegistry::new();
1185 registry.insert(rel_users());
1186 let mut state = TxnState::default();
1187 let mut out = vec![];
1188
1189 handle_event(begin_event(0x16A_4F88), &mut registry, &mut state, &mut out).unwrap();
1190 handle_event(
1191 xlogdata(update_full_payload(
1192 16384,
1193 &[("id", "1"), ("name", "alice")],
1194 &[("id", "1"), ("name", "alice2")],
1195 )),
1196 &mut registry,
1197 &mut state,
1198 &mut out,
1199 )
1200 .unwrap();
1201 handle_event(
1202 commit_event(0x16A_4F88),
1203 &mut registry,
1204 &mut state,
1205 &mut out,
1206 )
1207 .unwrap();
1208
1209 assert_eq!(out.len(), 1);
1210 assert_eq!(out[0]["op"], "update");
1211 assert_eq!(out[0]["before"]["id"], 1);
1212 assert_eq!(out[0]["before"]["name"], "alice");
1213 assert_eq!(out[0]["after"]["name"], "alice2");
1214 }
1215
1216 #[test]
1217 fn delete_with_replica_identity_full_emits_before_only() {
1218 let mut registry = RelationRegistry::new();
1219 registry.insert(rel_users());
1220 let mut state = TxnState::default();
1221 let mut out = vec![];
1222
1223 handle_event(begin_event(0x16A_4F88), &mut registry, &mut state, &mut out).unwrap();
1224 handle_event(
1225 xlogdata(delete_full_payload(
1226 16384,
1227 &[("id", "1"), ("name", "alice")],
1228 )),
1229 &mut registry,
1230 &mut state,
1231 &mut out,
1232 )
1233 .unwrap();
1234 handle_event(
1235 commit_event(0x16A_4F88),
1236 &mut registry,
1237 &mut state,
1238 &mut out,
1239 )
1240 .unwrap();
1241
1242 assert_eq!(out.len(), 1);
1243 assert_eq!(out[0]["op"], "delete");
1244 assert_eq!(out[0]["before"]["id"], 1);
1245 assert_eq!(out[0]["before"]["name"], "alice");
1246 assert_eq!(out[0]["after"], Value::Null);
1247 }
1248
1249 #[test]
1250 fn truncate_emits_one_record_per_relation() {
1251 let mut registry = RelationRegistry::new();
1252 registry.insert(rel_users());
1253 let mut second = rel_users();
1255 second.oid = 16385;
1256 second.name = "orders".into();
1257 registry.insert(second);
1258
1259 let mut state = TxnState::default();
1260 let mut out = vec![];
1261
1262 handle_event(begin_event(0x16A_4F88), &mut registry, &mut state, &mut out).unwrap();
1263 handle_event(
1264 xlogdata(truncate_payload(&[16384, 16385])),
1265 &mut registry,
1266 &mut state,
1267 &mut out,
1268 )
1269 .unwrap();
1270 handle_event(
1271 commit_event(0x16A_4F88),
1272 &mut registry,
1273 &mut state,
1274 &mut out,
1275 )
1276 .unwrap();
1277
1278 assert_eq!(out.len(), 2);
1279 assert!(out.iter().all(|r| r["op"] == "truncate"));
1280 let tables: Vec<_> = out.iter().map(|r| r["table"].as_str().unwrap()).collect();
1281 assert!(tables.contains(&"users"));
1282 assert!(tables.contains(&"orders"));
1283 }
1284
1285 #[test]
1286 fn unchanged_toast_in_before_surfaces_via_metadata() {
1287 let mut registry = RelationRegistry::new();
1291 registry.insert(rel_users());
1292 let mut state = TxnState::default();
1293 let mut out = vec![];
1294
1295 handle_event(begin_event(0x16A_4F88), &mut registry, &mut state, &mut out).unwrap();
1296 let mut buf: Vec<u8> = Vec::new();
1298 buf.push(b'U');
1299 buf.extend_from_slice(&16384u32.to_be_bytes());
1300 buf.push(b'O');
1301 buf.extend_from_slice(&2u16.to_be_bytes());
1302 text_cell(&mut buf, "1");
1304 buf.push(b'u');
1306 buf.push(b'N');
1308 buf.extend_from_slice(&2u16.to_be_bytes());
1309 text_cell(&mut buf, "1");
1310 text_cell(&mut buf, "alice2");
1311 handle_event(xlogdata(buf), &mut registry, &mut state, &mut out).unwrap();
1312 handle_event(
1313 commit_event(0x16A_4F88),
1314 &mut registry,
1315 &mut state,
1316 &mut out,
1317 )
1318 .unwrap();
1319
1320 assert_eq!(out.len(), 1);
1321 assert_eq!(out[0]["before"]["__unchanged_toast__"], json!(["name"]));
1322 assert!(out[0]["before"].get("name").is_none());
1323 assert_eq!(out[0]["before"]["id"], 1);
1324 assert_eq!(out[0]["after"]["name"], "alice2");
1325 }
1326
1327 #[test]
1328 fn unchanged_toast_field_name_is_pinned() {
1329 assert_eq!(UNCHANGED_TOAST_FIELD, "__unchanged_toast__");
1333 assert_eq!(
1337 UNCHANGED_TOAST_FIELD,
1338 faucet_core::stage::CDC_UNCHANGED_TOAST_FIELD
1339 );
1340 }
1341
1342 #[test]
1345 fn dataset_uri_strips_credentials() {
1346 let redacted = faucet_core::redact_uri_credentials("postgres://u:p@h:5432/db");
1347 let uri = format!("{}?publication={}", redacted, "my_pub");
1348 assert_eq!(uri, "postgres://h:5432/db?publication=my_pub");
1349 }
1350}