Skip to main content

BarrierCoordinator

Struct BarrierCoordinator 

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

Cross-instance barrier coordination.

Implementations§

Source§

impl BarrierCoordinator

Source

pub fn new(kv: Arc<dyn ClusterKv>) -> Self

Wrap a KV implementation.

Source

pub fn set_leader_lease_store(&self, store: Arc<LeaderLeaseStore>)

Install the durable authority used to validate clustered reversible barrier phases. Without it, clustered Prepare and Aligned traffic fails closed.

Source

pub fn checkpoint_authority( &self, ) -> Result<Arc<LeaderLeaseStore>, ClusterCheckpointAuthorityError>

Exact durable authority installed for clustered barriers and checkpoint decisions.

Embedded and single-node runtimes do not call this path. A cluster runtime that omitted authority wiring fails closed instead of falling back to standalone outcome objects.

§Errors

Returns NotConfigured when durable cluster checkpoint authority is not installed.

Source

pub fn set_leader_election( &mut self, instance_id: NodeId, members_rx: Receiver<Vec<NodeInfo>>, leader_eligible: Arc<AtomicBool>, )

Configure membership used to target active barrier peers. Gossip election is not a barrier authority boundary.

Source

pub fn prepare_received_at( &self, prepare: &BarrierAnnouncement, ) -> Option<Instant>

Local monotonic receipt time for this exact gRPC Prepare.

Source

pub async fn start_server( &self, bind_addr: SocketAddr, advertise_host: Option<String>, ) -> Result<SocketAddr, String>

Bind and run the follower’s direct gRPC barrier sync server.

§Errors

Returns an error string on bind or socket address retrieval failures.

Source

pub async fn confirm_remote_leader_proof( &self, proof: &LeaderProof, deadline: Instant, ) -> Result<bool, String>

Ask one exact remote process to confirm a proof already read from durable authority.

The response echoes only a fresh challenge id. It never returns a process-local or durable fencing token.

§Errors

Fails when the proof, peer address, RPC, acknowledgement, or deadline is invalid.

Source

pub async fn announce(&self, ann: &BarrierAnnouncement) -> Result<(), String>

Leader-side announcement for terminal, aligned, and assignment-less local/KV phases.

§Errors

Assignment-certified Prepare must use Self::announce_prepare. Assignment-certified reversible phases require a started leased barrier server, and Aligned requires the exact Prepare to have completed Self::wait_for_quorum. Other errors propagate validation, encoding, and publication failures.

Source

pub async fn announce_prepare( &self, ann: &BarrierAnnouncement, quorum_window: Duration, ) -> Result<(), String>

Durably publish one assignment-certified Prepare and immediately start its direct fan-out.

§Errors

Rejects a different phase, an assignment-less announcement, a zero/indivisible quorum window, malformed authority, conflicting in-flight Prepare state, or publication failure.

Source

pub fn announcement_watch( &self, ) -> Option<Receiver<Option<BarrierAnnouncement>>>

Watch over gRPC-delivered announcements, for push-driven waits (the decision wait and the Aligned resume gate). None until the gRPC server is started — gossip-KV-only deployments fall back to polling the merged gossip history.

Source

pub async fn ack(&self, ack: &BarrierAck) -> Result<(), String>

Follower-side ack.

§Errors

Returns a string on JSON encode failure.

Source

pub async fn wait_for_quorum( &self, prepare: &BarrierAnnouncement, expected: &[NodeId], deadline: Duration, ) -> QuorumOutcome

Leader-side: wait until quorum or deadline.

Trait Implementations§

Source§

impl Debug for BarrierCoordinator

Source§

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

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

impl Drop for BarrierCoordinator

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

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
§

impl<T> ArchivePointee for T

§

type ArchivedMetadata = ()

The archived version of the pointer metadata for this type.
§

fn pointer_metadata( _: &<T as ArchivePointee>::ArchivedMetadata, ) -> <T as Pointee>::Metadata

Converts some archived metadata to the pointer metadata for itself.
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<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

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
§

impl<T> IntoRequest<T> for T

§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
§

impl<L> LayerExt<L> for L

§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in [Layered].
§

impl<T> LayoutRaw for T

§

fn layout_raw(_: <T as Pointee>::Metadata) -> Result<Layout, LayoutError>

Returns the layout of the type.
§

impl<T, N1, N2> Niching<NichedOption<T, N1>> for N2
where T: SharedNiching<N1, N2>, N1: Niching<T>, N2: Niching<T>,

§

unsafe fn is_niched(niched: *const NichedOption<T, N1>) -> bool

Returns whether the given value has been niched. Read more
§

fn resolve_niched(out: Place<NichedOption<T, N1>>)

Writes data to out indicating that a T is niched.
§

impl<T> Pointee for T

§

type Metadata = ()

The metadata type for pointers and references to this type.
§

impl<T> PolicyExt for T
where T: ?Sized,

§

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
§

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

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

impl<T> WithSubscriber for T

§

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
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more