pub struct TestSink { /* private fields */ }Expand description
A sink that records everything written, optionally deduplicating by a key field (upsert), optionally advertising the atomic-watermark idempotent path.
Modes:
TestSink::new— append-only, non-idempotent.TestSink::keyed— dedups by key onwrite_batch(keyed-upsert /dedups_by_key), advertisesUpsert/Delete.TestSink::idempotent— additionally advertisessupports_idempotent_writesand stores a per-scope commit token, so the atomic-watermark path can be exercised.
Implementations§
Source§impl TestSink
impl TestSink
Sourcepub fn keyed(key_field: impl Into<String>) -> Self
pub fn keyed(key_field: impl Into<String>) -> Self
An upsert sink that dedups by key_field in write_batch.
Sourcepub fn idempotent(key_field: impl Into<String>) -> Self
pub fn idempotent(key_field: impl Into<String>) -> Self
An upsert sink that also commits an atomic watermark token per scope,
so it advertises (and honours) supports_idempotent_writes.
Sourcepub fn total_written(&self) -> usize
pub fn total_written(&self) -> usize
Total number of records passed to write_batch across all calls
(counts re-delivered duplicates).
Trait Implementations§
Source§impl Sink for TestSink
impl Sink for TestSink
Source§fn write_batch<'life0, 'life1, 'async_trait>(
&'life0 self,
records: &'life1 [Value],
) -> Pin<Box<dyn Future<Output = Result<usize, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn write_batch<'life0, 'life1, 'async_trait>(
&'life0 self,
records: &'life1 [Value],
) -> Pin<Box<dyn Future<Output = Result<usize, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Write a batch of records to the destination. Read more
Source§fn supports_idempotent_writes(&self) -> bool
fn supports_idempotent_writes(&self) -> bool
Whether this sink can durably commit a page’s rows and a commit token
in a single atomic transaction. Default:
false (at-least-once only). Read moreSource§fn dedups_by_key(&self) -> bool
fn dedups_by_key(&self) -> bool
Whether this sink instance is configured to dedup by key — i.e.
write_mode: upsert (or delete) with a non-empty key, so
re-applying a record with the same key converges instead of
duplicating. Default: false. Read moreSource§fn supported_write_modes(&self) -> &'static [WriteMode]
fn supported_write_modes(&self) -> &'static [WriteMode]
Write modes this sink can apply. Default: append-only. Sinks that
implement key-based merge override this to include
WriteMode::Upsert /
WriteMode::Delete. The CLI rejects a
configured mode that is not in this set at config-load time.Source§fn write_batch_idempotent<'life0, 'life1, 'life2, 'life3, 'async_trait>(
&'life0 self,
records: &'life1 [Value],
scope: &'life2 str,
token: &'life3 str,
) -> Pin<Box<dyn Future<Output = Result<usize, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
'life3: 'async_trait,
fn write_batch_idempotent<'life0, 'life1, 'life2, 'life3, 'async_trait>(
&'life0 self,
records: &'life1 [Value],
scope: &'life2 str,
token: &'life3 str,
) -> Pin<Box<dyn Future<Output = Result<usize, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
'life3: 'async_trait,
Source§fn last_committed_token<'life0, 'life1, 'async_trait>(
&'life0 self,
scope: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<String>, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn last_committed_token<'life0, 'life1, 'async_trait>(
&'life0 self,
scope: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<String>, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
The last token durably committed for
scope, or None if this sink has
never committed under that scope. Default: None.Source§fn connector_name(&self) -> &'static str
fn connector_name(&self) -> &'static str
Stable identifier used as the
connector label on metrics and the
connector attribute on spans. See Source::connector_name.Source§fn flush<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
Self: 'async_trait,
fn flush<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
Self: 'async_trait,
Flush any buffered data to the destination. Read more
Source§fn write_batch_partial<'life0, 'life1, 'async_trait>(
&'life0 self,
records: &'life1 [Value],
) -> Pin<Box<dyn Future<Output = Result<Vec<Result<(), FaucetError>>, FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
Self: 'async_trait,
fn write_batch_partial<'life0, 'life1, 'async_trait>(
&'life0 self,
records: &'life1 [Value],
) -> Pin<Box<dyn Future<Output = Result<Vec<Result<(), FaucetError>>, FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
Self: 'async_trait,
Write a batch and report per-row outcomes. Read more
Source§fn sink_guarantee(&self) -> SinkGuarantee
fn sink_guarantee(&self) -> SinkGuarantee
The strongest delivery guarantee this sink can uphold — see
SinkGuarantee. Read moreSource§fn current_schema<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Option<Value>, FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
Self: 'async_trait,
fn current_schema<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Option<Value>, FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
Self: 'async_trait,
The sink’s live destination schema as an
infer_schema-shaped object
({"type":"object","properties":{ <col>: <type-fragment>, … }}), or
None for a schemaless sink or a target that does not exist yet. Read moreSource§fn supports_schema_evolution(&self) -> bool
fn supports_schema_evolution(&self) -> bool
Whether this sink can apply additive/widening DDL via
evolve_schema. Default: false. The CLI rejects
on_drift: evolve against a sink that returns false at config-load.Source§fn evolve_schema<'life0, 'life1, 'async_trait>(
&'life0 self,
evolution: &'life1 SchemaEvolution,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
Self: 'async_trait,
fn evolve_schema<'life0, 'life1, 'async_trait>(
&'life0 self,
evolution: &'life1 SchemaEvolution,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
Self: 'async_trait,
Apply an additive schema evolution (new columns, lossless widenings,
nullability relaxations) to the destination. MUST be idempotent
(
ADD COLUMN IF NOT EXISTS semantics) so concurrent runs converge. Read moreSource§fn config_schema(&self) -> Value
fn config_schema(&self) -> Value
Return a JSON Schema describing the configuration this sink accepts. Read more
Source§fn dataset_uri(&self) -> String
fn dataset_uri(&self) -> String
Logical dataset identity for lineage emission, following OpenLineage
naming conventions (https://openlineage.io/docs/spec/naming). Read more
Source§fn check<'life0, 'life1, 'async_trait>(
&'life0 self,
_ctx: &'life1 CheckContext,
) -> Pin<Box<dyn Future<Output = Result<CheckReport, FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
Self: 'async_trait,
fn check<'life0, 'life1, 'async_trait>(
&'life0 self,
_ctx: &'life1 CheckContext,
) -> Pin<Box<dyn Future<Output = Result<CheckReport, FaucetError>> + Send + 'async_trait>>where
'life0: 'async_trait,
'life1: 'async_trait,
Self: 'async_trait,
Run a fast, non-mutating preflight probe (used by
faucet doctor). Read moreAuto Trait Implementations§
impl Freeze for TestSink
impl RefUnwindSafe for TestSink
impl Send for TestSink
impl Sync for TestSink
impl Unpin for TestSink
impl UnsafeUnpin for TestSink
impl UnwindSafe for TestSink
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more