Skip to main content

ReplicatedRecord

Struct ReplicatedRecord 

Source
pub struct ReplicatedRecord {
    pub topic: String,
    pub partition: i32,
    pub offset: i64,
    pub timestamp: i64,
    pub key: Option<Bytes>,
    pub value: Option<Bytes>,
    pub headers: Vec<(String, Option<Bytes>)>,
}
Expand description

One record from the source cluster, carried as the connect-value type V through the crabka_connect runtime.

Because crabka_connect::ConnectRecord has no topic/partition fields, all envelope metadata is bundled here alongside the raw payload.

Fields§

§topic: String

The source topic name.

§partition: i32

The source partition index.

§offset: i64

The source offset (0-based).

§timestamp: i64

Record timestamp in epoch milliseconds.

§key: Option<Bytes>

Record key, or None for a null key.

§value: Option<Bytes>

Record value, or None for a tombstone.

§headers: Vec<(String, Option<Bytes>)>

Per-record headers in declaration order.

Trait Implementations§

Source§

impl Clone for ReplicatedRecord

Source§

fn clone(&self) -> ReplicatedRecord

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 Debug for ReplicatedRecord

Source§

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

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

impl Eq for ReplicatedRecord

Source§

impl PartialEq for ReplicatedRecord

Source§

fn eq(&self, other: &ReplicatedRecord) -> bool

Tests for self and other values to be equal, and is used by ==.
1.0.0 (const: unstable) · Source§

fn ne(&self, other: &Rhs) -> bool

Tests for !=. The default implementation is almost always sufficient, and should not be overridden without very good reason.
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
Source§

impl Source<(), ReplicatedRecord> for SourceConsumer

Source§

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

Poll the source cluster for the next record.

Returns Ok(None) when the consumer is momentarily caught up (the runtime should back off and retry). Returns Ok(Some(_)) with the next ReplicatedRecord otherwise.

§Errors

Returns ConnectError::Backend if the underlying consumer poll fails.

Source§

fn checkpoint(&self) -> Option<SourceOffset>

Snapshot the current read positions for all partitions seen so far.

Returns None before the first successful poll (nothing to commit yet).

Source§

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

Restore the read position from a previously-checkpointed SourceOffset.

The runtime calls this once before the first poll, passing the position loaded from the durable checkpoint store on the target. Each position entry is keyed "<topic>-<partition>"OffsetValue::Long(next_offset) (the value checkpoint wrote: last_consumed + 1). We decode each key back into (topic, partition) and hand the offset to the consumer’s seek.

The consumer holds each seek as pending and materialises it at the top of the first poll that sees the partition assigned — after the group’s post-assignment offset prime, but before any Fetch — so the sought offset is the one fetched. That makes restart resume from the last fully-committed record rather than re-reading the topic from offset 0: no record below the sought offset is re-delivered, and none above it is skipped (no data gap). Delivery remains at-least-once — a crash between a sink flush and the checkpoint save can re-deliver the in-flight batch, but never lose a record.

A malformed key (no -, or a non-integer partition/offset) is skipped with a warning rather than failing the restore: one corrupt entry must not strand recovery for the partitions that decoded cleanly.

§Errors

Returns ConnectError::Backend if the consumer is already closed.

Source§

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

Close the underlying consumer, sending LeaveGroup so a restarted replicator can rejoin the group immediately instead of waiting out the departed member’s session timeout.

§Errors

Returns ConnectError::Backend if the consumer fails to close cleanly.

Source§

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

Acknowledge that offset is durable end-to-end. Read more
Source§

impl StructuralPartialEq for ReplicatedRecord

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> 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<Q, K> Equivalent<K> for Q
where Q: Eq + ?Sized, K: Borrow<Q> + ?Sized,

Source§

fn equivalent(&self, key: &K) -> bool

Checks if this value is equivalent to the given key. Read more
Source§

impl<Q, K> Equivalent<K> for Q
where Q: Eq + ?Sized, K: Borrow<Q> + ?Sized,

Source§

fn equivalent(&self, key: &K) -> bool

Checks if this value is equivalent to the given key. Read more
Source§

impl<Q, K> Equivalent<K> for Q
where Q: Eq + ?Sized, K: Borrow<Q> + ?Sized,

Source§

fn equivalent(&self, key: &K) -> bool

Compare self to key and return true if they are equal.
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> 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<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