Skip to main content

faucet_core/
traits.rs

1//! Shared traits for faucet sources and sinks.
2
3use crate::error::FaucetError;
4use crate::pipeline::StreamPage;
5use async_trait::async_trait;
6use futures_core::Stream;
7use serde_json::Value;
8use std::pin::Pin;
9
10/// A source fetches records from an external system.
11#[async_trait]
12pub trait Source: Send + Sync {
13    /// Primary fetch method. Receives context from a parent source's records.
14    ///
15    /// An empty context map means this is a root source (no parent).
16    /// Connectors that support being a child should use
17    /// [`substitute_context()`](crate::util::substitute_context) to resolve
18    /// `{placeholder}` tokens in their URL path, query parameters, headers,
19    /// or body. Connectors that don't need parent context ignore the map.
20    async fn fetch_with_context(
21        &self,
22        context: &std::collections::HashMap<String, Value>,
23    ) -> Result<Vec<Value>, FaucetError>;
24
25    /// Convenience: fetch with no parent context.
26    async fn fetch_all(&self) -> Result<Vec<Value>, FaucetError> {
27        self.fetch_with_context(&std::collections::HashMap::new())
28            .await
29    }
30
31    /// Incremental fetch with parent context support.
32    ///
33    /// Returns the records and an optional bookmark value for incremental
34    /// replication. The default delegates to `fetch_with_context` and
35    /// returns `None` for the bookmark.
36    async fn fetch_with_context_incremental(
37        &self,
38        context: &std::collections::HashMap<String, Value>,
39    ) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
40        let records = self.fetch_with_context(context).await?;
41        Ok((records, None))
42    }
43
44    /// Convenience: incremental fetch with no parent context.
45    async fn fetch_all_incremental(&self) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
46        self.fetch_with_context_incremental(&std::collections::HashMap::new())
47            .await
48    }
49
50    /// Stream records page-by-page so the pipeline can write to the sink as
51    /// pages arrive instead of buffering the full result set.
52    ///
53    /// `batch_size` is the *hint* the pipeline passes down; sources are free
54    /// to use a larger or smaller native chunk (e.g. one page per HTTP
55    /// response, one row-group per Parquet file) but should approximate it
56    /// where feasible. The special value `batch_size = 0` means "do not
57    /// batch — emit the entire result set in a single page." Sources that
58    /// stream natively should treat `0` as "skip the chunking layer and
59    /// yield one page after the underlying read completes" (useful for
60    /// small lookup tables or for sinks like SQL `COPY` / BigQuery load
61    /// jobs that prefer one large request).
62    ///
63    /// The default implementation fetches the full result set via
64    /// [`fetch_with_context_incremental`](Self::fetch_with_context_incremental)
65    /// and chunks it in memory by `batch_size`. The bookmark (when present)
66    /// is attached to the *final* page so the pipeline only persists after
67    /// the entire fetch has been written. Sources that can stream natively
68    /// override this method and may emit per-page bookmarks (e.g. CDC).
69    ///
70    /// An empty result with a `Some(bookmark)` still yields one empty page
71    /// carrying the bookmark, so incremental runs that produce no records
72    /// still advance their checkpoint.
73    fn stream_pages<'a>(
74        &'a self,
75        context: &'a std::collections::HashMap<String, Value>,
76        batch_size: usize,
77    ) -> Pin<Box<dyn Stream<Item = Result<StreamPage, FaucetError>> + Send + 'a>> {
78        Box::pin(async_stream::try_stream! {
79            let (records, bookmark) = self
80                .fetch_with_context_incremental(context)
81                .await?;
82            let total = records.len();
83            // batch_size == 0 means "no batching" — emit all records as one
84            // page. Otherwise chunk into pages of size `batch_size`.
85            let chunk = if batch_size == 0 { usize::MAX } else { batch_size };
86
87            if total == 0 {
88                if bookmark.is_some() {
89                    yield StreamPage {
90                        records: Vec::new(),
91                        bookmark,
92                    };
93                }
94                return;
95            }
96
97            let mut iter = records.into_iter();
98            let mut consumed = 0usize;
99            loop {
100                let batch: Vec<Value> = iter.by_ref().take(chunk).collect();
101                if batch.is_empty() {
102                    break;
103                }
104                consumed += batch.len();
105                let page_bookmark = if consumed >= total {
106                    bookmark.clone()
107                } else {
108                    None
109                };
110                yield StreamPage {
111                    records: batch,
112                    bookmark: page_bookmark,
113                };
114            }
115        })
116    }
117
118    /// Whether this source can emit **columnar** ([`ColumnarPage`](crate::columnar::ColumnarPage))
119    /// pages via [`stream_batches`](Self::stream_batches). Default: `false`.
120    ///
121    /// The pipeline uses the columnar fast path only when *both* the source and
122    /// sink return `true` here (and no `Value`-shaped stage needs to observe the
123    /// records), so an Arrow-native `parquet → parquet` chain never materializes
124    /// `Value`. Opt-in and additive — see [`crate::columnar`] (RFC 0002 / #375).
125    #[cfg(feature = "arrow")]
126    fn supports_columnar(&self) -> bool {
127        false
128    }
129
130    /// Stream the source natively as Arrow
131    /// [`ColumnarPage`](crate::columnar::ColumnarPage)s.
132    ///
133    /// Only invoked when [`supports_columnar`](Self::supports_columnar) returns
134    /// `true`; the default yields a single typed "unsupported" error so a source
135    /// that advertises support but forgets to override this fails loudly rather
136    /// than silently. Each page's `bookmark` carries the same checkpoint
137    /// semantics as [`StreamPage`].
138    #[cfg(feature = "arrow")]
139    fn stream_batches<'a>(
140        &'a self,
141        context: &'a std::collections::HashMap<String, Value>,
142        batch_size: usize,
143    ) -> Pin<Box<dyn Stream<Item = Result<crate::columnar::ColumnarPage, FaucetError>> + Send + 'a>>
144    {
145        let _ = (context, batch_size);
146        let name = self.connector_name();
147        let err: Result<crate::columnar::ColumnarPage, FaucetError> = Err(FaucetError::Source(
148            format!("source '{name}' does not support columnar streaming (stream_batches)"),
149        ));
150        Box::pin(futures::stream::once(async move { err }))
151    }
152
153    /// Wire formats this source can emit as raw bytes for the **native
154    /// byte-passthrough** fast path (#633), in preference order (first = best).
155    /// Default: `&[]` (no native fast path).
156    ///
157    /// The pipeline uses the path only when this returns a non-empty slice, a
158    /// sink advertises a matching [`NativeLoadCapability`](crate::NativeLoadCapability),
159    /// and every prerequisite holds (see [`plan_native_transfer`](crate::plan_native_transfer)).
160    /// A [`TransformingSource`](crate::TransformingSource) deliberately does not
161    /// override this, so any attached transform disables the fast path.
162    fn native_output_formats(&self) -> &'static [crate::native::NativeFormat] {
163        &[]
164    }
165
166    /// Stream the source natively as byte batches in `format` (#633).
167    ///
168    /// Only invoked after the pipeline negotiated `format` — one this source
169    /// advertised via [`native_output_formats`](Self::native_output_formats). The
170    /// default yields a single typed "unsupported" error so a source that
171    /// advertises support but forgets to override this fails loudly. Each batch's
172    /// `bookmark` carries the same checkpoint semantics as [`StreamPage`].
173    fn stream_native<'a>(
174        &'a self,
175        context: &'a std::collections::HashMap<String, Value>,
176        format: crate::native::NativeFormat,
177        batch_size: usize,
178    ) -> Pin<Box<dyn Stream<Item = Result<crate::native::NativeBatch, FaucetError>> + Send + 'a>>
179    {
180        let _ = (context, batch_size, format);
181        let name = self.connector_name();
182        let err: Result<crate::native::NativeBatch, FaucetError> = Err(FaucetError::Source(
183            format!("source '{name}' does not support native byte streaming (stream_native)"),
184        ));
185        Box::pin(futures::stream::once(async move { err }))
186    }
187
188    /// Return a JSON Schema describing the configuration this source accepts.
189    fn config_schema(&self) -> Value {
190        serde_json::json!({"type": "object", "properties": {}})
191    }
192
193    /// Stable key under which this source's incremental-replication bookmark
194    /// should be persisted in a [`StateStore`](crate::state::StateStore).
195    ///
196    /// Returning `Some(key)` opts this source into resumable runs: when the
197    /// pipeline is configured with a state store via
198    /// [`Pipeline::with_state_store`](crate::Pipeline::with_state_store), it
199    /// reads the bookmark at `key` before fetching and writes the new
200    /// bookmark back only after the sink confirms the batch was written.
201    ///
202    /// The default returns `None`, meaning the source is not persisted.
203    /// Keys must satisfy [`validate_state_key`](crate::state::validate_state_key).
204    fn state_key(&self) -> Option<String> {
205        None
206    }
207
208    /// Apply a bookmark loaded from a [`StateStore`](crate::state::StateStore)
209    /// as this run's starting point.
210    ///
211    /// The default implementation ignores the value, which keeps existing
212    /// sources backwards-compatible. Sources that support incremental
213    /// replication override this — typically by storing the value behind
214    /// interior mutability and consulting it inside
215    /// `fetch_with_context_incremental`.
216    async fn apply_start_bookmark(&self, _bookmark: Value) -> Result<(), FaucetError> {
217        Ok(())
218    }
219
220    /// Capture the source's current replication position **without consuming
221    /// any changes**, ensuring any server-side resource (e.g. a logical
222    /// replication slot) needed to later resume from that position exists.
223    ///
224    /// Returns the position as a bookmark [`Value`] — the same shape
225    /// [`apply_start_bookmark`](Self::apply_start_bookmark) accepts — or `None`
226    /// if this source does not support position capture.
227    ///
228    /// Used by the snapshot→CDC replication orchestrator (`faucet replicate`)
229    /// to anchor the CDC stream at-or-before the bulk snapshot's read point so
230    /// the handoff has no gap. The default returns `None`.
231    async fn capture_resume_position(&self) -> Result<Option<Value>, FaucetError> {
232        Ok(None)
233    }
234
235    /// How far this source is behind its head (#733) — unread WAL bytes, binlog
236    /// distance, unconsumed messages, or the age of the oldest unread change.
237    ///
238    /// The pipeline polls it at page boundaries and exports
239    /// `faucet_source_lag_{bytes,events,seconds}`; `faucet status --probe`
240    /// asks for it on demand after applying the stored bookmark. It must be a
241    /// cheap, read-only query. An `Err` is logged once and reported as no lag;
242    /// it never fails a run. Default: `Ok(None)` (no notion of a head).
243    /// Decorators must forward this.
244    async fn lag(&self) -> Result<Option<crate::lag::SourceLag>, FaucetError> {
245        Ok(None)
246    }
247
248    /// The version of this source's bookmark shape (#736). Stored state carries
249    /// it; bump it whenever the shape changes and teach
250    /// [`migrate_state`](Self::migrate_state) the step. Default `0`. Decorators
251    /// must forward this.
252    fn state_schema(&self) -> u32 {
253        0
254    }
255
256    /// Bring a bookmark stored at shape version `from` up to
257    /// [`state_schema`](Self::state_schema) (#736). Must be pure — the
258    /// migrated value is only persisted by the next bookmark write, so a crash
259    /// in between re-runs it. The default knows only its own version. Decorators
260    /// must forward this.
261    fn migrate_state(&self, from: u32, data: Value) -> Result<Value, FaucetError> {
262        if from == self.state_schema() {
263            Ok(data)
264        } else {
265            Err(FaucetError::State(format!(
266                "{} has no migration from bookmark schema {from} to {}",
267                self.connector_name(),
268                self.state_schema()
269            )))
270        }
271    }
272
273    /// Whether this source **deterministically replays** the same page sequence
274    /// from a given bookmark — the requirement for the atomic-watermark
275    /// effectively-once path (a non-deterministic replay could cause the pipeline
276    /// to skip a page whose contents differ from the one already committed).
277    /// Default: `false`.
278    ///
279    /// Sources with a durable monotonic position and per-page bookmarks (CDC)
280    /// override this to return `true`. The pipeline rejects
281    /// `DeliveryMode::ExactlyOnce` against a source that returns `false`.
282    fn supports_exactly_once(&self) -> bool {
283        false
284    }
285
286    /// The typed replay capability this source advertises — see
287    /// [`ReplayGuarantee`](crate::ReplayGuarantee).
288    ///
289    /// The default derives from [`supports_exactly_once`](Self::supports_exactly_once)
290    /// (the boolean stays the back-compat primitive: existing connectors that
291    /// override only the boolean automatically advertise `Deterministic`
292    /// here). Override this directly only to *diverge* from the boolean —
293    /// there is currently no reason to.
294    fn replay_guarantee(&self) -> crate::idempotency::ReplayGuarantee {
295        if self.supports_exactly_once() {
296            crate::idempotency::ReplayGuarantee::Deterministic
297        } else {
298            crate::idempotency::ReplayGuarantee::NonDeterministic
299        }
300    }
301
302    /// Whether this source can split its work into independent shards for
303    /// clustered (Mode B) execution. Default: `false` (single whole-dataset
304    /// shard). Sources with a natural partition (object-store prefixes, table
305    /// primary-key ranges) override this to `true` and implement
306    /// [`enumerate_shards`](Self::enumerate_shards) +
307    /// [`apply_shard`](Self::apply_shard).
308    fn is_shardable(&self) -> bool {
309        false
310    }
311
312    /// Enumerate the shards this source splits into, aiming for roughly `target`
313    /// of them (a hint — the source may return fewer, e.g. when the dataset is
314    /// small, or one per natural partition regardless of `target`).
315    ///
316    /// Called **once per run** by the cluster coordinator; enumeration must be
317    /// deterministic enough that re-enumeration yields a compatible set (stable
318    /// shard ids), since it may run on more than one instance and is reconciled
319    /// by idempotent insert. May perform read-only I/O (a `LIST`, a `MIN/MAX`
320    /// query). The default returns a single whole-dataset shard
321    /// ([`ShardSpec::whole`](crate::ShardSpec::whole)), preserving today's
322    /// single-worker behavior.
323    async fn enumerate_shards(
324        &self,
325        _target: usize,
326    ) -> Result<Vec<crate::shard::ShardSpec>, FaucetError> {
327        Ok(vec![crate::shard::ShardSpec::whole()])
328    }
329
330    /// Narrow this source instance to a single shard before streaming.
331    ///
332    /// Called on the worker that claims `shard`, after construction and before
333    /// any `stream_pages` call. Like [`apply_start_bookmark`](Self::apply_start_bookmark)
334    /// this takes `&self` and is expected to record the shard behind interior
335    /// mutability (the source consults it when building its query / listing).
336    /// The default ignores the shard — a non-shardable source only ever receives
337    /// [`ShardSpec::whole`](crate::ShardSpec::whole), so ignoring it streams the
338    /// whole dataset. Implementations should accept the whole shard as a no-op.
339    async fn apply_shard(&self, _shard: &crate::shard::ShardSpec) -> Result<(), FaucetError> {
340        Ok(())
341    }
342
343    /// Compute a [`ServerDigest`](crate::diff::ServerDigest) of the rows in
344    /// `range` **inside the backend**, so a matching range ships no rows
345    /// (#701). `columns` are the compared columns; `key` the integer key the
346    /// range is over. Default `Ok(None)`: not supported, the verifier streams
347    /// the range instead. Two digests compare only when both sides report the
348    /// same `algorithm`.
349    async fn range_digest(
350        &self,
351        range: &crate::diff::KeyRange,
352        key: &str,
353        columns: &[String],
354    ) -> Result<Option<crate::diff::ServerDigest>, FaucetError> {
355        let _ = (range, key, columns);
356        Ok(None)
357    }
358
359    /// Whether this source can enumerate the datasets behind its connection
360    /// via [`discover`](Self::discover). Default: `false`. Sources backed by
361    /// an introspectable catalog (database `information_schema`, MongoDB
362    /// collections, Elasticsearch indices, object-store prefixes) override
363    /// this to `true`.
364    fn supports_discover(&self) -> bool {
365        false
366    }
367
368    /// Enumerate the datasets living behind this source's connection — one
369    /// [`DatasetDescriptor`](crate::discover::DatasetDescriptor) per table /
370    /// collection / index / prefix, each carrying a partial config override
371    /// that selects it (used by `faucet discover` to scaffold one matrix row
372    /// per dataset).
373    ///
374    /// Must be **read-only and cheap**: catalog metadata queries and listings
375    /// only, never a data scan. Descriptors must never embed credentials.
376    /// The default returns a typed "unsupported" error; override it (and
377    /// return `true` from [`supports_discover`](Self::supports_discover))
378    /// only for sources with a real catalog to introspect.
379    async fn discover(&self) -> Result<Vec<crate::discover::DatasetDescriptor>, FaucetError> {
380        Err(FaucetError::Source(format!(
381            "source '{}' does not support dataset discovery",
382            self.connector_name()
383        )))
384    }
385
386    /// The dataset (table / collection) a change record belongs to, for a
387    /// change stream that carries several tables (#731). A multi-table
388    /// `faucet mirror` runs one stream and routes each record to its table's
389    /// pipeline by this name, which must match the name the paired bulk
390    /// source's [`discover`](Self::discover) reports (e.g. `public.orders`).
391    /// `None` means the record belongs to no table (a DDL / control event) —
392    /// and is the default, so a source that does not override this cannot be
393    /// used for a multi-table mirror.
394    fn record_table(&self, _record: &Value) -> Option<String> {
395        None
396    }
397
398    /// Order two of this source's bookmarks (#731): `Some(true)` when stream
399    /// position `a` is at or before `b` (every change up to `a` is also covered
400    /// by `b`), `Some(false)` when it is not, and `None` when the source cannot
401    /// tell (the default, or two positions it cannot relate). A multi-table
402    /// mirror uses it to resume one shared stream from the earliest table and
403    /// skip, per table, pages that table has already committed.
404    fn position_le(&self, _a: &Value, _b: &Value) -> Option<bool> {
405        None
406    }
407
408    /// The earliest stream position every one of `positions` can resume from
409    /// (#731): the shared change stream restarts there and each table skips
410    /// what it has already applied. The default picks the position that
411    /// [`position_le`](Self::position_le) orders at or before all the others,
412    /// and `None` when there is none (an empty slice, or positions this source
413    /// cannot order). Sources whose positions are only partially ordered (one
414    /// cursor per capture instance, say) override it with a component-wise
415    /// minimum.
416    fn position_min(&self, positions: &[Value]) -> Option<Value> {
417        let first = positions.first()?;
418        if positions.iter().all(|p| p == first) {
419            return Some(first.clone());
420        }
421        positions
422            .iter()
423            .find(|cand| {
424                positions
425                    .iter()
426                    .all(|other| self.position_le(cand, other) == Some(true))
427            })
428            .cloned()
429    }
430
431    /// Stable identifier used as the `connector` label on metrics and the
432    /// `connector` attribute on spans. Defaults to the final segment of
433    /// `std::any::type_name::<Self>()`, e.g. `"RestSource"`. Built-in
434    /// connectors override with a short, friendly snake_case name (e.g.
435    /// `"rest"`). Must return a non-empty string; observability decorators
436    /// fall back to `"unknown"` in release builds if it is empty (and
437    /// `debug_assert!` in debug builds).
438    fn connector_name(&self) -> &'static str {
439        crate::observability::strip_type_name(std::any::type_name::<Self>())
440    }
441
442    /// Receive a pre-labelled handle for counting **upstream round trips**
443    /// (#638) — calls this connector makes to its own backend.
444    ///
445    /// The pipeline calls this once before streaming, because it is the only
446    /// place that knows the `pipeline` / `row` / `connector` labels every other
447    /// metric carries. A connector opts in by storing the handle behind
448    /// interior mutability (the pattern the REST source already uses for
449    /// `runtime_start`) and calling `recorder.record("<op>")` at each real
450    /// backend call — including retries, which are real round trips.
451    ///
452    /// Defaulted to a no-op: an uninstrumented connector emits nothing, so
453    /// instrumentation lands connector by connector with no behaviour change
454    /// in between, and a third-party connector is unaffected.
455    fn set_roundtrip_recorder(
456        &self,
457        _recorder: std::sync::Arc<crate::observability::RoundtripRecorder>,
458    ) {
459    }
460
461    /// Receive the run clock (#769): the instant `${now.*}` renders from —
462    /// `faucet run --clock`, a schedule tick, a backfill unit's start — so a
463    /// source that bounds reads by "now" (REST window slicing) reproduces the
464    /// run as of that instant rather than the wall clock.
465    ///
466    /// Called once before streaming. Defaulted to a no-op; a connector opts in
467    /// by storing the instant behind interior mutability.
468    fn set_run_clock(&self, _now: chrono::DateTime<chrono::Utc>) {}
469
470    /// Logical dataset identity for lineage emission, following OpenLineage
471    /// naming conventions (<https://openlineage.io/docs/spec/naming>).
472    ///
473    /// The default returns `"<connector_name>://unknown"`. Built-in connectors
474    /// override with a credential-free URI derived from their config. Strip any
475    /// credentials with [`redact_uri_credentials`](crate::redact_uri_credentials).
476    /// Informational metadata only — never used for I/O.
477    fn dataset_uri(&self) -> String {
478        format!("{}://unknown", self.connector_name())
479    }
480
481    /// Run a fast, non-mutating preflight probe (used by `faucet doctor`).
482    ///
483    /// The default pulls a **single page** via
484    /// [`stream_pages`](Self::stream_pages) and reports success/failure — it
485    /// exercises the real read path (DNS, TLS, auth, the first request, the
486    /// first-record decode) but never paginates the full dataset and never
487    /// repeats. The page stream is dropped immediately after the first page.
488    ///
489    /// Sources whose first page *blocks* waiting for inbound data (webhook,
490    /// websocket) or has *side effects* (CDC consuming WAL) override this with a
491    /// cheaper, side-effect-free probe. Probe-level failures are returned as a
492    /// [`ProbeStatus::Fail`](crate::check::ProbeStatus) inside `Ok(report)`.
493    async fn check(
494        &self,
495        ctx: &crate::check::CheckContext,
496    ) -> Result<crate::check::CheckReport, FaucetError> {
497        use crate::check::{CheckReport, Probe};
498        use futures::StreamExt;
499
500        let empty = std::collections::HashMap::new();
501        let start = std::time::Instant::now();
502        let mut pages = self.stream_pages(&empty, 1);
503        let probe = match tokio::time::timeout(ctx.timeout, pages.next()).await {
504            Err(_) => Probe::fail("read", start.elapsed(), "timed out fetching first page"),
505            Ok(None) | Ok(Some(Ok(_))) => Probe::pass("read", start.elapsed()),
506            Ok(Some(Err(e))) => Probe::fail("read", start.elapsed(), e.to_string()),
507        };
508        Ok(CheckReport::single(probe))
509    }
510}
511
512/// Per-row outcome from [`Sink::write_batch_partial`].
513///
514/// `Ok(())` — the row was durably written to the sink.
515/// `Err(_)` — the row failed; the pipeline will route it to the DLQ when
516/// one is configured.
517pub type RowOutcome = Result<(), FaucetError>;
518
519/// A sink writes records to an external system.
520#[async_trait]
521pub trait Sink: Send + Sync {
522    /// Write a batch of records to the destination.
523    ///
524    /// Returns the number of records successfully written.
525    async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError>;
526
527    /// Flush any buffered data to the destination.
528    ///
529    /// The default implementation is a no-op (suitable for sinks that
530    /// write immediately in `write_batch`).
531    async fn flush(&self) -> Result<(), FaucetError> {
532        Ok(())
533    }
534
535    /// Write a batch and report per-row outcomes.
536    ///
537    /// Sinks whose underlying API exposes per-row results (BigQuery
538    /// `insertAll`, Elasticsearch `_bulk`) override this. The default
539    /// implementation delegates to [`Self::write_batch`] and maps a single success
540    /// onto a uniform all-`Ok(())` vector. An outer failure is bubbled up
541    /// unchanged so the pipeline's DLQ router can apply its `on_batch_error`
542    /// policy at a single decision point.
543    async fn write_batch_partial(&self, records: &[Value]) -> Result<Vec<RowOutcome>, FaucetError> {
544        self.write_batch(records).await?;
545        Ok(records.iter().map(|_| Ok(())).collect())
546    }
547
548    /// What a failed batch write leaves behind (#737): nothing
549    /// ([`Atomic`](crate::BatchAtomicity::Atomic)), per-row outcomes with
550    /// nothing committed on an outer `Err`
551    /// ([`PerRow`](crate::BatchAtomicity::PerRow)), or possibly some rows
552    /// ([`BestEffort`](crate::BatchAtomicity::BestEffort), the default).
553    ///
554    /// `on_batch_error: dlq_all` is refused on a best-effort sink unless it
555    /// writes by key or the user opts in to the duplicates. Decorators must
556    /// forward this.
557    fn batch_atomicity(&self) -> crate::dlq::BatchAtomicity {
558        crate::dlq::BatchAtomicity::BestEffort
559    }
560
561    /// Whether this sink can consume **columnar** (`arrow::RecordBatch`) writes
562    /// via [`write_batch_columnar`](Self::write_batch_columnar) without first
563    /// converting to `Value`. Default: `false`.
564    ///
565    /// The pipeline takes the columnar fast path only when *both* the source and
566    /// sink return `true` (RFC 0002 / #375).
567    #[cfg(feature = "arrow")]
568    fn supports_columnar(&self) -> bool {
569        false
570    }
571
572    /// Write a columnar `RecordBatch` to the destination, returning the number of
573    /// rows written.
574    ///
575    /// The default converts the batch to `Value` rows via the
576    /// [`columnar`](crate::columnar) shim and delegates to
577    /// [`write_batch`](Self::write_batch), so every sink participates correctly
578    /// even without a native columnar path. Sinks that write Arrow/Parquet
579    /// directly override this — and return `true` from
580    /// [`supports_columnar`](Self::supports_columnar) — to skip the conversion.
581    #[cfg(feature = "arrow")]
582    async fn write_batch_columnar(
583        &self,
584        batch: &arrow::array::RecordBatch,
585    ) -> Result<usize, FaucetError> {
586        let rows = crate::columnar::record_batch_to_values(batch)?;
587        self.write_batch(&rows).await
588    }
589
590    /// Native byte-passthrough load mechanisms this sink offers, each naming
591    /// the wire format and the write modes it can honor (#633). Default:
592    /// `vec![]` (no fast path; use [`write_batch`](Self::write_batch)).
593    ///
594    /// The pipeline takes the native path only when a source advertises a
595    /// format this returns and the run passes the pipeline-owned gates — no
596    /// transforms, no quality/contract/masking pass, no DLQ, at-least-once
597    /// delivery, no upsert/delete `write_mode`, and (for overwrite) an
598    /// all-or-nothing load session finalized by a single terminal `flush` —
599    /// see [`plan_native_transfer`](crate::plan_native_transfer). A capability
600    /// cannot waive those gates.
601    fn native_load_capabilities(&self) -> Vec<crate::native::NativeLoadCapability> {
602        Vec::new()
603    }
604
605    /// Bulk-load one native-format byte batch directly (#633), returning the
606    /// number of rows written.
607    ///
608    /// Only invoked when the pipeline negotiated one of this sink's
609    /// [`native_load_capabilities`](Self::native_load_capabilities) and the
610    /// pipeline-owned gates passed. `ctx.first_batch` lets an overwrite sink
611    /// truncate on the first load and append thereafter. The default returns a
612    /// typed "unsupported" error so a sink that advertises a capability but
613    /// forgets to override this fails loudly.
614    ///
615    /// **Contract:** an **empty payload must be a successful no-op** returning
616    /// `Ok(0)` — sources emit a trailing empty batch to carry the final
617    /// bookmark, and an error here would break bookmark persistence. Under
618    /// `Overwrite`, nothing may become visible in the destination until the
619    /// terminal [`flush`](Self::flush) — the pipeline flushes exactly once, on
620    /// success only, so a mechanism that commits per batch must omit
621    /// `Overwrite` from its capability's `write_modes`.
622    async fn load_native(
623        &self,
624        batch: crate::native::NativeBatch,
625        scope: &str,
626        ctx: crate::native::NativeLoadContext,
627    ) -> Result<usize, FaucetError> {
628        let _ = (batch, scope, ctx);
629        Err(FaucetError::Sink(format!(
630            "sink '{}' does not support native byte loading (load_native)",
631            self.connector_name()
632        )))
633    }
634
635    /// Whether this sink can durably commit a page's rows **and** a commit token
636    /// in a single atomic transaction. Default: `false` (at-least-once only).
637    ///
638    /// Only return `true` from a sink that genuinely commits both atomically —
639    /// see [`write_batch_idempotent`](Self::write_batch_idempotent). The pipeline
640    /// rejects `DeliveryMode::ExactlyOnce` against a sink that returns `false`.
641    fn supports_idempotent_writes(&self) -> bool {
642        false
643    }
644
645    /// The strongest delivery guarantee this sink can uphold — see
646    /// [`SinkGuarantee`](crate::SinkGuarantee).
647    ///
648    /// The default derives from the two back-compat primitives:
649    /// [`supports_idempotent_writes`](Self::supports_idempotent_writes) →
650    /// `AtomicWatermark`, else an upsert-capable
651    /// [`supported_write_modes`](Self::supported_write_modes) → `KeyedUpsert`,
652    /// else `AtLeastOnce`. Existing connectors that override only the
653    /// primitives automatically advertise the right capability here.
654    fn sink_guarantee(&self) -> crate::idempotency::SinkGuarantee {
655        if self.supports_idempotent_writes() {
656            crate::idempotency::SinkGuarantee::AtomicWatermark
657        } else if self
658            .supported_write_modes()
659            .contains(&crate::write_mode::WriteMode::Upsert)
660        {
661            crate::idempotency::SinkGuarantee::KeyedUpsert
662        } else {
663            crate::idempotency::SinkGuarantee::AtLeastOnce
664        }
665    }
666
667    /// Whether a **plain [`write_batch`](Self::write_batch) is safe to replay**
668    /// — i.e. re-sending the same page after an ambiguous failure converges
669    /// instead of duplicating rows. Default: `false`.
670    ///
671    /// This is the property the pipeline needs before it may retry a
672    /// non-idempotent write, and it is **not** the same as
673    /// [`supports_idempotent_writes`](Self::supports_idempotent_writes), which
674    /// promises only that rows *and a commit token* can be committed in one
675    /// transaction (the [`write_batch_idempotent`](Self::write_batch_idempotent)
676    /// path). A sink can offer that token protocol and still have a plain
677    /// `write_batch` that is a bare multi-row `INSERT` — replaying it after the
678    /// server committed but the response was lost duplicates every row, which
679    /// is this repo's worst bug class (F29/F32).
680    ///
681    /// The default implementation is therefore conservative and derives from
682    /// the live config: a keyed upsert/delete converges on replay, anything
683    /// else does not. Override only for a sink whose plain write is genuinely
684    /// replay-safe by construction.
685    fn write_batch_is_replay_safe(&self) -> bool {
686        self.dedups_by_key()
687    }
688
689    /// Whether this sink instance is **configured** to dedup by key — i.e.
690    /// `write_mode: upsert` (or `delete`) with a non-empty `key`, so
691    /// re-applying a record with the same key converges instead of
692    /// duplicating. Default: `false`.
693    ///
694    /// Distinct from [`sink_guarantee`](Self::sink_guarantee) (capability):
695    /// this reflects the *live config*. Sinks that flatten a
696    /// [`WriteSpec`](crate::write_mode::WriteSpec) into their config override
697    /// it as `self.config.write.dedups_by_key()`. The pipeline consults it to
698    /// derive the keyed-upsert effectively-once mechanism at run time.
699    fn dedups_by_key(&self) -> bool {
700        false
701    }
702
703    /// Write modes this sink can apply. Default: append-only. Sinks that
704    /// implement key-based merge override this to include
705    /// [`WriteMode::Upsert`](crate::write_mode::WriteMode) /
706    /// [`WriteMode::Delete`](crate::write_mode::WriteMode). The CLI rejects a
707    /// configured mode that is not in this set at config-load time.
708    fn supported_write_modes(&self) -> &'static [crate::write_mode::WriteMode] {
709        &[crate::write_mode::WriteMode::Append]
710    }
711
712    /// The sink's live destination schema as an `infer_schema`-shaped object
713    /// (`{"type":"object","properties":{ <col>: <type-fragment>, … }}`), or
714    /// `None` for a schemaless sink or a target that does not exist yet.
715    ///
716    /// Used by the schema-drift policy to diff each page's shape against the
717    /// real destination. Default: `Ok(None)` (drift handling is inert).
718    async fn current_schema(&self) -> Result<Option<Value>, FaucetError> {
719        Ok(None)
720    }
721
722    /// Whether this sink can apply additive/widening DDL via
723    /// [`evolve_schema`](Self::evolve_schema). Default: `false`. The CLI rejects
724    /// `on_drift: evolve` against a sink that returns `false` at config-load.
725    fn supports_schema_evolution(&self) -> bool {
726        false
727    }
728
729    /// Apply an additive schema evolution (new columns, lossless widenings,
730    /// nullability relaxations) to the destination. MUST be idempotent
731    /// (`ADD COLUMN IF NOT EXISTS` semantics) so concurrent runs converge.
732    ///
733    /// Default: a typed "unsupported" error. Override only when the backend
734    /// supports in-place additive DDL (and return `true` from
735    /// `supports_schema_evolution`).
736    async fn evolve_schema(
737        &self,
738        evolution: &crate::drift::SchemaEvolution,
739    ) -> Result<(), FaucetError> {
740        let _ = evolution;
741        Err(FaucetError::Sink(format!(
742            "sink '{}' does not support schema evolution",
743            self.connector_name()
744        )))
745    }
746
747    /// Whether this sink can delete a scoped set of rows for scoped cleanup
748    /// (#478). Default `false`; the upsert-capable sinks override it.
749    fn supports_cleanup(&self) -> bool {
750        false
751    }
752
753    /// Whether this sink can bulk-load via an object-store stage — write the
754    /// page to S3/GCS/Azure, then pull it with the warehouse's native load
755    /// command (`COPY` / `COPY INTO` / `s3()` table function) (#528). Default
756    /// `false`; warehouse sinks that honour a [`StagingSpec`](crate::StagingSpec)
757    /// override it. Object-safe (no args, no generics) so `Box<dyn Sink>` is
758    /// unaffected.
759    fn supports_staged_load(&self) -> bool {
760        false
761    }
762
763    /// Delete rows matching `scope` whose key is **not** in `seen`.
764    ///
765    /// Called at most once per invocation, only after the run completed
766    /// successfully and uncancelled — see [`crate::cleanup`] for why the timing
767    /// is load-bearing. `scope` is a set of equality predicates in destination
768    /// column terms, AND-ed together; `seen` holds the key tuples this run wrote.
769    ///
770    /// Returns the number of rows deleted. Implementations **must** be
771    /// all-or-nothing where the backend allows it: a partial delete would remove
772    /// rows the run actually wrote.
773    ///
774    /// The default is a typed "unsupported" error, so no existing or third-party
775    /// connector breaks.
776    async fn cleanup_scope(
777        &self,
778        scope: &std::collections::BTreeMap<String, Value>,
779        seen: &crate::cleanup::SeenKeys,
780    ) -> Result<u64, FaucetError> {
781        let _ = (scope, seen);
782        Err(FaucetError::Sink(format!(
783            "sink '{}' does not support scoped cleanup",
784            self.connector_name()
785        )))
786    }
787
788    /// Write `records` AND durably record `token` for `scope`, atomically.
789    ///
790    /// `scope` namespaces the watermark (the pipeline passes the per-row state
791    /// key, e.g. `"{name}::{row_id}"`). `token` is a monotonic, fixed-width
792    /// string (see [`format_token`](crate::format_token)).
793    ///
794    /// The default is **not** idempotent: it ignores the token and delegates to
795    /// [`write_batch`](Self::write_batch). Override only when the commit is
796    /// genuinely atomic (and return `true` from `supports_idempotent_writes`).
797    async fn write_batch_idempotent(
798        &self,
799        records: &[Value],
800        scope: &str,
801        token: &str,
802    ) -> Result<usize, FaucetError> {
803        let _ = (scope, token);
804        self.write_batch(records).await
805    }
806
807    /// The last token durably committed for `scope`, or `None` if this sink has
808    /// never committed under that scope. Default: `None`.
809    async fn last_committed_token(&self, scope: &str) -> Result<Option<String>, FaucetError> {
810        let _ = scope;
811        Ok(None)
812    }
813
814    /// Whether this sink instance is configured for full-destination
815    /// replacement ([`WriteMode::Overwrite`](crate::write_mode::WriteMode)).
816    ///
817    /// The pipeline consults this to drive the overwrite lifecycle:
818    /// [`begin_overwrite`](Self::begin_overwrite) before the first page, then
819    /// [`commit_overwrite`](Self::commit_overwrite) once the run finishes
820    /// successfully, or [`abort_overwrite`](Self::abort_overwrite) on
821    /// failure/cancel. Sinks that flatten a [`WriteSpec`](crate::write_mode::WriteSpec)
822    /// into their config return `self.config.write.is_overwrite()`. Default `false`.
823    fn is_overwrite(&self) -> bool {
824        false
825    }
826
827    /// Prepare a staging target for an overwrite run.
828    ///
829    /// Called once, before the first [`write_batch`](Self::write_batch), only
830    /// when [`is_overwrite`](Self::is_overwrite) is true. The sink stages this
831    /// run's writes (a temp table / new index / temp prefix) so the existing
832    /// destination is untouched until the run succeeds. Subsequent
833    /// `write_batch` calls for this sink must land in the staging target.
834    ///
835    /// Default: a typed "unsupported" error, so a sink that advertises
836    /// `WriteMode::Overwrite` but forgets to implement the lifecycle fails
837    /// loudly rather than silently appending. Never called for sinks whose
838    /// `is_overwrite()` is false.
839    async fn begin_overwrite(&self) -> Result<(), FaucetError> {
840        Err(FaucetError::Sink(format!(
841            "sink '{}' does not support write_mode: overwrite",
842            self.connector_name()
843        )))
844    }
845
846    /// Atomically replace the destination with the staged data.
847    ///
848    /// Called **once, only after the run completed successfully and
849    /// uncancelled**. Implementations MUST swap staging → destination
850    /// atomically (or as close as the backend allows) so a reader never sees a
851    /// half-replaced dataset, and MUST NOT have destroyed the prior contents
852    /// before this point — a failed run leaves the old data in place.
853    ///
854    /// Default: a typed "unsupported" error (unreachable for a correct sink
855    /// whose `is_overwrite()` is false).
856    async fn commit_overwrite(&self) -> Result<(), FaucetError> {
857        Err(FaucetError::Sink(format!(
858            "sink '{}' does not support write_mode: overwrite",
859            self.connector_name()
860        )))
861    }
862
863    /// Whether this sink can undo a run it wrote (#706) through
864    /// [`rollback_run`](Self::rollback_run). Default `false`.
865    fn supports_rollback(&self) -> bool {
866        false
867    }
868
869    /// Undo everything `run_id` wrote to this destination, per
870    /// `opts.mode`: delete the run's rows (append), restore journaled
871    /// before-images (upsert / delete), or swap the kept previous table back
872    /// (overwrite). Must be all-or-nothing per destination where the backend
873    /// allows it, and must leave the destination untouched when `opts.dry_run`
874    /// is set or when conflicts block it (see
875    /// [`RollbackOutcome`](crate::rollback::RollbackOutcome)).
876    ///
877    /// Default: a typed "unsupported" error.
878    async fn rollback_run(
879        &self,
880        run_id: &str,
881        opts: &crate::rollback::RollbackOptions,
882    ) -> Result<crate::rollback::RollbackOutcome, FaucetError> {
883        let _ = (run_id, opts);
884        Err(FaucetError::Sink(format!(
885            "sink '{}' does not support rollback",
886            self.connector_name()
887        )))
888    }
889
890    /// Drop whatever this sink kept to make `run_id` undoable (journal rows, a
891    /// previous table) once the run is past the retention window. Default:
892    /// no-op.
893    async fn forget_run(&self, run_id: &str) -> Result<(), FaucetError> {
894        let _ = run_id;
895        Ok(())
896    }
897
898    /// Rewind the exactly-once watermark for `scope` to `token` (`None` clears
899    /// it), so a rolled-back run's pages are not skipped on the next run.
900    /// Default: a typed "unsupported" error.
901    async fn rewind_commit_token(
902        &self,
903        scope: &str,
904        token: Option<&str>,
905    ) -> Result<(), FaucetError> {
906        let _ = (scope, token);
907        Err(FaucetError::Sink(format!(
908            "sink '{}' does not support rewinding its commit token",
909            self.connector_name()
910        )))
911    }
912
913    /// A source `(kind, config)` that reads this destination back — what
914    /// `faucet verify` (#701) compares the pipeline's source against. Default
915    /// `None`: the user names the destination reader in `verify.destination`.
916    /// The config must not need a credential the sink config lacks.
917    fn readback_source(&self) -> Option<(String, Value)> {
918        None
919    }
920
921    /// Discard the staging target after a failed or cancelled overwrite run.
922    ///
923    /// Called (best-effort) when an overwrite run does not reach
924    /// [`commit_overwrite`](Self::commit_overwrite). The destination must be
925    /// left exactly as it was before the run. Default: no-op — a leftover
926    /// staging object is untidy but never data loss, so a sink may skip it.
927    async fn abort_overwrite(&self) -> Result<(), FaucetError> {
928        Ok(())
929    }
930
931    /// Finalize this run's output after [`Pipeline::run`](crate::Pipeline::run)
932    /// finished **successfully and uncancelled**, on every transfer path
933    /// (`Value`, columnar, native), after the terminal flush.
934    ///
935    /// A sink whose destination must reflect *this* run even when it wrote
936    /// nothing uses it: the file sinks with `append: false` truncate (or, for
937    /// a fixed-path Parquet file, remove) the previous run's output here when
938    /// no record arrived, so a source that became empty never leaves stale
939    /// rows presented as current (#753). Never called after a failed or
940    /// cancelled run, so the previous good output survives those. Decorators
941    /// must forward it. Default: no-op.
942    async fn complete_run(&self) -> Result<(), FaucetError> {
943        Ok(())
944    }
945
946    /// Whether an overwrite staging object (the `…__faucet_ovw` table or
947    /// collection a [`begin_overwrite`](Self::begin_overwrite) creates) exists
948    /// right now — a read-only probe used by `faucet status --probe` to report
949    /// staging a crashed or aborted overwrite left behind.
950    ///
951    /// `Ok(None)` means the sink cannot tell (the default); overwrite-capable
952    /// sinks return `Ok(Some(_))`.
953    async fn overwrite_staging_exists(&self) -> Result<Option<bool>, FaucetError> {
954        Ok(None)
955    }
956
957    /// Return a JSON Schema describing the configuration this sink accepts.
958    ///
959    /// The schema is auto-generated from the config struct using `schemars`.
960    /// Callers can inspect it to discover required fields, types, defaults,
961    /// and descriptions before constructing the sink.
962    ///
963    /// The default returns an empty object schema.
964    fn config_schema(&self) -> Value {
965        serde_json::json!({"type": "object", "properties": {}})
966    }
967
968    /// Stable identifier used as the `connector` label on metrics and the
969    /// `connector` attribute on spans. See `Source::connector_name`.
970    fn connector_name(&self) -> &'static str {
971        crate::observability::strip_type_name(std::any::type_name::<Self>())
972    }
973
974    /// Receive a pre-labelled handle for counting **upstream round trips**
975    /// (#638). See [`Source::set_roundtrip_recorder`]; defaulted to a no-op.
976    fn set_roundtrip_recorder(
977        &self,
978        _recorder: std::sync::Arc<crate::observability::RoundtripRecorder>,
979    ) {
980    }
981
982    /// Logical dataset identity for lineage emission, following OpenLineage
983    /// naming conventions (<https://openlineage.io/docs/spec/naming>).
984    ///
985    /// The default returns `"<connector_name>://unknown"`. Built-in connectors
986    /// override with a credential-free URI derived from their config. Strip any
987    /// credentials with [`redact_uri_credentials`](crate::redact_uri_credentials).
988    /// Informational metadata only — never used for I/O.
989    fn dataset_uri(&self) -> String {
990        format!("{}://unknown", self.connector_name())
991    }
992
993    /// The concrete **local files** this sink instance opened during the run.
994    ///
995    /// The provenance record faucet's local-output retention GC (#587) deletes
996    /// from: it may only remove files faucet recorded here, never a glob or a
997    /// directory. A sink that writes local files (jsonl, csv, parquet to a local
998    /// destination) accumulates a
999    /// [`LocalOutputLog`](crate::local_outputs::LocalOutputLog) as it opens them
1000    /// and returns its snapshot; every other sink keeps the empty default, which
1001    /// simply means "nothing local to collect".
1002    ///
1003    /// Two properties the GC depends on, so an implementation must preserve them:
1004    ///
1005    /// 1. **Every file, individually named.** A rolling parquet sink returns one
1006    ///    entry per UUID-named part — the directory itself is not an output.
1007    /// 2. **[`LocalOutput::pre_existing`](crate::local_outputs::LocalOutput) is
1008    ///    honest**, captured before the first open. A file faucet appended to but
1009    ///    did not create is never collected.
1010    ///
1011    /// Called after the run (and by `faucet cleanup`), never on the data path, so
1012    /// an implementation may take a lock. Decorating sinks **must forward** this
1013    /// to their inner sink, exactly as they do
1014    /// [`dataset_uri`](Self::dataset_uri) — a decorator that returns the default
1015    /// hides its inner sink's files from the GC and they are never reclaimed.
1016    async fn local_outputs(&self) -> Vec<crate::local_outputs::LocalOutput> {
1017        Vec::new()
1018    }
1019
1020    /// Run a fast, non-mutating preflight probe (used by `faucet doctor`).
1021    ///
1022    /// Unlike sources, a sink has no non-mutating "first page" equivalent
1023    /// (`write_batch` mutates the destination), so the default returns
1024    /// [`CheckReport::not_implemented`](crate::check::CheckReport::not_implemented).
1025    /// Built-in sinks override this with a connect / auth / metadata probe.
1026    ///
1027    /// The probe **MUST be idempotent and side-effect-free** — no inserts, no
1028    /// residual rows or objects — and must never put credentials or connection
1029    /// strings in a probe `reason`/`hint`.
1030    async fn check(
1031        &self,
1032        _ctx: &crate::check::CheckContext,
1033    ) -> Result<crate::check::CheckReport, FaucetError> {
1034        Ok(crate::check::CheckReport::not_implemented())
1035    }
1036}
1037
1038#[cfg(test)]
1039mod tests {
1040    use super::*;
1041    use serde_json::json;
1042
1043    // ── Mock Source ──────────────────────────────────────────────────────────
1044
1045    struct MockSource {
1046        records: Vec<Value>,
1047    }
1048
1049    #[async_trait]
1050    impl Source for MockSource {
1051        async fn fetch_with_context(
1052            &self,
1053            _context: &std::collections::HashMap<String, Value>,
1054        ) -> Result<Vec<Value>, FaucetError> {
1055            Ok(self.records.clone())
1056        }
1057    }
1058
1059    #[test]
1060    fn default_migrate_state_accepts_only_its_own_schema() {
1061        let src = MockSource { records: vec![] };
1062        assert_eq!(src.state_schema(), 0);
1063        assert_eq!(
1064            src.migrate_state(0, json!({"id": 1})).unwrap(),
1065            json!({"id": 1})
1066        );
1067        let err = src.migrate_state(3, json!({})).unwrap_err().to_string();
1068        assert!(
1069            err.contains("no migration from bookmark schema 3 to 0"),
1070            "{err}"
1071        );
1072    }
1073
1074    #[test]
1075    fn default_multi_table_hooks_route_nothing_and_order_nothing() {
1076        let src = MockSource { records: vec![] };
1077        assert_eq!(src.record_table(&json!({"table": "t"})), None);
1078        assert_eq!(src.position_le(&json!(1), &json!(2)), None);
1079        assert_eq!(src.position_min(&[]), None);
1080        assert_eq!(src.position_min(&[json!(3), json!(3)]), Some(json!(3)));
1081        assert_eq!(src.position_min(&[json!(1), json!(2)]), None);
1082    }
1083
1084    #[tokio::test]
1085    async fn ordered_double_fetches_nothing() {
1086        use crate::Source as _;
1087        assert!(
1088            OrderedSource
1089                .fetch_with_context(&Default::default())
1090                .await
1091                .unwrap()
1092                .is_empty()
1093        );
1094    }
1095
1096    struct OrderedSource;
1097
1098    #[async_trait]
1099    impl Source for OrderedSource {
1100        async fn fetch_with_context(
1101            &self,
1102            _context: &std::collections::HashMap<String, Value>,
1103        ) -> Result<Vec<Value>, FaucetError> {
1104            Ok(vec![])
1105        }
1106
1107        fn position_le(&self, a: &Value, b: &Value) -> Option<bool> {
1108            Some(a.as_u64()? <= b.as_u64()?)
1109        }
1110    }
1111
1112    #[test]
1113    fn default_position_min_uses_position_le() {
1114        let src = OrderedSource;
1115        assert_eq!(
1116            src.position_min(&[json!(5), json!(2), json!(9)]),
1117            Some(json!(2))
1118        );
1119        assert_eq!(src.position_min(&[json!(5), json!("x")]), None);
1120    }
1121
1122    struct IncrementalSource {
1123        records: Vec<Value>,
1124        bookmark: Value,
1125    }
1126
1127    #[async_trait]
1128    impl Source for IncrementalSource {
1129        async fn fetch_with_context(
1130            &self,
1131            _context: &std::collections::HashMap<String, Value>,
1132        ) -> Result<Vec<Value>, FaucetError> {
1133            Ok(self.records.clone())
1134        }
1135
1136        async fn fetch_with_context_incremental(
1137            &self,
1138            _context: &std::collections::HashMap<String, Value>,
1139        ) -> Result<(Vec<Value>, Option<Value>), FaucetError> {
1140            Ok((self.records.clone(), Some(self.bookmark.clone())))
1141        }
1142    }
1143
1144    struct FailingSource;
1145
1146    #[async_trait]
1147    impl Source for FailingSource {
1148        async fn fetch_with_context(
1149            &self,
1150            _context: &std::collections::HashMap<String, Value>,
1151        ) -> Result<Vec<Value>, FaucetError> {
1152            Err(FaucetError::Auth("no credentials".into()))
1153        }
1154    }
1155
1156    // ── Mock Sink ───────────────────────────────────────────────────────────
1157
1158    struct MockSink {
1159        written: std::sync::Mutex<Vec<Value>>,
1160    }
1161
1162    impl MockSink {
1163        fn new() -> Self {
1164            Self {
1165                written: std::sync::Mutex::new(Vec::new()),
1166            }
1167        }
1168    }
1169
1170    #[async_trait]
1171    impl Sink for MockSink {
1172        async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
1173            let mut w = self.written.lock().unwrap();
1174            w.extend(records.iter().cloned());
1175            Ok(records.len())
1176        }
1177    }
1178
1179    struct FailingSink;
1180
1181    #[async_trait]
1182    impl Sink for FailingSink {
1183        async fn write_batch(&self, _records: &[Value]) -> Result<usize, FaucetError> {
1184            Err(FaucetError::Sink("write failed".into()))
1185        }
1186    }
1187
1188    // ── Native byte-passthrough defaults (#633) ──────────────────────────────
1189
1190    /// A source/sink pair that advertises nothing native, so the defaulted
1191    /// trait methods are what answer.
1192    struct PlainSink;
1193
1194    #[async_trait]
1195    impl Sink for PlainSink {
1196        async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
1197            Ok(records.len())
1198        }
1199    }
1200
1201    #[tokio::test]
1202    async fn native_defaults_advertise_nothing_and_fail_loudly() {
1203        use futures::StreamExt as _;
1204
1205        // A source that never overrides the native methods advertises no
1206        // formats, so the pipeline can never negotiate the byte path…
1207        let src = MockSource { records: vec![] };
1208        assert!(src.native_output_formats().is_empty());
1209        // …and calling it anyway yields one typed error rather than silently
1210        // producing an empty stream (which would look like "no rows").
1211        let ctx = std::collections::HashMap::new();
1212        let mut batches = src.stream_native(&ctx, crate::native::NativeFormat::NdJson, 10);
1213        let err = batches
1214            .next()
1215            .await
1216            .expect("one item")
1217            .expect_err("default must error");
1218        assert!(
1219            err.to_string()
1220                .contains("does not support native byte streaming"),
1221            "{err}"
1222        );
1223        assert!(batches.next().await.is_none(), "exactly one item");
1224
1225        // Same on the sink side: no capabilities, and `load_native` is a typed
1226        // error so a sink that advertises but forgets to implement is obvious.
1227        let sink = PlainSink;
1228        assert!(sink.native_load_capabilities().is_empty());
1229        let err = sink
1230            .load_native(
1231                crate::native::NativeBatch::bytes(
1232                    crate::native::NativeFormat::NdJson,
1233                    b"{}\n".to_vec(),
1234                ),
1235                "scope",
1236                crate::native::NativeLoadContext {
1237                    write_mode: crate::write_mode::WriteMode::Append,
1238                    first_batch: true,
1239                },
1240            )
1241            .await
1242            .expect_err("default must error");
1243        assert!(
1244            err.to_string()
1245                .contains("does not support native byte loading"),
1246            "{err}"
1247        );
1248    }
1249
1250    #[tokio::test]
1251    async fn default_rollback_and_readback_hooks() {
1252        let sink = MockSink::new();
1253        assert!(!sink.supports_rollback());
1254        let opts = crate::rollback::RollbackOptions {
1255            run_id_column: "_faucet_run_id".into(),
1256            mode: crate::rollback::RollbackMode::Append,
1257            force: false,
1258            dry_run: false,
1259        };
1260        let err = sink.rollback_run("r1", &opts).await.unwrap_err();
1261        assert!(
1262            err.to_string().contains("does not support rollback"),
1263            "{err}"
1264        );
1265        assert!(sink.forget_run("r1").await.is_ok());
1266        let err = sink.rewind_commit_token("s", None).await.unwrap_err();
1267        assert!(err.to_string().contains("commit token"), "{err}");
1268        assert!(sink.readback_source().is_none());
1269        let source = MockSource { records: vec![] };
1270        let digest = source
1271            .range_digest(&crate::diff::KeyRange::ALL, "id", &["v".into()])
1272            .await
1273            .unwrap();
1274        assert!(digest.is_none(), "not supported by default");
1275    }
1276
1277    #[tokio::test]
1278    async fn default_overwrite_methods_reject_or_noop() {
1279        // A sink that does not opt into overwrite (MockSink uses the defaults):
1280        // `is_overwrite` is false, begin/commit return the typed "unsupported"
1281        // error, and abort is a no-op success.
1282        let sink = MockSink::new();
1283        assert!(!sink.is_overwrite());
1284        assert!(sink.begin_overwrite().await.is_err());
1285        assert!(sink.commit_overwrite().await.is_err());
1286        assert!(sink.abort_overwrite().await.is_ok());
1287        assert_eq!(sink.overwrite_staging_exists().await.unwrap(), None);
1288    }
1289
1290    // ── Source tests ────────────────────────────────────────────────────────
1291
1292    #[tokio::test]
1293    async fn source_fetch_all_returns_records() {
1294        let source = MockSource {
1295            records: vec![json!({"id": 1}), json!({"id": 2})],
1296        };
1297        let records = source.fetch_all().await.unwrap();
1298        assert_eq!(records.len(), 2);
1299        assert_eq!(records[0]["id"], 1);
1300    }
1301
1302    #[tokio::test]
1303    async fn source_fetch_all_empty() {
1304        let source = MockSource { records: vec![] };
1305        let records = source.fetch_all().await.unwrap();
1306        assert!(records.is_empty());
1307    }
1308
1309    #[tokio::test]
1310    async fn source_default_incremental_returns_none_bookmark() {
1311        let source = MockSource {
1312            records: vec![json!({"id": 1})],
1313        };
1314        let (records, bookmark) = source.fetch_all_incremental().await.unwrap();
1315        assert_eq!(records.len(), 1);
1316        assert!(bookmark.is_none());
1317    }
1318
1319    #[tokio::test]
1320    async fn source_custom_incremental_returns_bookmark() {
1321        let source = IncrementalSource {
1322            records: vec![json!({"id": 1})],
1323            bookmark: json!("2024-12-01"),
1324        };
1325        let (records, bookmark) = source.fetch_all_incremental().await.unwrap();
1326        assert_eq!(records.len(), 1);
1327        assert_eq!(bookmark, Some(json!("2024-12-01")));
1328    }
1329
1330    #[tokio::test]
1331    async fn source_error_propagates() {
1332        let source = FailingSource;
1333        let result = source.fetch_all().await;
1334        assert!(result.is_err());
1335        assert!(matches!(result, Err(FaucetError::Auth(_))));
1336    }
1337
1338    #[tokio::test]
1339    async fn source_as_trait_object() {
1340        let source: Box<dyn Source> = Box::new(MockSource {
1341            records: vec![json!({"id": 42})],
1342        });
1343        let records = source.fetch_all().await.unwrap();
1344        assert_eq!(records[0]["id"], 42);
1345    }
1346
1347    // ── Sink tests ──────────────────────────────────────────────────────────
1348
1349    #[tokio::test]
1350    async fn sink_write_batch_returns_count() {
1351        let sink = MockSink::new();
1352        let records = vec![json!({"id": 1}), json!({"id": 2}), json!({"id": 3})];
1353        let count = sink.write_batch(&records).await.unwrap();
1354        assert_eq!(count, 3);
1355    }
1356
1357    #[tokio::test]
1358    async fn sink_write_batch_empty() {
1359        let sink = MockSink::new();
1360        let count = sink.write_batch(&[]).await.unwrap();
1361        assert_eq!(count, 0);
1362    }
1363
1364    #[tokio::test]
1365    async fn sink_accumulates_records() {
1366        let sink = MockSink::new();
1367        sink.write_batch(&[json!({"a": 1})]).await.unwrap();
1368        sink.write_batch(&[json!({"b": 2})]).await.unwrap();
1369        let written = sink.written.lock().unwrap();
1370        assert_eq!(written.len(), 2);
1371    }
1372
1373    #[tokio::test]
1374    async fn sink_default_flush_is_noop() {
1375        let sink = MockSink::new();
1376        assert!(sink.flush().await.is_ok());
1377    }
1378
1379    #[tokio::test]
1380    async fn sink_error_propagates() {
1381        let sink = FailingSink;
1382        let result = sink.write_batch(&[json!({"id": 1})]).await;
1383        assert!(result.is_err());
1384        assert!(matches!(result, Err(FaucetError::Sink(_))));
1385    }
1386
1387    #[tokio::test]
1388    async fn sink_as_trait_object() {
1389        let sink: Box<dyn Sink> = Box::new(MockSink::new());
1390        let count = sink.write_batch(&[json!({"id": 1})]).await.unwrap();
1391        assert_eq!(count, 1);
1392    }
1393
1394    // ── stream_pages tests ──────────────────────────────────────────────────
1395
1396    use crate::pipeline::DEFAULT_BATCH_SIZE;
1397    use futures::StreamExt;
1398
1399    #[tokio::test]
1400    async fn default_stream_pages_chunks_records() {
1401        let source = MockSource {
1402            records: (0..5).map(|i| json!({"i": i})).collect(),
1403        };
1404        let ctx = std::collections::HashMap::new();
1405        let mut pages = source.stream_pages(&ctx, 2);
1406        let mut all = Vec::new();
1407        while let Some(page) = pages.next().await {
1408            all.push(page.unwrap());
1409        }
1410        // 5 records, batch_size=2 → pages of [2, 2, 1]
1411        assert_eq!(all.len(), 3);
1412        assert_eq!(all[0].records.len(), 2);
1413        assert_eq!(all[1].records.len(), 2);
1414        assert_eq!(all[2].records.len(), 1);
1415    }
1416
1417    #[tokio::test]
1418    async fn default_stream_pages_attaches_bookmark_to_final_page_only() {
1419        let source = IncrementalSource {
1420            records: (0..5).map(|i| json!({"i": i})).collect(),
1421            bookmark: json!("v1"),
1422        };
1423        let ctx = std::collections::HashMap::new();
1424        let mut pages = source.stream_pages(&ctx, 2);
1425        let mut collected = Vec::new();
1426        while let Some(page) = pages.next().await {
1427            collected.push(page.unwrap());
1428        }
1429        assert_eq!(collected.len(), 3);
1430        assert!(collected[0].bookmark.is_none());
1431        assert!(collected[1].bookmark.is_none());
1432        assert_eq!(collected[2].bookmark, Some(json!("v1")));
1433    }
1434
1435    #[tokio::test]
1436    async fn default_stream_pages_single_page_when_batch_size_exceeds_total() {
1437        let source = MockSource {
1438            records: vec![json!({"id": 1}), json!({"id": 2})],
1439        };
1440        let ctx = std::collections::HashMap::new();
1441        let mut pages = source.stream_pages(&ctx, 100);
1442        let mut collected = Vec::new();
1443        while let Some(page) = pages.next().await {
1444            collected.push(page.unwrap());
1445        }
1446        assert_eq!(collected.len(), 1);
1447        assert_eq!(collected[0].records.len(), 2);
1448    }
1449
1450    #[tokio::test]
1451    async fn default_stream_pages_batch_size_zero_emits_single_page() {
1452        // batch_size = 0 is the "no batching" sentinel — yields every record
1453        // in one page regardless of total count.
1454        let source = MockSource {
1455            records: (0..50_000).map(|i| json!({"i": i})).collect(),
1456        };
1457        let ctx = std::collections::HashMap::new();
1458        let mut pages = source.stream_pages(&ctx, 0);
1459        let mut collected = Vec::new();
1460        while let Some(page) = pages.next().await {
1461            collected.push(page.unwrap());
1462        }
1463        assert_eq!(
1464            collected.len(),
1465            1,
1466            "batch_size=0 must emit exactly one page"
1467        );
1468        assert_eq!(collected[0].records.len(), 50_000);
1469    }
1470
1471    #[tokio::test]
1472    async fn default_stream_pages_batch_size_zero_attaches_bookmark_to_sole_page() {
1473        let source = IncrementalSource {
1474            records: (0..3).map(|i| json!({"i": i})).collect(),
1475            bookmark: json!("v1"),
1476        };
1477        let ctx = std::collections::HashMap::new();
1478        let mut pages = source.stream_pages(&ctx, 0);
1479        let page = pages.next().await.unwrap().unwrap();
1480        assert_eq!(page.records.len(), 3);
1481        assert_eq!(page.bookmark, Some(json!("v1")));
1482        assert!(pages.next().await.is_none());
1483    }
1484
1485    #[tokio::test]
1486    async fn default_stream_pages_empty_source_yields_no_pages() {
1487        let source = MockSource { records: vec![] };
1488        let ctx = std::collections::HashMap::new();
1489        let mut pages = source.stream_pages(&ctx, DEFAULT_BATCH_SIZE);
1490        assert!(pages.next().await.is_none());
1491    }
1492
1493    #[tokio::test]
1494    async fn default_stream_pages_empty_source_with_bookmark_yields_single_empty_page() {
1495        let source = IncrementalSource {
1496            records: vec![],
1497            bookmark: json!("v0"),
1498        };
1499        let ctx = std::collections::HashMap::new();
1500        let mut pages = source.stream_pages(&ctx, DEFAULT_BATCH_SIZE);
1501        let mut collected = Vec::new();
1502        while let Some(page) = pages.next().await {
1503            collected.push(page.unwrap());
1504        }
1505        // One empty-records page that carries the bookmark, so the pipeline
1506        // still persists progress on otherwise-empty incremental runs.
1507        assert_eq!(collected.len(), 1);
1508        assert!(collected[0].records.is_empty());
1509        assert_eq!(collected[0].bookmark, Some(json!("v0")));
1510    }
1511
1512    #[tokio::test]
1513    async fn default_stream_pages_propagates_fetch_errors() {
1514        let source = FailingSource;
1515        let ctx = std::collections::HashMap::new();
1516        let mut pages = source.stream_pages(&ctx, DEFAULT_BATCH_SIZE);
1517        let first = pages.next().await.unwrap();
1518        assert!(matches!(first, Err(FaucetError::Auth(_))));
1519    }
1520
1521    #[test]
1522    fn source_default_connector_name_is_stripped_type_name() {
1523        // MockSource lives at `faucet_core::traits::tests::MockSource`; the
1524        // stripped type_name yields the trailing segment.
1525        let source = MockSource { records: vec![] };
1526        assert_eq!(source.connector_name(), "MockSource");
1527    }
1528
1529    #[test]
1530    fn sink_default_connector_name_is_stripped_type_name() {
1531        let sink = MockSink::new();
1532        assert_eq!(sink.connector_name(), "MockSink");
1533    }
1534
1535    #[test]
1536    fn source_default_dataset_uri_uses_connector_name() {
1537        let source = MockSource { records: vec![] };
1538        assert_eq!(source.dataset_uri(), "MockSource://unknown");
1539    }
1540
1541    #[test]
1542    fn sink_default_dataset_uri_uses_connector_name() {
1543        let sink = MockSink::new();
1544        assert_eq!(sink.dataset_uri(), "MockSink://unknown");
1545    }
1546
1547    // ── write_batch_partial tests ───────────────────────────────────────────
1548
1549    #[tokio::test]
1550    async fn default_write_batch_partial_success_returns_all_ok() {
1551        let sink = MockSink::new();
1552        let records = vec![json!({"id": 1}), json!({"id": 2}), json!({"id": 3})];
1553        let outcomes = sink.write_batch_partial(&records).await.unwrap();
1554        assert_eq!(outcomes.len(), 3);
1555        assert!(outcomes.iter().all(|o| o.is_ok()));
1556        assert_eq!(sink.written.lock().unwrap().len(), 3);
1557    }
1558
1559    #[tokio::test]
1560    async fn default_write_batch_partial_bubbles_outer_err() {
1561        let sink = FailingSink;
1562        let records = vec![json!({"id": 1}), json!({"id": 2})];
1563        let result = sink.write_batch_partial(&records).await;
1564        assert!(matches!(result, Err(FaucetError::Sink(_))));
1565    }
1566
1567    #[tokio::test]
1568    async fn default_write_batch_partial_empty_returns_empty_vec() {
1569        let sink = MockSink::new();
1570        let outcomes = sink.write_batch_partial(&[]).await.unwrap();
1571        assert!(outcomes.is_empty());
1572    }
1573
1574    #[tokio::test]
1575    async fn default_write_batch_partial_callable_through_trait_object() {
1576        let sink: Box<dyn Sink> = Box::new(MockSink::new());
1577        let records = vec![json!({"id": 1}), json!({"id": 2})];
1578        let outcomes = sink.write_batch_partial(&records).await.unwrap();
1579        assert_eq!(outcomes.len(), 2);
1580        assert!(outcomes.iter().all(|o| o.is_ok()));
1581    }
1582
1583    // ── check() tests ─────────────────────────────────────────────────────────
1584
1585    #[tokio::test]
1586    async fn source_default_check_pulls_first_page_and_passes() {
1587        let source = MockSource {
1588            records: vec![json!({"id": 1}), json!({"id": 2})],
1589        };
1590        let report = source
1591            .check(&crate::check::CheckContext::default())
1592            .await
1593            .unwrap();
1594        assert_eq!(report.failed_count(), 0);
1595        assert!(
1596            report
1597                .probes
1598                .iter()
1599                .any(|p| p.name == "read" && matches!(p.status, crate::check::ProbeStatus::Pass))
1600        );
1601    }
1602
1603    #[tokio::test]
1604    async fn source_default_check_passes_on_empty_source() {
1605        let source = MockSource { records: vec![] };
1606        let report = source
1607            .check(&crate::check::CheckContext::default())
1608            .await
1609            .unwrap();
1610        // Reachable but empty is still a healthy source.
1611        assert_eq!(report.failed_count(), 0);
1612    }
1613
1614    #[tokio::test]
1615    async fn source_default_check_fails_when_fetch_errors() {
1616        let source = FailingSource;
1617        let report = source
1618            .check(&crate::check::CheckContext::default())
1619            .await
1620            .unwrap();
1621        assert_eq!(report.failed_count(), 1);
1622        assert!(report.probes.iter().any(
1623            |p| p.name == "read" && matches!(p.status, crate::check::ProbeStatus::Fail { .. })
1624        ));
1625    }
1626
1627    #[tokio::test]
1628    async fn sink_default_check_is_not_implemented_skip() {
1629        let sink = MockSink::new();
1630        let report = sink
1631            .check(&crate::check::CheckContext::default())
1632            .await
1633            .unwrap();
1634        assert_eq!(report.probes.len(), 1);
1635        assert!(matches!(
1636            report.probes[0].status,
1637            crate::check::ProbeStatus::Skip { .. }
1638        ));
1639    }
1640
1641    #[tokio::test]
1642    async fn source_check_callable_through_trait_object() {
1643        let source: Box<dyn Source> = Box::new(MockSource {
1644            records: vec![json!({"id": 1})],
1645        });
1646        let report = source
1647            .check(&crate::check::CheckContext::default())
1648            .await
1649            .unwrap();
1650        assert_eq!(report.failed_count(), 0);
1651    }
1652
1653    // ── idempotent-write / exactly-once capability tests ──────────────────────
1654
1655    #[tokio::test]
1656    async fn sink_default_is_not_idempotent() {
1657        let sink = MockSink::new();
1658        assert!(!sink.supports_idempotent_writes());
1659        // Default write_batch_idempotent ignores the token and delegates.
1660        let n = sink
1661            .write_batch_idempotent(&[json!({"id": 1})], "scope::a", "00000000000000000001")
1662            .await
1663            .unwrap();
1664        assert_eq!(n, 1);
1665        assert_eq!(sink.last_committed_token("scope::a").await.unwrap(), None);
1666        assert_eq!(sink.written.lock().unwrap().len(), 1);
1667    }
1668
1669    #[test]
1670    fn source_default_does_not_support_exactly_once() {
1671        let source = MockSource { records: vec![] };
1672        assert!(!source.supports_exactly_once());
1673    }
1674
1675    #[test]
1676    fn sink_default_supported_write_modes_is_append_only() {
1677        use crate::write_mode::WriteMode;
1678        let sink = MockSink::new();
1679        assert_eq!(sink.supported_write_modes(), &[WriteMode::Append]);
1680    }
1681
1682    #[test]
1683    fn supported_write_modes_callable_through_trait_object() {
1684        use crate::write_mode::WriteMode;
1685        let sink: Box<dyn Sink> = Box::new(MockSink::new());
1686        assert!(sink.supported_write_modes().contains(&WriteMode::Append));
1687    }
1688
1689    #[tokio::test]
1690    async fn sink_default_current_schema_is_none() {
1691        let sink = MockSink::new();
1692        assert_eq!(sink.current_schema().await.unwrap(), None);
1693    }
1694
1695    #[test]
1696    fn sink_default_does_not_support_schema_evolution() {
1697        let sink = MockSink::new();
1698        assert!(!sink.supports_schema_evolution());
1699    }
1700
1701    #[tokio::test]
1702    async fn sink_default_evolve_schema_is_unsupported_error() {
1703        let sink = MockSink::new();
1704        let evo = crate::drift::SchemaEvolution::default();
1705        let err = sink.evolve_schema(&evo).await.unwrap_err();
1706        assert!(matches!(err, FaucetError::Sink(_)));
1707        assert!(err.to_string().contains("schema evolution"));
1708    }
1709
1710    #[tokio::test]
1711    async fn source_default_capture_resume_position_is_none() {
1712        let source = MockSource { records: vec![] };
1713        assert_eq!(source.capture_resume_position().await.unwrap(), None);
1714    }
1715
1716    #[tokio::test]
1717    async fn capture_resume_position_callable_through_trait_object() {
1718        let source: Box<dyn Source> = Box::new(MockSource { records: vec![] });
1719        assert!(source.capture_resume_position().await.unwrap().is_none());
1720    }
1721
1722    #[tokio::test]
1723    async fn source_default_does_not_support_discover() {
1724        let source: Box<dyn Source> = Box::new(MockSource { records: vec![] });
1725        assert!(!source.supports_discover());
1726        let err = source.discover().await.unwrap_err();
1727        assert!(matches!(err, FaucetError::Source(_)));
1728        assert!(
1729            err.to_string().contains("dataset discovery"),
1730            "typed unsupported error: {err}"
1731        );
1732    }
1733
1734    #[tokio::test]
1735    async fn source_default_is_not_shardable() {
1736        let source: Box<dyn Source> = Box::new(MockSource { records: vec![] });
1737        assert!(!source.is_shardable());
1738    }
1739
1740    #[tokio::test]
1741    async fn source_default_enumerates_single_whole_shard() {
1742        // A non-shardable source enumerates to exactly one whole-dataset shard,
1743        // regardless of the requested target — preserving single-worker behavior.
1744        let source: Box<dyn Source> = Box::new(MockSource { records: vec![] });
1745        let shards = source.enumerate_shards(8).await.unwrap();
1746        assert_eq!(shards.len(), 1);
1747        assert!(shards[0].is_whole());
1748    }
1749
1750    #[tokio::test]
1751    async fn source_default_apply_shard_is_noop() {
1752        let source: Box<dyn Source> = Box::new(MockSource { records: vec![] });
1753        // Applying the whole shard is a no-op and must not error.
1754        source
1755            .apply_shard(&crate::shard::ShardSpec::whole())
1756            .await
1757            .unwrap();
1758    }
1759}