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}