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 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 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 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 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 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, <T as TryFrom<U>>::Error>

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