Skip to main content

Tier

Struct Tier 

Source
pub struct Tier<R, St, S> { /* private fields */ }
Expand description

Buffering tier over a pluggable Store.

Records ingest into the wall-clock-aligned window align(meta.end) (aligned to the policy’s every); when the policy fires, stored windows drain downstream oldest-first, each removed from the store only after the downstream sink accepted it — a transient outage of the terminal sink does not lose data held in this tier.

One window key can drain downstream more than once: the max_records valve fires mid-window and later records re-open the same key, and a failed drain is retried after the window has grown. Both deliveries carry the same WindowMeta, so a terminal sink that keys storage by window start alone would overwrite the earlier chunk — key by content as well, as HfSink (feature huggingface) does.

A window whose delivery keeps failing (for example a record the terminal sink deterministically rejects) does not stall the pipeline: each drain pass attempts every closed window oldest-first, retains the failed ones for the next pass, and surfaces the first error. Such a poisoned window is retried on every pass and retained in the store indefinitely — remove its segment by hand (for JsonlStore, delete the window’s .jsonl file) if it must be discarded.

The policy is checked on each ingest (collector ticks are frequent compared to flush windows, so no extra timer task is needed); flush force-drains. Whatever a previous run left in the store is replayed downstream on the first ingest or flush call, with WindowMeta reconstructed from the window key alone — replayed windows land at the same storage path (idempotent).

Implementations§

Source§

impl<R, St, S> Tier<R, St, S>

Source

pub fn new(store: St, policy: FlushPolicy, inner: S) -> Self

Create a tier holding records in store until policy fires.

Source

pub fn with_pipeline_name(self, name: impl Into<String>) -> Self

Override the pipeline name used in drained WindowMeta (default: the first live meta seen, then the store’s pipeline_hint).

Source

pub fn inner(&self) -> &S

Access the wrapped sink.

Trait Implementations§

Source§

impl<R, St, S> Sink<R> for Tier<R, St, S>
where R: Send + 'static, St: Store<R>, S: Sink<R>,

Source§

type Error = TierError<<St as Store<R>>::Error, <S as Sink<R>>::Error>

Concrete error type (a thiserror enum, not a boxed error).
Source§

async fn ingest( &mut self, meta: &WindowMeta, records: Vec<R>, ) -> Result<(), Self::Error>

Hand records to this layer. A buffering layer may hold them; a terminal sink ships them immediately.
Source§

async fn flush(&mut self) -> Result<(), Self::Error>

Force-drain this layer and everything downstream (shutdown, final flush, startup recovery).

Auto Trait Implementations§

§

impl<R, St, S> Freeze for Tier<R, St, S>
where St: Freeze, S: Freeze,

§

impl<R, St, S> RefUnwindSafe for Tier<R, St, S>

§

impl<R, St, S> Send for Tier<R, St, S>
where St: Send, S: Send,

§

impl<R, St, S> Sync for Tier<R, St, S>
where St: Sync, S: Sync,

§

impl<R, St, S> Unpin for Tier<R, St, S>
where St: Unpin, S: Unpin,

§

impl<R, St, S> UnsafeUnpin for Tier<R, St, S>
where St: UnsafeUnpin, S: UnsafeUnpin,

§

impl<R, St, S> UnwindSafe for Tier<R, St, S>
where St: UnwindSafe, S: UnwindSafe,

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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

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> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<R, S> SinkExt<R> for S
where S: Sink<R>,

Source§

fn tee<B: Sink<R>>(self, other: B) -> Tee<Self, B>

Fan out: every batch is ingested into both self and other.
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