pub struct PgCursor {
pub sequence: i64,
pub partition: Partition,
}Fields§
§sequence: i64§partition: PartitionImplementations§
Trait Implementations§
impl Copy for PgCursor
Source§impl<'de> Deserialize<'de> for PgCursor
impl<'de> Deserialize<'de> for PgCursor
Source§fn deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>where
__D: Deserializer<'de>,
fn deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>where
__D: Deserializer<'de>,
Deserialize this value from the given Serde deserializer. Read more
impl Eq for PgCursor
Source§impl Ord for PgCursor
impl Ord for PgCursor
1.21.0 (const: unstable) · Source§fn max(self, other: Self) -> Selfwhere
Self: Sized,
fn max(self, other: Self) -> Selfwhere
Self: Sized,
Compares and returns the maximum of two values. Read more
Source§impl PartialOrd for PgCursor
impl PartialOrd for PgCursor
Source§impl PartitionCoordinator<PgCursor> for PgPartitionCoordinator
impl PartitionCoordinator<PgCursor> for PgPartitionCoordinator
async fn heartbeat<'a>( &'a self, scope: &'a CheckpointScope, owner_id: &'a OwnerId, lease_duration: Duration, ) -> Result<()>
async fn live_consumers<'a>( &'a self, scope: &'a CheckpointScope, ) -> Result<usize>
Source§async fn release_consumer<'a>(
&'a self,
scope: &'a CheckpointScope,
owner_id: &'a OwnerId,
) -> Result<()>
async fn release_consumer<'a>( &'a self, scope: &'a CheckpointScope, owner_id: &'a OwnerId, ) -> Result<()>
Remove this consumer’s registration so other consumers stop counting it
as live when computing
target_partition_count. Called by
CoordinatedReader on stream shutdown. Should be idempotent and must
not return an error when the consumer record is missing.Source§async fn claim<'a>(
&'a self,
scope: &'a CheckpointScope,
owner_id: &'a OwnerId,
partition: Partition,
lease_duration: Duration,
) -> Result<Option<PartitionLease<PgCursor>>>
async fn claim<'a>( &'a self, scope: &'a CheckpointScope, owner_id: &'a OwnerId, partition: Partition, lease_duration: Duration, ) -> Result<Option<PartitionLease<PgCursor>>>
Attempt to take a partition. Returns
Ok(Some(lease)) on success or
Ok(None) if another live owner holds it. Increments generation on
every successful claim.Source§async fn renew<'a>(
&'a self,
lease: &'a PartitionLease<PgCursor>,
lease_duration: Duration,
) -> Result<()>
async fn renew<'a>( &'a self, lease: &'a PartitionLease<PgCursor>, lease_duration: Duration, ) -> Result<()>
Extend
lease_until only if (owner_id, generation) still matches.
Returns Err(Error::OwnershipLost(_)) on mismatch.Source§async fn release<'a>(
&'a self,
lease: &'a PartitionLease<PgCursor>,
) -> Result<()>
async fn release<'a>( &'a self, lease: &'a PartitionLease<PgCursor>, ) -> Result<()>
Clear
owner_id and lease_until only if (owner_id, generation)
matches. Increments generation. Returns Err(Error::OwnershipLost(_))
on mismatch.Source§async fn checkpoint<'a>(
&'a self,
lease: &'a PartitionLease<PgCursor>,
cursor: PgCursor,
) -> Result<()>
async fn checkpoint<'a>( &'a self, lease: &'a PartitionLease<PgCursor>, cursor: PgCursor, ) -> Result<()>
Write
cursor only if (owner_id, generation) matches. The update
is monotonic via Cursor::order_key(). Returns
Err(Error::OwnershipLost(_)) on mismatch.Source§impl PartitionableSubscription<PgCursor> for PgSubscription
impl PartitionableSubscription<PgCursor> for PgSubscription
Source§fn with_partitions(self, group: PartitionGroup) -> Self
fn with_partitions(self, group: PartitionGroup) -> Self
Restrict this subscription to a validated group of partitions sharing
the same partition_count. The reader emits a single SQL query per
poll using partition_id = ANY($::bigint[]) instead of one query per
partition. Single-partition uses fall through the default trait impl
which wraps in a singleton group; partition_id = ANY(ARRAY[$1])
plans identically to partition_id = $1 on modern Postgres.
Source§fn with_partition(self, partition: Partition) -> Selfwhere
Self: Sized,
fn with_partition(self, partition: Partition) -> Selfwhere
Self: Sized,
Restrict the subscription to a single partition. The default impl
wraps the partition in a singleton
PartitionGroup and delegates to
PartitionableSubscription::with_partitions. Backends override
only when a single-partition path is materially cheaper than the
multi-partition path (rare).Source§impl StartableSubscription<PgCursor> for PgSubscription
impl StartableSubscription<PgCursor> for PgSubscription
fn with_start(self, start: StartFrom<PgCursor>) -> Self
Source§fn with_starts(self, starts: Vec<StartFrom<C>>) -> Selfwhere
C: Ord,
fn with_starts(self, starts: Vec<StartFrom<C>>) -> Selfwhere
C: Ord,
Seed this subscription with a collection of candidate start
positions. Default behavior: pick the smallest
StartFrom::After(c)
from the vec and delegate to with_start. Other variants
(Earliest, Latest, Timestamp) are ignored by the default impl —
readers that support fan-in, dual historic+live consumption, or
topology-aware resume must override.impl StructuralPartialEq for PgCursor
Auto Trait Implementations§
impl Freeze for PgCursor
impl RefUnwindSafe for PgCursor
impl Send for PgCursor
impl Sync for PgCursor
impl Unpin for PgCursor
impl UnsafeUnpin for PgCursor
impl UnwindSafe for PgCursor
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Mutably borrows from an owned value. Read more
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<Q, K> Comparable<K> for Q
impl<Q, K> Comparable<K> for Q
impl<T> DeserializeOwned for Twhere
T: for<'de> Deserialize<'de>,
Source§impl<Q, K> Equivalent<K> for Q
impl<Q, K> Equivalent<K> for Q
Source§fn equivalent(&self, key: &K) -> bool
fn equivalent(&self, key: &K) -> bool
Compare self to
key and return true if they are equal.Source§impl<Q, K> Equivalent<K> for Q
impl<Q, K> Equivalent<K> for Q
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
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