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 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 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 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 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 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