Skip to main content

CleanupTracker

Struct CleanupTracker 

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

Source

pub fn new(inner: &'a S, policy: &CleanupPolicy) -> Self

Source

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.

Source

pub fn tracked(&self) -> usize

Number of keys tracked so far (for logging).

Trait Implementations§

Source§

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,

Write a batch of records 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<RowOutcome>, FaucetError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Write a batch and report per-row outcomes. Read more
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,

Write records AND durably record token for scope, atomically. Read more
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,

Flush any buffered data to the destination. Read more
Source§

fn supports_cleanup(&self) -> bool

Whether this sink can delete a scoped set of rows for scoped cleanup (#478). Default 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,

Delete rows matching scope whose key is not in seen. Read more
Source§

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

The last token durably committed for scope, or None if this sink has never committed under that scope. Default: None.
Source§

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 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 more
Source§

fn batch_atomicity(&self) -> BatchAtomicity

What a failed batch write leaves behind (#737): nothing (Atomic), per-row outcomes with nothing committed on an outer Err (PerRow), or possibly some rows (BestEffort, the default). Read more
Source§

fn sink_guarantee(&self) -> SinkGuarantee

The strongest delivery guarantee this sink can uphold — see SinkGuarantee. Read more
Source§

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,

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 more
Source§

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 Self: 'async_trait, 'life0: 'async_trait, 'life1: '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 more
Source§

fn config_schema(&self) -> Value

Return a JSON Schema describing the configuration this sink accepts. Read more
Source§

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 dataset_uri(&self) -> String

Logical dataset identity for lineage emission, following OpenLineage naming conventions (https://openlineage.io/docs/spec/naming). Read more
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,

The concrete local files this sink instance opened during the run. Read more
Source§

fn is_overwrite(&self) -> bool

Whether this sink instance is configured for full-destination replacement (WriteMode::Overwrite). Read more
Source§

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,

Prepare a staging target for an overwrite run. Read more
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,

Atomically replace the destination with the staged data. Read more
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,

Discard the staging target after a failed or cancelled overwrite run. Read more
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,

Finalize this run’s output after Pipeline::run finished successfully and uncancelled, on every transfer path (Value, columnar, native), after the terminal flush. Read more
Source§

fn supports_columnar(&self) -> bool

Available on crate feature arrow only.
Whether this sink can consume columnar (arrow::RecordBatch) writes via write_batch_columnar without first converting to Value. Default: false. Read more
Source§

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,

Available on crate feature arrow only.
Write a columnar RecordBatch to the destination, returning the number of rows written. Read more
Source§

fn native_load_capabilities(&self) -> Vec<NativeLoadCapability>

Native byte-passthrough load mechanisms this sink offers, each naming the wire format and the write modes it can honor (#633). Default: vec![] (no fast path; use write_batch). Read more
Source§

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,

Bulk-load one native-format byte batch directly (#633), returning the number of rows written. Read more
Source§

fn write_batch_is_replay_safe(&self) -> bool

Whether a plain 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 more
Source§

fn supports_staged_load(&self) -> bool

Whether this sink can bulk-load via an object-store stage — write the page to S3/GCS/Azure, then pull it with the warehouse’s native load command (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

Whether this sink can undo a run it wrote (#706) through 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,

Undo everything 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 more
Source§

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,

Drop whatever this sink kept to make 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,

Rewind the exactly-once watermark for 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)>

A source (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,

Whether an overwrite staging object (the …__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 more
Source§

fn set_roundtrip_recorder(&self, _recorder: Arc<RoundtripRecorder>)

Receive a pre-labelled handle for counting upstream round trips (#638). See 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,

Run a fast, non-mutating preflight probe (used by faucet doctor). Read more

Auto 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>
where &'a S: Send, S: ?Sized,

§

impl<'a, S> Sync for CleanupTracker<'a, S>
where &'a S: Sync, S: ?Sized,

§

impl<'a, S> Unpin for CleanupTracker<'a, S>
where &'a S: Unpin, S: ?Sized,

§

impl<'a, S> UnsafeUnpin for CleanupTracker<'a, S>
where &'a S: UnsafeUnpin, S: ?Sized,

§

impl<'a, S> UnwindSafe for CleanupTracker<'a, S>
where &'a S: UnwindSafe, S: ?Sized,

Blanket Implementations§

Source§

impl<T> Allocation for T
where T: RefUnwindSafe + Send + Sync,

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FutureExt for T

Source§

fn with_context(self, otel_cx: Context) -> WithContext<Self> ⓘ

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
Source§

fn with_current_context(self) -> WithContext<Self> ⓘ

Attaches the current Context to this type, returning a WithContext wrapper. Read more
Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts 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 more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts 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 more
Source§

impl<T> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more