Skip to main content

faucet_source_postgres_cdc/
stream.rs

1//! `PostgresCdcSource` — public `Source` implementation.
2
3use 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
22/// Field added to a `before`/`after` image listing the columns Postgres elided
23/// because their out-of-line (TOAST) value wasn't rewritten by the change.
24///
25/// Reserved, and **part of the emitted record** — it is an ordinary key in the
26/// row image, so a downstream sink in `auto_map` mode sees it as a column and a
27/// schema-drift policy sees it as an addition. Named here rather than inlined at
28/// the insertion site so the reserved key has one definition and a consumer can
29/// filter on it (#654 L27). Postgres-specific (TOAST is a Postgres storage
30/// concept), so it lives in this crate rather than `faucet-core`.
31pub const UNCHANGED_TOAST_FIELD: &str = "__unchanged_toast__";
32
33pub struct PostgresCdcSource {
34    config: PostgresCdcSourceConfig,
35    state_key_value: String,
36    /// Bookmark provided by `apply_start_bookmark`, applied at the start of
37    /// the next fetch cycle. Becomes the new "confirmed_flush_lsn" we
38    /// advertise to Postgres.
39    pending_bookmark: Mutex<Option<Bookmark>>,
40    /// Last LSN we have told Postgres is durable (advertised as the slot's
41    /// `confirmed_flush_lsn`, which authorises WAL recycling up to it).
42    ///
43    /// Advanced **only** by `apply_start_bookmark` — i.e. after the pipeline
44    /// has durably persisted a bookmark — or seeded from `start_lsn` config.
45    /// It is never advanced from decoded WAL (commit-decode time) or from
46    /// keepalive `wal_end`, because those positions are not yet durable
47    /// downstream; doing so would let Postgres discard WAL for unwritten
48    /// changes and lose data on a crash (#78/#1).
49    confirmed_lsn: Mutex<u64>,
50    /// Highest LSN handed to the pipeline on a page bookmark this run — the
51    /// position [`Source::lag`] measures from while the slot's
52    /// `confirmed_flush_lsn` still trails (#733).
53    emitted_lsn: std::sync::atomic::AtomicU64,
54}
55
56/// Unread WAL: the server's current position minus the furthest point the
57/// consumer is known to have reached — the slot's `confirmed_flush_lsn` or,
58/// when ahead of it, this run's own position.
59fn 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    /// Drop this source's replication slot on the server, freeing the WAL it
81    /// pins. A no-op if the slot doesn't exist; errors if the slot is still
82    /// active (in use by a live replication connection). Call this when
83    /// decommissioning a `permanent` slot so it doesn't leak WAL (#78/#12).
84    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    /// Drain the replication stream into a single `Vec<Value>` plus the
105    /// bookmark of the most recent COMMIT.
106    ///
107    /// Implemented by collecting [`Source::stream_pages`] with the
108    /// `batch_size = 0` sentinel — that sentinel coalesces every transaction
109    /// in the run window into a single trailing page, which exactly matches
110    /// the historical `fetch_with_context_incremental` contract (one
111    /// aggregated buffer, one max-LSN bookmark). The streaming pipeline
112    /// (`Pipeline::run` / `run_stream`) drives `stream_pages` directly with
113    /// the per-source `batch_size` config field instead so it gets
114    /// per-transaction durability.
115    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    /// Per-transaction streaming.
134    ///
135    /// Each committed transaction is emitted as its own
136    /// [`StreamPage`] with `bookmark = Some(commit_lsn)`. Because the
137    /// pipeline flushes the sink and persists the bookmark on every page
138    /// that carries one, a mid-stream crash recovers from the last fully-
139    /// committed transaction with no partial-transaction leakage.
140    ///
141    /// **Atomic transactions.** A transaction is never split across pages.
142    /// If a single transaction's record count exceeds
143    /// [`PostgresCdcSourceConfig::batch_size`] it is still emitted as one
144    /// page; `batch_size` is advisory.
145    ///
146    /// **`batch_size = 0`** is the "no batching" sentinel: every committed
147    /// transaction during the run window is accumulated into a single
148    /// trailing page with `bookmark = max(commit_lsn)`. This negates
149    /// per-transaction durability and is only useful for tests / initial
150    /// snapshot runs.
151    ///
152    /// The trait-level `batch_size` argument is intentionally ignored in
153    /// favour of the config field (matches the convention used by the
154    /// query-mode postgres source and the rest source).
155    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        // Update confirmed_lsn so the next initial status update (in the next
175        // fetch cycle) advertises the correct position.
176        *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        // A temporary slot is dropped when the capture connection closes, so it
183        // cannot retain WAL across the snapshot — require a permanent slot.
184        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        // Durable monotonic LSN + deterministic replay from it + per-transaction
219        // (per-page) bookmarks persisted only after the sink confirms — the
220        // requirements for exactly-once delivery.
221        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    /// Preflight probe that does **not** start replication.
246    ///
247    /// The default [`Source::check`] would call `stream_pages`, which opens the
248    /// replication stream and consumes/holds WAL (a side effect that pins server
249    /// resources) — unacceptable as a preflight. Instead we open a *normal*
250    /// (non-replication) SQL connection — the same connection style
251    /// `ensure_slot` uses — and inspect the slot catalog without touching the
252    /// replication protocol:
253    ///
254    /// - connection fails → `auth` probe `Fail` (could not connect / authenticate),
255    /// - connected but the slot row is absent → `slot` probe `Skip`
256    ///   (faucet can create it on the first run),
257    /// - slot present → `slot` probe `Pass`.
258    ///
259    /// The whole call is bounded by `ctx.timeout`.
260    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        // A bad connection URL is a config error, not an unreachable server.
271        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        // Bound the whole connect+query under ctx.timeout so an unreachable
284        // host doesn't hang the probe.
285        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    /// Shared streaming implementation used by both `stream_pages` (per-
336    /// transaction page emission when `batch_size > 0`) and the legacy
337    /// `fetch_with_context_incremental` (single-page aggregation when
338    /// `batch_size == 0`).
339    ///
340    /// The slot lifecycle (`connect` → `ensure_slot` → `start_replication`)
341    /// and Standby Status Update bootstrap are identical to the pre-Plan-15
342    /// behaviour. The only difference is when records cross the page
343    /// boundary: on every COMMIT (`batch_size > 0`) or once at the end of
344    /// the run window (`batch_size == 0`).
345    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            // 1. Resolve start_lsn for THIS fetch cycle.
356            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            // 2. Open replication connection + ensure slot + START_REPLICATION.
373            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(&params).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            // Advance the slot's confirmed_flush_lsn to the durable resume
397            // point BEFORE opening the stream. For a logical slot,
398            // START_REPLICATION resumes decoding from confirmed_flush_lsn (the
399            // client-supplied start LSN does not skip already-committed
400            // transactions), so this is the only way to skip changes we have
401            // already consumed and durably persisted. `start_lsn` is set only
402            // from a durable bookmark (apply_start_bookmark) or the start_lsn
403            // config override; when it is `None` (no durable resume point) we
404            // do NOT advance, so an interrupted run with no persisted bookmark
405            // redelivers rather than loses data (#78/#1).
406            if let Some(lsn) = start_lsn {
407                // Both the slot advance and START_REPLICATION require the slot
408                // to be inactive; on a rapid re-run the prior connection may not
409                // have released it yet ("slot is active"). Retry with bounded
410                // backoff instead of failing the whole run (#146 M12).
411                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, &params),
425            )
426            .await?;
427
428            // Send an initial Standby Status Update advertising the same
429            // durable position so the server's bookkeeping is consistent.
430            let initial_confirmed = *self.confirmed_lsn.lock().await;
431            send_status_update(&mut duplex, initial_confirmed, false).await?;
432
433            // 4. Drain the replication stream until idle_timeout, ctrl_c, or
434            //    max_messages. Per-transaction mode (`batch_size > 0`) emits
435            //    one page per COMMIT carrying `bookmark = Some(commit_lsn)`.
436            //    Aggregated mode (`batch_size == 0`) accumulates every
437            //    transaction's records into one buffer and emits a single
438            //    trailing page with `bookmark = max(commit_lsn)`.
439            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                // Tracks per-iteration outcomes that we cannot propagate
455                // directly from inside `tokio::select!` arms while remaining
456                // inside `try_stream!`.
457                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                                // Reset on any server activity, not just committed
471                                // output, so idle_timeout never fires mid-transaction
472                                // while WAL is still flowing.
473                                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                                    // A COMMIT was just processed and
489                                    // `committed_records` holds the drained
490                                    // staged records.
491                                    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                    // NOTE: we deliberately do NOT advance `confirmed_lsn` here.
523                    // `confirmed_lsn` is the position advertised to Postgres as
524                    // `confirmed_flush_lsn` (which lets the server recycle WAL),
525                    // and it must only ever reflect data the consumer has
526                    // *durably* persisted. The only durable signal is the
527                    // bookmark the pipeline persists after the sink flush, which
528                    // arrives back via `apply_start_bookmark` at the start of the
529                    // next run. Advancing here at commit-decode time would tell
530                    // Postgres to discard WAL for changes that were never written
531                    // downstream — a crash in that window loses data (#78/#1).
532                    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            // 5. In aggregated mode emit the single trailing page (carrying
554            //    the max LSN seen). In per-transaction mode the trailing
555            //    `state.staged` is an *uncommitted* partial transaction and
556            //    must be dropped — Postgres will redeliver after the next
557            //    START_REPLICATION.
558            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
579/// One decoded tuple, split into its values and the names of any unchanged
580/// (large/TOAST) columns whose value the server didn't re-send.
581struct TupleRow {
582    values: Map<String, Value>,
583    unchanged_toast: Vec<String>,
584}
585
586/// In-flight transaction state while draining the replication stream.
587#[derive(Default)]
588struct TxnState {
589    /// Records produced inside the current BEGIN..COMMIT, buffered until
590    /// COMMIT is seen so partial transactions never leak into the output.
591    /// On COMMIT, `handle_event` drains this into the caller-supplied
592    /// `out: &mut Vec<Value>`.
593    staged: Vec<Value>,
594    /// commit_lsn of the most recently fully-applied transaction.
595    last_committed: Option<u64>,
596    /// commit_ts (Postgres epoch micros) of the in-progress transaction,
597    /// set by BEGIN.
598    in_progress_ts: i64,
599    /// commit_lsn announced by the in-progress BEGIN (== final_lsn).
600    in_progress_lsn: u64,
601    /// Whether we are currently inside a BEGIN..COMMIT pair.
602    in_txn: bool,
603    /// Optional cap on `staged.len()` for a single transaction. `None` means
604    /// unbounded. Exceeding it aborts the run with a typed error instead of
605    /// risking an OOM-kill on a huge bulk transaction (see
606    /// [`PostgresCdcSourceConfig::max_staged_records`]).
607    max_staged_records: Option<usize>,
608}
609
610impl TxnState {
611    /// Stage one decoded change record, enforcing
612    /// [`max_staged_records`](Self::max_staged_records).
613    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                // A second BEGIN without an intervening COMMIT is a protocol
642                // desync: silently discarding the staged records would drop
643                // committed-but-unemitted changes. Fail fast so the run
644                // restarts cleanly from the durable bookmark (#78/#46).
645                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            // Drain staged records into the caller-supplied output buffer.
667            // The caller is responsible for emitting a StreamPage with these
668            // records and the bookmark.
669            //
670            // The bookmark is the commit's `end_lsn` — the WAL position
671            // immediately AFTER the commit record — not the commit_lsn. To
672            // resume *past* a consumed transaction the slot's
673            // confirmed_flush_lsn must be set to a position beyond the commit
674            // record; advancing only to commit_lsn leaves the commit record
675            // unconfirmed and Postgres redelivers the whole transaction
676            // (#78/#1). `end_lsn` is the position a standby reports as flushed.
677            //
678            // The exactly-once-across-resume boundary (a transaction whose
679            // commit lands exactly at the persisted bookmark is delivered once,
680            // not skipped or duplicated) is exercised by the Docker integration
681            // tests `resume_from_bookmark_skips_already_consumed` and
682            // `lsn_not_advanced_without_durable_bookmark_redelivers` (#78 LOW).
683            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            // pg_logical_emit_message — a user-emitted logical message, never a
693            // table change. Intentionally ignored by this row-oriented source.
694        }
695        // KeepAlive and StoppedAt are filtered inside recv(). Any other variant
696        // is one this decoder doesn't understand — it may carry table data, so
697        // dropping it silently risks data loss. Fail fast instead (#78/#46).
698        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 => {} // ignored
716        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        // Begin/Commit pgoutput messages should never arrive here — the
721        // pgwire-replication library decodes them into structured
722        // ReplicationEvent::Begin / Commit variants, handled in handle_event.
723        // If we see one, log a warning and ignore.
724        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
812/// Convert a tuple's text cells to a [`TupleRow`].
813fn 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
844/// `schema.table` of a change envelope, the name the `postgres` source's
845/// discovery reports for the same table.
846fn 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    /// 'U' relation 'O' fullold 'N' new — exercises REPLICA IDENTITY FULL.
939    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    /// 'D' relation 'O' fullold — REPLICA IDENTITY FULL delete.
961    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    /// 'T' flags=0 oids... — truncate without cascade/restart_identity.
974    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); // flags
979        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        // The bookmark is the commit's `end_lsn` (the resume position just
1044        // past the commit record), not the commit_lsn. `commit_event` sets
1045        // end_lsn = commit_lsn + 0x10. The record's own "lsn" field above is
1046        // still the commit_lsn (display/identity), which is intentional.
1047        assert_eq!(state.last_committed, Some(0x16A_4F88 + 0x10));
1048    }
1049
1050    #[test]
1051    fn staging_beyond_max_staged_records_aborts() {
1052        // Regression for #78/#2: a single transaction that stages more records
1053        // than `max_staged_records` must abort with a typed error rather than
1054        // buffering unboundedly (OOM risk).
1055        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        // First two inserts stage fine.
1065        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        // The third exceeds the cap and must error.
1075        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        // With max_staged_records = None (default) an arbitrarily large
1092        // transaction stages without error.
1093        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        // Regression for #78/#46: a second BEGIN without an intervening COMMIT
1144        // is a stream desync and must abort, not silently discard staged rows.
1145        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        // Insert references relation 99999 which is not in the registry.
1172        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        // Build a second relation so the truncate-list has two known OIDs.
1254        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        // Exercise Fix 1: a REPLICA IDENTITY FULL update with an UnchangedToast
1288        // cell in the old tuple must record the column name in
1289        // before.__unchanged_toast__.
1290        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        // Hand-build an UPDATE where the OLD tuple's `name` cell is 'u' (unchanged TOAST).
1297        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        // id = text "1"
1303        text_cell(&mut buf, "1");
1304        // name = unchanged TOAST
1305        buf.push(b'u');
1306        // New tuple: id=1, name="alice2"
1307        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        // The name is part of the emitted record, so downstream configs (a sink
1330        // column list, a contract, a drift allowlist) are written against this
1331        // literal — renaming it is a wire-format change, not a refactor.
1332        assert_eq!(UNCHANGED_TOAST_FIELD, "__unchanged_toast__");
1333        // `faucet-core` cannot depend on a connector, so `cdc_unwrap` carries
1334        // its own copy of this reserved name to strip. If the two ever drift,
1335        // the marker silently reaches the sink as a real column (#670 L27).
1336        assert_eq!(
1337            UNCHANGED_TOAST_FIELD,
1338            faucet_core::stage::CDC_UNCHANGED_TOAST_FIELD
1339        );
1340    }
1341
1342    // dataset_uri is a pure-config method; the source requires a live DB to
1343    // construct (async + real PG connection), so we verify the logic directly.
1344    #[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}