Skip to main content

BufferedSink

Struct BufferedSink 

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

A sink that buffers writes and only makes them durable on flush — modelling a real buffered sink whose output is committed at flush time (a Parquet footer, an S3 multipart completion). Used to exercise assert_cancellation_flushes: the pipeline must flush at the cancellation page boundary or the staged rows are lost.

Construct with BufferedSink::new for a faithful sink (flush commits the staging buffer) or BufferedSink::broken for one whose flush silently drops the buffer (used to prove the check fails when a cancel does not yield durable output).

Implementations§

Source§

impl BufferedSink

Source

pub fn new() -> Self

A buffered sink whose flush durably commits everything staged so far.

Source

pub fn broken() -> Self

A broken buffered sink whose flush is a silent no-op — staged rows never become durable, so any output buffered when the run ends is lost.

Source

pub fn durable_len(&self) -> usize

Number of rows that are durable (committed via a flush).

Source

pub fn staged_len(&self) -> usize

Number of rows currently staged but not yet flushed.

Trait Implementations§

Source§

impl Clone for BufferedSink

Source§

fn clone(&self) -> BufferedSink

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Default for BufferedSink

Source§

fn default() -> Self

Returns the “default value” for a type. Read more
Source§

impl Sink for BufferedSink

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 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 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 write_batch_partial<'life0, 'life1, 'async_trait>( &'life0 self, records: &'life1 [Value], ) -> Pin<Box<dyn Future<Output = Result<Vec<Result<(), FaucetError>>, FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Write a batch and report per-row outcomes. 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 sink_guarantee(&self) -> SinkGuarantee

The strongest delivery guarantee this sink can uphold — see SinkGuarantee. Read more
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 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 current_schema<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Option<Value>, FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: '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 'life0: 'async_trait, 'life1: 'async_trait, Self: '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 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 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, Self: 'async_trait,

Delete rows matching scope whose key is not in seen. 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 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait, Self: 'async_trait,

Write records AND durably record token for scope, atomically. 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 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

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

fn config_schema(&self) -> Value

Return a JSON Schema describing the configuration this sink accepts. Read more
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 check<'life0, 'life1, 'async_trait>( &'life0 self, _ctx: &'life1 CheckContext, ) -> Pin<Box<dyn Future<Output = Result<CheckReport, FaucetError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

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

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> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> DynClone for T
where T: Clone,

Source§

fn __clone_box(&self, _: Private) -> *mut ()

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

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> 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> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

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

Source§

type Error = Infallible

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