Skip to main content

TargetSink

Struct TargetSink 

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

A Sink that applies residency filtering, topic renaming, and loop-guard logic before producing records to the target cluster.

On each flush it drains pending produce acknowledgements, sets the downstream offset on each OffsetSync, and writes those syncs back to the target’s offset-syncs topic (MM2-compatible).

Implementations§

Source§

impl TargetSink

Source

pub async fn start(params: SinkParams) -> Result<Self, ConnectError>

Build a TargetSink, ensure the offset-syncs topic exists, and connect the producer to the target cluster.

§Errors

Returns ConnectError::Backend if the producer cannot connect or the offset-syncs topic cannot be created.

Source

pub fn drain_offset_syncs(&mut self) -> Vec<OffsetSync>

Return and clear all completed OffsetSync records accumulated since the last call (or since construction).

Trait Implementations§

Source§

impl Sink<(), ReplicatedRecord> for TargetSink

Source§

fn put<'life0, 'async_trait>( &'life0 mut self, records: Vec<ConnectRecord<(), ReplicatedRecord>>, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Accept a batch of replicated records, applying filtering and loop-guard logic, then enqueue produce calls for the accepted records.

Records are dropped (not buffered) when:

  • value is None (tombstone or no payload).
  • Identity-naming loop-guard fires: the record’s __crabka_origin header matches our own source_alias.
  • Residency gate blocks the topic for the target’s zones.
Source§

fn flush<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Await all pending produce acks, write offset-syncs to the target, and flush the producer.

Source§

fn supports_transactions(&self) -> bool

Whether this sink writes atomically and so can take part in exactly-once delivery. When false (the default) the runtime drives it at-least-once and never calls begin / abort.
Source§

fn begin<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Open a transaction enclosing the puts that follow, up to the next commit or abort. Called by the runtime only when supports_transactions is true and only before a non-empty batch. Default: no-op. Read more
Source§

fn commit<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Atomically commit the records written since begin, making them durable. The default (no transaction support) delegates to flush. Read more
Source§

fn abort<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Discard the records written since begin without making them durable. Called by the runtime to roll back a failed interval. Default: no-op. Read more
Source§

fn close<'life0, 'async_trait>( &'life0 mut self, ) -> Pin<Box<dyn Future<Output = Result<(), ConnectError>> + Send + 'async_trait>>
where 'life0: 'async_trait, Self: 'async_trait,

Release any resources held by the sink (connections, buffers). The default is a no-op. After close, no further records are written. 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<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

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