Skip to main content

S3Source

Struct S3Source 

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

Coordinated object-storage backfill source. See the crate docs for the delivery model and the split identity/drift contract.

Implementations§

Source§

impl S3Source

Source

pub fn new(config: S3SourceConfig, io: Handle) -> S3Source

A source over config, doing its network I/O on io — pass the pipeline’s I/O runtime handle (Pipeline::io_handle).

The handle must belong to a multi-thread runtime that outlives the pipeline (the pipeline’s I/O runtime satisfies both): the source, its lanes, and the coordinator briefly block pipeline threads on it, which a current-thread runtime cannot drive.

Source

pub fn from_component_config( section: &ComponentConfig, io: Handle, ) -> Result<S3Source, ConfigError>

Build from the pipeline’s opaque source: { s3: ... } section.

Source

pub fn with_framer<F>(self, factory: F) -> S3Source
where F: Fn() -> Box<dyn RecordFramer> + Send + Sync + 'static,

Set the record framer that cuts each object’s byte stream into records. spate-s3 is a transport and owns no framing, so this is required: supply the framer for the objects’ format — e.g. spate-json’s NdjsonFramer for NDJSON — before the pipeline opens the source.

factory builds a fresh RecordFramer per object (framers are per-object stateful and each lane frames its own split). A framed source always emits one record per payload, so its FramingContract is PerRecord and the paired deserializer decodes a single unit.

Source

pub fn with_coordinator( self, coordinator: Box<dyn SplitCoordinator>, ) -> S3Source

Hand the source its coordinator — the multi-instance seam. Build any SplitCoordinator at assembly time (e.g. StoreCoordinator over the NATS store from spate-coordination, with your own tuning) and inject it here; run more replicas of the same pipeline against the same backend and they share the backfill. Must be called before the pipeline opens the source.

Without it the source runs solo over an in-process store: correct and self-terminating, but progress is ephemeral — a restart replays the whole prefix (a startup WARN says so).

Trait Implementations§

Source§

impl Debug for S3Source

Source§

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

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

impl Drop for S3Source

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more
Source§

impl Source for S3Source

Source§

type Lane = S3Lane

The lane type this source produces.
Source§

fn component_type(&self) -> &str

The component_type metric label for this source (e.g. "kafka"), mirroring SinkParts::with_component_type on the sink side. It is also the namespace of the source’s custom-metrics Meter: declaring "kafka" scopes the source’s own families under spate_kafka_source_*. The default "source" is a reserved root, so a source that does not override this gets no custom Meter (its framework stage metrics are unaffected).
Source§

fn framing_contract(&self) -> FramingContract

How the payloads this source emits are framed, so the framework can pair it with a deserializer without the two being coordinated by hand (see FramingContract). A source that splits its own bytes into one record per payload returns FramingContract::PerRecord; the default is FramingContract::WholePayload — the source emits whole payloads and the deserializer owns framing (Kafka, and any source that does not frame).
Source§

fn open(&mut self, ctx: SourceCtx) -> Result<(), SourceError>

Connect and prepare. Called once before any other method.
Source§

fn poll_events( &mut self, timeout: Duration, ) -> Result<SourceEvent<S3Lane>, SourceError>

Service control-plane work (rebalance callbacks, statistics) and return the next event, waiting at most timeout. Must be called regularly regardless of backpressure state.
Source§

fn commit( &mut self, watermarks: &[(PartitionId, i64)], ) -> Result<(), SourceError>

Store per-partition committable positions (each is the offset one past the last acknowledged record). Positions are durable per the source’s own policy (e.g. interval auto-commit of stored offsets).
Source§

fn pause(&mut self, lanes: &[LaneId]) -> Result<(), SourceError>

Stop fetching for lanes (backpressure). Optional capability: sources that cannot pause rely on bounded-queue pushback alone.
Source§

fn resume(&mut self, lanes: &[LaneId]) -> Result<(), SourceError>

Resume fetching for lanes.
Source§

fn flush_commits(&mut self) -> Result<(), SourceError>

Synchronously flush stored positions (shutdown, revocation).

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

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