Skip to main content

PostgresCdcSource

Struct PostgresCdcSource 

Source
pub struct PostgresCdcSource { /* private fields */ }

Implementations§

Source§

impl PostgresCdcSource

Source

pub async fn new(config: PostgresCdcSourceConfig) -> Result<Self, FaucetError>

Source

pub async fn drop_slot(&self) -> Result<(), FaucetError>

Drop this source’s replication slot on the server, freeing the WAL it pins. A no-op if the slot doesn’t exist; errors if the slot is still active (in use by a live replication connection). Call this when decommissioning a permanent slot so it doesn’t leak WAL (#78/#12).

Trait Implementations§

Source§

impl Source for PostgresCdcSource

Source§

fn fetch_with_context_incremental<'life0, 'life1, 'async_trait>( &'life0 self, ctx: &'life1 HashMap<String, Value>, ) -> Pin<Box<dyn Future<Output = Result<(Vec<Value>, Option<Value>), FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Drain the replication stream into a single Vec<Value> plus the bookmark of the most recent COMMIT.

Implemented by collecting Source::stream_pages with the batch_size = 0 sentinel — that sentinel coalesces every transaction in the run window into a single trailing page, which exactly matches the historical fetch_with_context_incremental contract (one aggregated buffer, one max-LSN bookmark). The streaming pipeline (Pipeline::run / run_stream) drives stream_pages directly with the per-source batch_size config field instead so it gets per-transaction durability.

Source§

fn stream_pages<'a>( &'a self, ctx: &'a HashMap<String, Value>, _batch_size: usize, ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>>

Per-transaction streaming.

Each committed transaction is emitted as its own StreamPage with bookmark = Some(commit_lsn). Because the pipeline flushes the sink and persists the bookmark on every page that carries one, a mid-stream crash recovers from the last fully- committed transaction with no partial-transaction leakage.

Atomic transactions. A transaction is never split across pages. If a single transaction’s record count exceeds PostgresCdcSourceConfig::batch_size it is still emitted as one page; batch_size is advisory.

batch_size = 0 is the “no batching” sentinel: every committed transaction during the run window is accumulated into a single trailing page with bookmark = max(commit_lsn). This negates per-transaction durability and is only useful for tests / initial snapshot runs.

The trait-level batch_size argument is intentionally ignored in favour of the config field (matches the convention used by the query-mode postgres source and the rest source).

Source§

fn check<'life0, 'life1, 'async_trait>( &'life0 self, ctx: &'life1 CheckContext, ) -> Pin<Box<dyn Future<Output = Result<CheckReport, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Preflight probe that does not start replication.

The default Source::check would call stream_pages, which opens the replication stream and consumes/holds WAL (a side effect that pins server resources) — unacceptable as a preflight. Instead we open a normal (non-replication) SQL connection — the same connection style ensure_slot uses — and inspect the slot catalog without touching the replication protocol:

  • connection fails → auth probe Fail (could not connect / authenticate),
  • connected but the slot row is absent → slot probe Skip (faucet can create it on the first run),
  • slot present → slot probe Pass.

The whole call is bounded by ctx.timeout.

Source§

fn fetch_with_context<'life0, 'life1, 'async_trait>( &'life0 self, ctx: &'life1 HashMap<String, Value>, ) -> Pin<Box<dyn Future<Output = Result<Vec<Value>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Primary fetch method. Receives context from a parent source’s records. Read more
Source§

fn config_schema(&self) -> Value

Return a JSON Schema describing the configuration this source accepts.
Source§

fn state_key(&self) -> Option<String>

Stable key under which this source’s incremental-replication bookmark should be persisted in a StateStore. Read more
Source§

fn apply_start_bookmark<'life0, 'async_trait>( &'life0 self, bookmark: Value, ) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Apply a bookmark loaded from a StateStore as this run’s starting point. Read more
Source§

fn capture_resume_position<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Option<Value>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Capture the source’s current replication position without consuming any changes, ensuring any server-side resource (e.g. a logical replication slot) needed to later resume from that position exists. Read more
Source§

fn lag<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Option<SourceLag>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

How far this source is behind its head (#733) — unread WAL bytes, binlog distance, unconsumed messages, or the age of the oldest unread change. Read more
Source§

fn supports_exactly_once(&self) -> bool

Whether this source deterministically replays the same page sequence from a given bookmark — the requirement for the atomic-watermark effectively-once path (a non-deterministic replay could cause the pipeline to skip a page whose contents differ from the one already committed). Default: false. Read more
Source§

fn connector_name(&self) -> &'static str

Stable identifier used as the connector label on metrics and the connector attribute on spans. Defaults to the final segment of std::any::type_name::<Self>(), e.g. "RestSource". Built-in connectors override with a short, friendly snake_case name (e.g. "rest"). Must return a non-empty string; observability decorators fall back to "unknown" in release builds if it is empty (and debug_assert! in debug builds).
Source§

fn record_table(&self, record: &Value) -> Option<String>

The dataset (table / collection) a change record belongs to, for a change stream that carries several tables (#731). A multi-table faucet mirror runs one stream and routes each record to its table’s pipeline by this name, which must match the name the paired bulk source’s discover reports (e.g. public.orders). None means the record belongs to no table (a DDL / control event) — and is the default, so a source that does not override this cannot be used for a multi-table mirror.
Source§

fn position_le(&self, a: &Value, b: &Value) -> Option<bool>

Order two of this source’s bookmarks (#731): Some(true) when stream position a is at or before b (every change up to a is also covered by b), Some(false) when it is not, and None when the source cannot tell (the default, or two positions it cannot relate). A multi-table mirror uses it to resume one shared stream from the earliest table and skip, per table, pages that table has already committed.
Source§

fn dataset_uri(&self) -> String

Logical dataset identity for lineage emission, following OpenLineage naming conventions (https://openlineage.io/docs/spec/naming). Read more
Source§

fn fetch_all<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Vec<Value>, FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Convenience: fetch with no parent context.
Source§

fn fetch_all_incremental<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<(Vec<Value>, Option<Value>), FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Convenience: incremental fetch with no parent context.
Source§

fn native_output_formats(&self) -> &'static [NativeFormat]

Wire formats this source can emit as raw bytes for the native byte-passthrough fast path (#633), in preference order (first = best). Default: &[] (no native fast path). Read more
Source§

fn stream_native<'a>( &'a self, context: &'a HashMap<String, Value>, format: NativeFormat, batch_size: usize, ) -> Pin<Box<dyn Stream<Item = Result<NativeBatch, FaucetError>> + Send + 'a>>

Stream the source natively as byte batches in format (#633). Read more
Source§

fn state_schema(&self) -> u32

The version of this source’s bookmark shape (#736). Stored state carries it; bump it whenever the shape changes and teach migrate_state the step. Default 0. Decorators must forward this.
Source§

fn migrate_state(&self, from: u32, data: Value) -> Result<Value, FaucetError>

Bring a bookmark stored at shape version from up to state_schema (#736). Must be pure — the migrated value is only persisted by the next bookmark write, so a crash in between re-runs it. The default knows only its own version. Decorators must forward this.
Source§

fn replay_guarantee(&self) -> ReplayGuarantee

The typed replay capability this source advertises — see ReplayGuarantee. Read more
Source§

fn is_shardable(&self) -> bool

Whether this source can split its work into independent shards for clustered (Mode B) execution. Default: false (single whole-dataset shard). Sources with a natural partition (object-store prefixes, table primary-key ranges) override this to true and implement enumerate_shards + apply_shard.
Source§

fn enumerate_shards<'life0, 'async_trait>( &'life0 self, _target: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<ShardSpec>, FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Enumerate the shards this source splits into, aiming for roughly target of them (a hint — the source may return fewer, e.g. when the dataset is small, or one per natural partition regardless of target). Read more
Source§

fn apply_shard<'life0, 'life1, 'async_trait>( &'life0 self, _shard: &'life1 ShardSpec, ) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Narrow this source instance to a single shard before streaming. Read more
Source§

fn range_digest<'life0, 'life1, 'life2, 'life3, 'async_trait>( &'life0 self, range: &'life1 KeyRange, key: &'life2 str, columns: &'life3 [String], ) -> Pin<Box<dyn Future<Output = Result<Option<ServerDigest>, FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait, Self: 'async_trait,

Compute a ServerDigest of the rows in range inside the backend, so a matching range ships no rows (#701). columns are the compared columns; key the integer key the range is over. Default Ok(None): not supported, the verifier streams the range instead. Two digests compare only when both sides report the same algorithm.
Source§

fn supports_discover(&self) -> bool

Whether this source can enumerate the datasets behind its connection via discover. Default: false. Sources backed by an introspectable catalog (database information_schema, MongoDB collections, Elasticsearch indices, object-store prefixes) override this to true.
Source§

fn discover<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Vec<DatasetDescriptor>, FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Enumerate the datasets living behind this source’s connection — one DatasetDescriptor per table / collection / index / prefix, each carrying a partial config override that selects it (used by faucet discover to scaffold one matrix row per dataset). Read more
Source§

fn position_min(&self, positions: &[Value]) -> Option<Value>

The earliest stream position every one of positions can resume from (#731): the shared change stream restarts there and each table skips what it has already applied. The default picks the position that position_le orders at or before all the others, and None when there is none (an empty slice, or positions this source cannot order). Sources whose positions are only partially ordered (one cursor per capture instance, say) override it with a component-wise minimum.
Source§

fn set_roundtrip_recorder(&self, _recorder: Arc<RoundtripRecorder>)

Receive a pre-labelled handle for counting upstream round trips (#638) — calls this connector makes to its own backend. Read more
Source§

fn set_run_clock(&self, _now: DateTime<Utc>)

Receive the run clock (#769): the instant ${now.*} renders from — faucet run --clock, a schedule tick, a backfill unit’s start — so a source that bounds reads by “now” (REST window slicing) reproduces the run as of that instant rather than the wall clock. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more