Skip to main content

MetadataSink

Struct MetadataSink 

Source
pub struct MetadataSink { /* private fields */ }
Expand description

A Sink decorator that stamps _faucet_* metadata columns onto every row before delegating the write.

Implementations§

Source§

impl MetadataSink

Source

pub fn new( inner: Box<dyn Sink>, meta: CompiledMetadata, ctx: MetadataContext, ) -> Self

Wrap a sink so every written row carries the configured metadata columns.

Trait Implementations§

Source§

impl Debug for MetadataSink

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Sink for MetadataSink

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

Auto Trait Implementations§

Blanket Implementations§

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