pub struct CleanupTracker<'a, S: Sink + ?Sized> { /* private fields */ }Expand description
A sink wrapper that records the key tuples written through it (#478).
This is how scoped cleanup tracks “what this run wrote” without adding a
field to RunStreamOptions, which is an
externally-constructible struct whose shape is part of the public API. It is
also the more honest home for the bookkeeping: the thing that writes the rows
is the thing that knows which rows were written, and it composes with the
existing sink-decorator pattern (InstrumentedSink) rather than threading a
second concern through the page loop.
Every write path is counted, including write_batch_partial.
That is deliberate: a row handed to the sink that fails and lands in the DLQ
is still a record the source claimed present, so it must count as seen or the
cleanup would delete its destination row.
Records that never reach the sink at all — quarantined by a quality, contract, or drift policy — are consequently not counted, which is why those combinations are rejected at config-load time rather than silently deleting the quarantined rows’ destination counterparts.
Implementations§
Source§impl<'a, S: Sink + ?Sized> CleanupTracker<'a, S>
impl<'a, S: Sink + ?Sized> CleanupTracker<'a, S>
pub fn new(inner: &'a S, policy: &CleanupPolicy) -> Self
Sourcepub async fn finish(&self, policy: &CleanupPolicy) -> Result<u64, FaucetError>
pub async fn finish(&self, policy: &CleanupPolicy) -> Result<u64, FaucetError>
Run the scoped delete against the wrapped sink. Call only after a fully successful, uncancelled run — see the module docs.
Trait Implementations§
Source§impl<S: Sink + ?Sized> Sink for CleanupTracker<'_, S>
impl<S: Sink + ?Sized> Sink for CleanupTracker<'_, S>
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,
Source§fn write_batch_partial<'life0, 'life1, 'async_trait>(
&'life0 self,
records: &'life1 [Value],
) -> Pin<Box<dyn Future<Output = Result<Vec<RowOutcome>, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn write_batch_partial<'life0, 'life1, 'async_trait>(
&'life0 self,
records: &'life1 [Value],
) -> Pin<Box<dyn Future<Output = Result<Vec<RowOutcome>, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
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 flush<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn flush<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Source§fn supports_cleanup(&self) -> bool
fn supports_cleanup(&self) -> bool
false; the upsert-capable sinks override it.Source§fn cleanup_scope<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
scope: &'life1 BTreeMap<String, Value>,
seen: &'life2 SeenKeys,
) -> Pin<Box<dyn Future<Output = Result<u64, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
fn cleanup_scope<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
scope: &'life1 BTreeMap<String, Value>,
seen: &'life2 SeenKeys,
) -> Pin<Box<dyn Future<Output = Result<u64, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
Source§fn supports_idempotent_writes(&self) -> bool
fn supports_idempotent_writes(&self) -> bool
false (at-least-once only). Read moreSource§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,
scope, or None if this sink has
never committed under that scope. Default: None.Source§fn supported_write_modes(&self) -> &'static [WriteMode]
fn supported_write_modes(&self) -> &'static [WriteMode]
WriteMode::Upsert /
WriteMode::Delete. The CLI rejects a
configured mode that is not in this set at config-load time.Source§fn dedups_by_key(&self) -> bool
fn dedups_by_key(&self) -> bool
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 batch_atomicity(&self) -> BatchAtomicity
fn batch_atomicity(&self) -> BatchAtomicity
Atomic), per-row outcomes with
nothing committed on an outer Err
(PerRow), or possibly some rows
(BestEffort, the default). Read moreSource§fn sink_guarantee(&self) -> SinkGuarantee
fn sink_guarantee(&self) -> SinkGuarantee
SinkGuarantee. Read moreSource§fn current_schema<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Option<Value>, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn current_schema<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Option<Value>, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
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
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
Self: 'async_trait,
'life0: 'async_trait,
'life1: '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
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
ADD COLUMN IF NOT EXISTS semantics) so concurrent runs converge. Read moreSource§fn config_schema(&self) -> Value
fn config_schema(&self) -> Value
Source§fn connector_name(&self) -> &'static str
fn connector_name(&self) -> &'static str
connector label on metrics and the
connector attribute on spans. See Source::connector_name.Source§fn dataset_uri(&self) -> String
fn dataset_uri(&self) -> String
Source§fn local_outputs<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Vec<LocalOutput>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn local_outputs<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Vec<LocalOutput>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Source§fn is_overwrite(&self) -> bool
fn is_overwrite(&self) -> bool
WriteMode::Overwrite). Read moreSource§fn begin_overwrite<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn begin_overwrite<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Source§fn commit_overwrite<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn commit_overwrite<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Source§fn abort_overwrite<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn abort_overwrite<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Source§fn complete_run<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn complete_run<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Pipeline::run
finished successfully and uncancelled, on every transfer path
(Value, columnar, native), after the terminal flush. Read moreSource§fn supports_columnar(&self) -> bool
fn supports_columnar(&self) -> bool
arrow only.arrow::RecordBatch) writes
via write_batch_columnar without first
converting to Value. Default: false. Read moreSource§fn write_batch_columnar<'life0, 'life1, 'async_trait>(
&'life0 self,
batch: &'life1 RecordBatch,
) -> Pin<Box<dyn Future<Output = Result<usize, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn write_batch_columnar<'life0, 'life1, 'async_trait>(
&'life0 self,
batch: &'life1 RecordBatch,
) -> Pin<Box<dyn Future<Output = Result<usize, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
arrow only.RecordBatch to the destination, returning the number of
rows written. Read moreSource§fn native_load_capabilities(&self) -> Vec<NativeLoadCapability>
fn native_load_capabilities(&self) -> Vec<NativeLoadCapability>
vec![] (no fast path; use write_batch). Read moreSource§fn load_native<'life0, 'life1, 'async_trait>(
&'life0 self,
batch: NativeBatch,
scope: &'life1 str,
ctx: NativeLoadContext,
) -> Pin<Box<dyn Future<Output = Result<usize, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn load_native<'life0, 'life1, 'async_trait>(
&'life0 self,
batch: NativeBatch,
scope: &'life1 str,
ctx: NativeLoadContext,
) -> Pin<Box<dyn Future<Output = Result<usize, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Source§fn write_batch_is_replay_safe(&self) -> bool
fn write_batch_is_replay_safe(&self) -> bool
write_batch is safe to replay
— i.e. re-sending the same page after an ambiguous failure converges
instead of duplicating rows. Default: false. Read moreSource§fn supports_staged_load(&self) -> bool
fn supports_staged_load(&self) -> bool
COPY / COPY INTO / s3() table function) (#528). Default
false; warehouse sinks that honour a StagingSpec
override it. Object-safe (no args, no generics) so Box<dyn Sink> is
unaffected.Source§fn supports_rollback(&self) -> bool
fn supports_rollback(&self) -> bool
rollback_run. Default false.Source§fn rollback_run<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
run_id: &'life1 str,
opts: &'life2 RollbackOptions,
) -> Pin<Box<dyn Future<Output = Result<RollbackOutcome, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
fn rollback_run<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
run_id: &'life1 str,
opts: &'life2 RollbackOptions,
) -> Pin<Box<dyn Future<Output = Result<RollbackOutcome, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
run_id wrote to this destination, per
opts.mode: delete the run’s rows (append), restore journaled
before-images (upsert / delete), or swap the kept previous table back
(overwrite). Must be all-or-nothing per destination where the backend
allows it, and must leave the destination untouched when opts.dry_run
is set or when conflicts block it (see
RollbackOutcome). Read moreSource§fn forget_run<'life0, 'life1, 'async_trait>(
&'life0 self,
run_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn forget_run<'life0, 'life1, 'async_trait>(
&'life0 self,
run_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
run_id undoable (journal rows, a
previous table) once the run is past the retention window. Default:
no-op.Source§fn rewind_commit_token<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
scope: &'life1 str,
token: Option<&'life2 str>,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
fn rewind_commit_token<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
scope: &'life1 str,
token: Option<&'life2 str>,
) -> Pin<Box<dyn Future<Output = Result<(), FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
scope to token (None clears
it), so a rolled-back run’s pages are not skipped on the next run.
Default: a typed “unsupported” error.Source§fn readback_source(&self) -> Option<(String, Value)>
fn readback_source(&self) -> Option<(String, Value)>
(kind, config) that reads this destination back — what
faucet verify (#701) compares the pipeline’s source against. Default
None: the user names the destination reader in verify.destination.
The config must not need a credential the sink config lacks.Source§fn overwrite_staging_exists<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Option<bool>, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn overwrite_staging_exists<'life0, 'async_trait>(
&'life0 self,
) -> Pin<Box<dyn Future<Output = Result<Option<bool>, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
…__faucet_ovw table or
collection a begin_overwrite creates) exists
right now — a read-only probe used by faucet status --probe to report
staging a crashed or aborted overwrite left behind. Read moreSource§fn set_roundtrip_recorder(&self, _recorder: Arc<RoundtripRecorder>)
fn set_roundtrip_recorder(&self, _recorder: Arc<RoundtripRecorder>)
Source::set_roundtrip_recorder; defaulted to a no-op.Source§fn check<'life0, 'life1, 'async_trait>(
&'life0 self,
_ctx: &'life1 CheckContext,
) -> Pin<Box<dyn Future<Output = Result<CheckReport, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn check<'life0, 'life1, 'async_trait>(
&'life0 self,
_ctx: &'life1 CheckContext,
) -> Pin<Box<dyn Future<Output = Result<CheckReport, FaucetError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
faucet doctor). Read moreAuto Trait Implementations§
impl<'a, S> !Freeze for CleanupTracker<'a, S>
impl<'a, S> RefUnwindSafe for CleanupTracker<'a, S>
impl<'a, S> Send for CleanupTracker<'a, S>
impl<'a, S> Sync for CleanupTracker<'a, S>
impl<'a, S> Unpin for CleanupTracker<'a, S>
impl<'a, S> UnsafeUnpin for CleanupTracker<'a, S>
impl<'a, S> UnwindSafe for CleanupTracker<'a, S>
Blanket Implementations§
impl<T> Allocation for T
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
Source§impl<T> FutureExt for T
impl<T> FutureExt for T
Source§fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ
Source§fn with_current_context(self) -> WithContext<Self> ⓘ
fn with_current_context(self) -> WithContext<Self> ⓘ
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left is true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
self into a Left variant of Either<Self, Self>
if into_left(&self) returns true.
Converts self into a Right variant of Either<Self, Self>
otherwise. Read moreSource§impl<T> IntoRequest<T> for T
impl<T> IntoRequest<T> for T
Source§fn into_request(self) -> Request<T>
fn into_request(self) -> Request<T>
T in a tonic::Request