Skip to main content

laminar_core/cluster/control/
barrier.rs

1//! Cross-instance barrier protocol. Direct gRPC leader-to-follower calls
2//! under `cluster`, falling back to gossip-KV announce/ack/poll.
3
4#[cfg(feature = "cluster")]
5use std::sync::atomic::AtomicBool;
6use std::sync::Arc;
7use std::time::{Duration, Instant};
8
9use async_trait::async_trait;
10use parking_lot::Mutex;
11use rustc_hash::{FxHashMap, FxHashSet};
12use serde::{Deserialize, Serialize};
13
14use crate::checkpoint::CheckpointWatermark;
15use crate::cluster::discovery::NodeId;
16#[cfg(feature = "cluster")]
17use crate::cluster::discovery::{NodeInfo, NodeState};
18use crate::state::CheckpointAttempt;
19#[cfg(feature = "cluster")]
20use tokio::sync::watch;
21
22/// KV key for the leader's barrier announcement.
23pub const ANNOUNCEMENT_KEY: &str = "control:barrier";
24
25/// KV key for a follower's barrier ack.
26pub const ACK_KEY: &str = "control:barrier-ack";
27
28/// Gossip KV key used by follower barrier servers to advertise their bound address.
29#[cfg(feature = "cluster")]
30pub const BARRIER_ADDR_KEY: &str = "barrier:addr";
31
32#[cfg(feature = "cluster")]
33const BARRIER_ENDPOINT_VERSION: u8 = 1;
34
35#[cfg(feature = "cluster")]
36const MAX_BARRIER_ENDPOINT_BYTES: usize = 1_024;
37
38/// Process identity attached to one control-plane endpoint. The durable process lease remains the
39/// authority; this value prevents a stable node id from resolving to a different process boot.
40#[cfg(feature = "cluster")]
41#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
42#[serde(deny_unknown_fields)]
43struct BarrierProcessIdentity {
44    node_id: u64,
45    boot_incarnation: uuid::Uuid,
46    process_term: u64,
47}
48
49#[cfg(feature = "cluster")]
50impl BarrierProcessIdentity {
51    fn from_process_lease(lease: &super::ProcessLease) -> Result<Self, String> {
52        lease
53            .validate(lease.node)
54            .map_err(|error| error.to_string())?;
55        Ok(Self {
56            node_id: lease.node.0,
57            boot_incarnation: lease.owner,
58            process_term: lease.term,
59        })
60    }
61
62    const fn is_canonical(self) -> bool {
63        self.node_id != 0 && !self.boot_incarnation.is_nil() && self.process_term != 0
64    }
65
66    fn matches_participant(self, participant: &crate::checkpoint::CheckpointParticipant) -> bool {
67        self.node_id == participant.node_id && self.boot_incarnation == participant.boot_incarnation
68    }
69}
70
71#[cfg(feature = "cluster")]
72#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
73#[serde(deny_unknown_fields)]
74struct BarrierEndpointRecord {
75    version: u8,
76    address: String,
77    process: BarrierProcessIdentity,
78}
79
80#[cfg(feature = "cluster")]
81impl BarrierEndpointRecord {
82    fn new(address: String, process: BarrierProcessIdentity) -> Result<Self, String> {
83        let record = Self {
84            version: BARRIER_ENDPOINT_VERSION,
85            address,
86            process,
87        };
88        record.validate()?;
89        Ok(record)
90    }
91
92    fn validate(&self) -> Result<(), String> {
93        if self.version != BARRIER_ENDPOINT_VERSION {
94            return Err(format!(
95                "unsupported cluster control endpoint version {}",
96                self.version
97            ));
98        }
99        if self.address.is_empty() || self.address.len() > MAX_BARRIER_ENDPOINT_BYTES / 2 {
100            return Err("cluster control endpoint address is empty or oversized".into());
101        }
102        if !self.process.is_canonical() {
103            return Err("cluster control endpoint process identity is not canonical".into());
104        }
105        Ok(())
106    }
107
108    fn encode(&self) -> Result<String, String> {
109        let encoded = serde_json::to_string(self).map_err(|error| error.to_string())?;
110        if encoded.len() > MAX_BARRIER_ENDPOINT_BYTES {
111            return Err("cluster control endpoint advertisement is oversized".into());
112        }
113        Ok(encoded)
114    }
115}
116
117#[cfg(feature = "cluster")]
118#[derive(Debug, Clone, Copy)]
119struct ExpectedBarrierProcess {
120    node_id: u64,
121    boot_incarnation: uuid::Uuid,
122    process_term: Option<u64>,
123}
124
125#[cfg(feature = "cluster")]
126impl ExpectedBarrierProcess {
127    const fn participant(node_id: u64, boot_incarnation: uuid::Uuid) -> Self {
128        Self {
129            node_id,
130            boot_incarnation,
131            process_term: None,
132        }
133    }
134
135    const fn exact(process: &crate::checkpoint::LeaderProofOwner) -> Self {
136        Self {
137            node_id: process.node_id,
138            boot_incarnation: process.boot_id,
139            process_term: Some(process.process_term),
140        }
141    }
142
143    fn matches(self, actual: BarrierProcessIdentity) -> bool {
144        self.node_id == actual.node_id
145            && self.boot_incarnation == actual.boot_incarnation
146            && self
147                .process_term
148                .is_none_or(|term| term == actual.process_term)
149    }
150}
151
152#[cfg(feature = "cluster")]
153fn decode_barrier_endpoint(raw: &str) -> Result<(String, Option<BarrierProcessIdentity>), String> {
154    if raw.len() > MAX_BARRIER_ENDPOINT_BYTES {
155        return Err("cluster control endpoint advertisement is oversized".into());
156    }
157    if !raw.trim_start().starts_with('{') {
158        if raw.is_empty() {
159            return Err("cluster control endpoint address is empty".into());
160        }
161        return Ok((raw.to_string(), None));
162    }
163    let record: BarrierEndpointRecord = serde_json::from_str(raw)
164        .map_err(|error| format!("invalid cluster control endpoint advertisement: {error}"))?;
165    record.validate()?;
166    Ok((record.address, Some(record.process)))
167}
168
169#[cfg(feature = "cluster")]
170fn encode_barrier_endpoint(
171    address: &str,
172    process: Option<BarrierProcessIdentity>,
173) -> Result<String, String> {
174    if address.is_empty() || address.len() > MAX_BARRIER_ENDPOINT_BYTES / 2 {
175        return Err("cluster control endpoint address is empty or oversized".into());
176    }
177    process.map_or_else(
178        || Ok(address.to_string()),
179        |process| BarrierEndpointRecord::new(address.to_string(), process)?.encode(),
180    )
181}
182
183/// Upper bound for a non-Prepare phase notification round. The durable KV
184/// announcement is authoritative; direct gRPC delivery is only the low-latency
185/// path and must never hold the checkpoint coordinator indefinitely.
186#[cfg(feature = "cluster")]
187const PHASE_RPC_TIMEOUT: Duration = Duration::from_secs(3);
188
189#[cfg(feature = "cluster")]
190const PREPARE_RETRY_INITIAL_BACKOFF: Duration = Duration::from_millis(10);
191
192#[cfg(feature = "cluster")]
193const PREPARE_RETRY_MAX_BACKOFF: Duration = Duration::from_millis(250);
194
195/// Barrier phase.
196#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
197pub enum Phase {
198    /// Align the shuffle, capture state locally, ack. The durable tail
199    /// (sink pre-commit, manifest, uploads) runs after the ack.
200    Prepare,
201    /// Every node has aligned + captured this epoch (full-membership
202    /// capture quorum). Pipelines may resume the next epoch; the epoch
203    /// is NOT yet restorable.
204    Aligned,
205    /// Durability gate passed; commit sinks. The epoch is restorable.
206    Commit,
207    /// Prepare failed; roll back.
208    Abort,
209}
210
211const fn is_terminal_phase(phase: Phase) -> bool {
212    matches!(phase, Phase::Commit | Phase::Abort)
213}
214
215/// Leader-written barrier announcement.
216#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
217pub struct BarrierAnnouncement {
218    /// Monotonic checkpoint ID retained in the wire field named `epoch`.
219    pub epoch: u64,
220    /// The same nonzero coordinator-assigned checkpoint ID.
221    pub checkpoint_id: u64,
222    /// Exact clustered assignment cut captured when this attempt was admitted. Required on every
223    /// clustered `Prepare` and retained on terminal phases for exact follower validation.
224    #[serde(default)]
225    pub assignment_fence: Option<super::CheckpointAssignmentFence>,
226    /// Exact durable leader term that issued this announcement. Clustered reversible phases
227    /// (`Prepare` and `Aligned`) are rejected unless this proof is present and still live.
228    /// Terminal notifications carry it only for diagnostics; their authority comes from the
229    /// immutable durable checkpoint outcome.
230    #[serde(default)]
231    pub leader_proof: Option<super::LeaderProof>,
232    /// Phase this announcement signals.
233    pub phase: Phase,
234    /// Reserved for unaligned/other flags.
235    pub flags: u64,
236}
237
238fn announcement_attempt(ann: &BarrierAnnouncement) -> CheckpointAttempt {
239    CheckpointAttempt::new(ann.epoch, ann.checkpoint_id)
240}
241
242fn validate_announcement_attempt(ann: &BarrierAnnouncement) -> Result<(), String> {
243    if announcement_attempt(ann).is_canonical() {
244        Ok(())
245    } else {
246        Err(format!(
247            "barrier announcement must use one nonzero canonical checkpoint ID; received epoch {} and checkpoint ID {}",
248            ann.epoch, ann.checkpoint_id
249        ))
250    }
251}
252
253fn validate_ack_attempt(ack: &BarrierAck) -> Result<(), String> {
254    if CheckpointAttempt::new(ack.epoch, ack.checkpoint_id).is_canonical() {
255        Ok(())
256    } else {
257        Err(format!(
258            "barrier acknowledgement must use one nonzero canonical checkpoint ID; received epoch {} and checkpoint ID {}",
259            ack.epoch, ack.checkpoint_id
260        ))
261    }
262}
263
264#[cfg(feature = "cluster")]
265fn validate_wire_checkpoint_attempt(epoch: u64, checkpoint_id: u64) -> Result<(), tonic::Status> {
266    if CheckpointAttempt::new(epoch, checkpoint_id).is_canonical() {
267        Ok(())
268    } else {
269        Err(tonic::Status::invalid_argument(
270            "Barrier request must use one nonzero canonical checkpoint ID",
271        ))
272    }
273}
274
275fn same_announcement_identity(left: &BarrierAnnouncement, right: &BarrierAnnouncement) -> bool {
276    announcement_attempt(left) == announcement_attempt(right)
277        && left.assignment_fence == right.assignment_fence
278        && left.leader_proof == right.leader_proof
279        && left.flags == right.flags
280}
281
282/// Merge one exact attempt from a single history. A successor may settle a reversible phase under
283/// a new assignment and leader proof, but reversible equivocation and opposing decisions fail
284/// closed.
285fn merge_history_exact(
286    current: BarrierAnnouncement,
287    incoming: BarrierAnnouncement,
288    history: &str,
289) -> Result<BarrierAnnouncement, String> {
290    if current.flags != incoming.flags {
291        return Err(format!(
292            "conflicting {history} barrier flags for exact attempt ({}, {})",
293            current.epoch, current.checkpoint_id
294        ));
295    }
296    if current.phase == incoming.phase && current != incoming {
297        return Err(format!(
298            "conflicting {history} {:?} payloads for exact attempt ({}, {})",
299            current.phase, current.epoch, current.checkpoint_id
300        ));
301    }
302    if is_terminal_phase(current.phase) && is_terminal_phase(incoming.phase) {
303        if current.phase != incoming.phase {
304            return Err(format!(
305                "conflicting {history} terminal phases for exact attempt ({}, {})",
306                current.epoch, current.checkpoint_id
307            ));
308        }
309        return Ok(current);
310    }
311    if is_terminal_phase(current.phase) {
312        return Ok(current);
313    }
314    if is_terminal_phase(incoming.phase) {
315        return Ok(incoming);
316    }
317    if !same_announcement_identity(&current, &incoming) {
318        return Err(format!(
319            "conflicting {history} barrier certificates for exact attempt ({}, {})",
320            current.epoch, current.checkpoint_id
321        ));
322    }
323    if current.phase == Phase::Aligned && incoming.phase == Phase::Prepare {
324        Ok(current)
325    } else {
326        Ok(incoming)
327    }
328}
329
330/// Merge direct deliveries without allowing a delayed phase to regress the same exact attempt.
331#[cfg(feature = "cluster")]
332fn merge_direct_announcement(
333    current: BarrierAnnouncement,
334    incoming: BarrierAnnouncement,
335) -> Result<BarrierAnnouncement, String> {
336    validate_announcement_attempt(&current)?;
337    validate_announcement_attempt(&incoming)?;
338    match incoming.checkpoint_id.cmp(&current.checkpoint_id) {
339        std::cmp::Ordering::Greater => Ok(incoming),
340        std::cmp::Ordering::Less => Ok(current),
341        std::cmp::Ordering::Equal => merge_history_exact(current, incoming, "direct"),
342    }
343}
344
345/// Merge the low-latency direct value with the durable leader announcement.
346/// At the same exact attempt a terminal KV value is the decision authority;
347/// otherwise phase progress is monotonic while gossip catches up.
348fn merge_observed_announcement(
349    grpc: BarrierAnnouncement,
350    durable: BarrierAnnouncement,
351) -> Result<BarrierAnnouncement, String> {
352    validate_announcement_attempt(&grpc)?;
353    validate_announcement_attempt(&durable)?;
354    match durable.checkpoint_id.cmp(&grpc.checkpoint_id) {
355        std::cmp::Ordering::Greater => Ok(durable),
356        std::cmp::Ordering::Less => Ok(grpc),
357        std::cmp::Ordering::Equal => {
358            if grpc.phase == durable.phase && !is_terminal_phase(grpc.phase) && grpc != durable {
359                Err(format!(
360                    "conflicting direct and durable {:?} payloads for exact attempt ({}, {})",
361                    grpc.phase, grpc.epoch, grpc.checkpoint_id
362                ))
363            } else if is_terminal_phase(durable.phase) {
364                Ok(durable)
365            } else if is_terminal_phase(grpc.phase) {
366                Ok(grpc)
367            } else if !same_announcement_identity(&grpc, &durable) {
368                Err(format!(
369                    "conflicting direct and durable certificates for exact attempt ({}, {})",
370                    grpc.epoch, grpc.checkpoint_id
371                ))
372            } else if grpc.phase == Phase::Aligned && durable.phase == Phase::Prepare {
373                Ok(grpc)
374            } else {
375                Ok(durable)
376            }
377        }
378    }
379}
380
381/// Merge per-node durable histories for leader reclamation. A successor terminal may carry a new
382/// certificate; reversible phases must retain one certificate and terminal outcomes must agree.
383fn merge_scanned_announcement(
384    current: BarrierAnnouncement,
385    incoming: BarrierAnnouncement,
386) -> Result<BarrierAnnouncement, String> {
387    validate_announcement_attempt(&current)?;
388    validate_announcement_attempt(&incoming)?;
389    match incoming.checkpoint_id.cmp(&current.checkpoint_id) {
390        std::cmp::Ordering::Greater => Ok(incoming),
391        std::cmp::Ordering::Less => Ok(current),
392        std::cmp::Ordering::Equal => merge_history_exact(current, incoming, "durable"),
393    }
394}
395
396fn validate_publication_order(
397    current: &BarrierAnnouncement,
398    incoming: &BarrierAnnouncement,
399) -> Result<(), String> {
400    validate_announcement_attempt(current)?;
401    validate_announcement_attempt(incoming)?;
402    match incoming.checkpoint_id.cmp(&current.checkpoint_id) {
403        std::cmp::Ordering::Greater => Ok(()),
404        std::cmp::Ordering::Less => Err(format!(
405            "stale barrier publication ({}, {}) cannot replace newer admitted attempt ({}, {})",
406            incoming.epoch, incoming.checkpoint_id, current.epoch, current.checkpoint_id
407        )),
408        std::cmp::Ordering::Equal if current == incoming => Ok(()),
409        std::cmp::Ordering::Equal => {
410            let merged =
411                merge_history_exact(current.clone(), incoming.clone(), "local publication")?;
412            if merged == *incoming {
413                Ok(())
414            } else {
415                Err(format!(
416                    "barrier publication cannot regress exact attempt ({}, {}) from {:?} to {:?}",
417                    current.epoch, current.checkpoint_id, current.phase, incoming.phase
418                ))
419            }
420        }
421    }
422}
423
424fn validate_scanned_announcements(
425    mut announcements: Vec<BarrierAnnouncement>,
426) -> Result<Option<BarrierAnnouncement>, String> {
427    for announcement in &announcements {
428        validate_announcement_attempt(announcement)?;
429    }
430    // Group exact attempts and audit every reversible certificate before a terminal can absorb it.
431    // Canonical attempt IDs are the only ordering authority.
432    announcements.sort_unstable_by_key(|announcement| {
433        (
434            announcement.checkpoint_id,
435            is_terminal_phase(announcement.phase),
436        )
437    });
438    announcements
439        .into_iter()
440        .try_fold(None, |highest, announcement| {
441            Ok(Some(match highest {
442                None => announcement,
443                Some(current) => merge_scanned_announcement(current, announcement)?,
444            }))
445        })
446}
447
448/// Follower ack. `ok = false` forces the leader to abort instead of wait.
449#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
450pub struct BarrierAck {
451    /// Canonical checkpoint ID retained in the wire field named `epoch`.
452    pub epoch: u64,
453    /// The same nonzero coordinator-assigned checkpoint ID being acknowledged.
454    #[serde(default)]
455    pub checkpoint_id: u64,
456    /// SHA-256 binding of the announcement's assignment certificate.
457    #[serde(default)]
458    pub assignment_digest: Option<[u8; 32]>,
459    /// `false` = snapshot failed locally; leader should abort.
460    pub ok: bool,
461    /// Free-text error; populated when `ok = false`.
462    pub error: Option<String>,
463    /// Follower event-time state at ack time. Uninitialized inputs block advancement; only
464    /// explicitly idle inputs are excluded from the active cluster minimum.
465    #[serde(default)]
466    pub watermark: CheckpointWatermark,
467}
468
469/// Outcome of `wait_for_quorum`.
470#[derive(Debug, Clone, PartialEq, Eq)]
471pub enum QuorumOutcome {
472    /// All expected peers acked with `ok = true`.
473    Reached {
474        /// Peers that acked successfully.
475        acks: Vec<NodeId>,
476        /// Safe aggregate across all required followers.
477        follower_watermark: CheckpointWatermark,
478    },
479    /// Deadline expired with at least one peer silent.
480    TimedOut {
481        /// Peers that did ack.
482        got: Vec<NodeId>,
483        /// Peers that didn't.
484        missing: Vec<NodeId>,
485    },
486    /// At least one peer acked `ok = false`.
487    Failed {
488        /// `(peer, error_message)` for every failed ack.
489        failures: Vec<(NodeId, String)>,
490    },
491}
492
493/// Gossip-KV seam.
494#[async_trait]
495pub trait ClusterKv: Send + Sync + 'static {
496    /// Write `value` to this instance's `key` slot (overwrites).
497    async fn write(&self, key: &str, value: String);
498    /// Write with transport failure reporting when the backend supports it.
499    ///
500    /// Fast gossip implementations may use the default because their write API has no result.
501    /// Durable control implementations override this so recovery never treats a dropped write as
502    /// successful.
503    ///
504    /// # Errors
505    /// Durable implementations return a transport or storage error when the value was not
506    /// accepted by their authority.
507    async fn write_checked(&self, key: &str, value: String) -> Result<(), String> {
508        self.write(key, value).await;
509        Ok(())
510    }
511    /// Read `key` from `who`'s slot.
512    async fn read_from(&self, who: NodeId, key: &str) -> Option<String>;
513    /// Read with transport failure reporting when the backend supports it.
514    ///
515    /// # Errors
516    /// Durable implementations return a transport or storage error. A genuinely absent key is
517    /// `Ok(None)`.
518    async fn read_from_checked(&self, who: NodeId, key: &str) -> Result<Option<String>, String> {
519        Ok(self.read_from(who, key).await)
520    }
521    /// Every visible instance's value for `key`.
522    async fn scan(&self, key: &str) -> Vec<(NodeId, String)>;
523    /// Scan with transport failure reporting when the backend supports it.
524    ///
525    /// # Errors
526    /// Durable implementations fail the whole scan when any visible participant cannot be read.
527    async fn scan_checked(&self, key: &str) -> Result<Vec<(NodeId, String)>, String> {
528        Ok(self.scan(key).await)
529    }
530}
531
532/// In-memory KV for tests.
533#[derive(Debug)]
534pub struct InMemoryKv {
535    local_id: NodeId,
536    state: Mutex<FxHashMap<(NodeId, String), String>>,
537}
538
539impl InMemoryKv {
540    /// Create a new in-memory KV identified as `local_id`.
541    #[must_use]
542    pub fn new(local_id: NodeId) -> Self {
543        Self {
544            local_id,
545            state: Mutex::new(FxHashMap::default()),
546        }
547    }
548
549    /// Seed a remote peer's state for tests.
550    pub fn seed(&self, peer: NodeId, key: &str, value: String) {
551        self.state.lock().insert((peer, key.to_string()), value);
552    }
553}
554
555#[async_trait]
556impl ClusterKv for InMemoryKv {
557    async fn write(&self, key: &str, value: String) {
558        self.state
559            .lock()
560            .insert((self.local_id, key.to_string()), value);
561    }
562
563    async fn read_from(&self, who: NodeId, key: &str) -> Option<String> {
564        self.state.lock().get(&(who, key.to_string())).cloned()
565    }
566
567    async fn scan(&self, key: &str) -> Vec<(NodeId, String)> {
568        self.state
569            .lock()
570            .iter()
571            .filter(|((_, k), _)| k == key)
572            .map(|((n, _), v)| (*n, v.clone()))
573            .collect()
574    }
575}
576
577#[cfg(feature = "cluster")]
578#[allow(
579    clippy::doc_markdown,
580    clippy::default_trait_access,
581    clippy::missing_const_for_fn,
582    clippy::must_use_candidate,
583    clippy::too_many_lines,
584    missing_docs
585)]
586pub(crate) mod barrier_v1 {
587    tonic::include_proto!("laminar.barrier.v1");
588}
589
590#[cfg(feature = "cluster")]
591type BarrierFlavor = crossfire::mpsc::Array<BarrierAnnouncement>;
592
593/// Full retry identity for direct barrier traffic. `CheckpointAttempt` alone is
594/// insufficient because an assignment rotation can leave delayed traffic with
595/// the same epoch/checkpoint pair but a different participant cut.
596#[cfg(feature = "cluster")]
597#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
598struct BarrierIdentity {
599    attempt: CheckpointAttempt,
600    assignment_digest: Option<[u8; 32]>,
601}
602
603#[cfg(feature = "cluster")]
604impl BarrierIdentity {
605    fn from_announcement(ann: &BarrierAnnouncement) -> Self {
606        Self {
607            attempt: CheckpointAttempt::new(ann.epoch, ann.checkpoint_id),
608            assignment_digest: ann
609                .assignment_fence
610                .as_ref()
611                .map(super::CheckpointAssignmentFence::digest),
612        }
613    }
614
615    const fn from_ack(ack: &BarrierAck) -> Self {
616        Self {
617            attempt: CheckpointAttempt::new(ack.epoch, ack.checkpoint_id),
618            assignment_digest: ack.assignment_digest,
619        }
620    }
621}
622
623#[cfg(feature = "cluster")]
624const MAX_RETAINED_BARRIER_IDENTITIES: usize = 256;
625#[cfg(feature = "cluster")]
626const MAX_PREPARE_WAITERS_PER_IDENTITY: usize = 32;
627#[cfg(feature = "cluster")]
628const PREPARE_RPC_TIMEOUT: Duration = Duration::from_secs(30);
629
630/// Follower-side state shared by the Prepare RPC and the local checkpoint
631/// completion path. One exact Prepare may have several retrying RPC waiters;
632/// they must all receive the same immutable local result.
633#[cfg(feature = "cluster")]
634#[derive(Default)]
635struct PrepareAckState {
636    pending: FxHashMap<BarrierIdentity, Vec<PendingPrepareWaiter>>,
637    completed: FxHashMap<BarrierIdentity, BarrierAck>,
638    received_at: FxHashMap<BarrierIdentity, std::time::Instant>,
639    next_waiter_id: u64,
640}
641
642#[cfg(feature = "cluster")]
643struct PendingPrepareWaiter {
644    id: u64,
645    response: tokio::sync::oneshot::Sender<BarrierAck>,
646}
647
648/// Cancellation-safe removal for one Prepare RPC registration. Tonic may drop
649/// a handler future when the client deadline expires, so cleanup cannot rely on
650/// reaching the explicit timeout branch below.
651#[cfg(feature = "cluster")]
652struct PrepareWaiterRegistration {
653    state: Arc<parking_lot::Mutex<PrepareAckState>>,
654    identity: BarrierIdentity,
655    waiter_id: u64,
656}
657
658#[cfg(feature = "cluster")]
659impl Drop for PrepareWaiterRegistration {
660    fn drop(&mut self) {
661        let mut state = self.state.lock();
662        let remove_entry = state
663            .pending
664            .get_mut(&self.identity)
665            .is_some_and(|waiters| {
666                waiters.retain(|waiter| waiter.id != self.waiter_id);
667                waiters.is_empty()
668            });
669        if remove_entry {
670            state.pending.remove(&self.identity);
671        }
672    }
673}
674
675#[cfg(feature = "cluster")]
676impl PrepareAckState {
677    fn trim_bounded<K, V>(entries: &mut FxHashMap<K, V>)
678    where
679        K: Copy + Eq + std::hash::Hash,
680    {
681        while entries.len() > MAX_RETAINED_BARRIER_IDENTITIES {
682            let Some(victim) = entries.keys().next().copied() else {
683                break;
684            };
685            entries.remove(&victim);
686        }
687    }
688
689    fn next_waiter_id(&mut self) -> u64 {
690        loop {
691            self.next_waiter_id = self.next_waiter_id.wrapping_add(1);
692            let candidate = self.next_waiter_id;
693            if self
694                .pending
695                .values()
696                .flatten()
697                .all(|waiter| waiter.id != candidate)
698            {
699                return candidate;
700            }
701        }
702    }
703
704    fn record_ack(&mut self, identity: BarrierIdentity, ack: &BarrierAck) -> BarrierAck {
705        use std::collections::hash_map::Entry;
706
707        let cached = match self.completed.entry(identity) {
708            Entry::Vacant(entry) => entry.insert(ack.clone()),
709            Entry::Occupied(mut entry) => {
710                if entry.get().ok && !ack.ok {
711                    entry.insert(ack.clone());
712                }
713                entry.into_mut()
714            }
715        }
716        .clone();
717        Self::trim_bounded(&mut self.completed);
718        cached
719    }
720
721    fn record_receipt(&mut self, identity: BarrierIdentity) {
722        self.received_at
723            .entry(identity)
724            .or_insert_with(std::time::Instant::now);
725        Self::trim_bounded(&mut self.received_at);
726    }
727}
728
729/// Per-peer barrier gRPC client pool. A stable node id may be reused by a replacement process,
730/// so every cached channel is bound to the boot incarnation certified by its assignment or
731/// leader proof.
732#[cfg(feature = "cluster")]
733#[derive(Clone)]
734struct BarrierClientEntry {
735    process: Option<BarrierProcessIdentity>,
736    client: barrier_v1::barrier_sync_client::BarrierSyncClient<tonic::transport::Channel>,
737}
738
739#[cfg(feature = "cluster")]
740type BarrierClientPool = Arc<parking_lot::Mutex<FxHashMap<NodeId, BarrierClientEntry>>>;
741
742#[cfg(feature = "cluster")]
743#[derive(Debug)]
744enum BarrierClientResolutionError {
745    ProcessMismatch,
746    Invalid(String),
747}
748
749#[cfg(feature = "cluster")]
750impl std::fmt::Display for BarrierClientResolutionError {
751    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
752        match self {
753            Self::ProcessMismatch => formatter.write_str("endpoint belongs to a different process"),
754            Self::Invalid(error) => formatter.write_str(error),
755        }
756    }
757}
758
759#[cfg(feature = "cluster")]
760fn barrier_client_process_matches(
761    expected: Option<ExpectedBarrierProcess>,
762    actual: Option<BarrierProcessIdentity>,
763) -> bool {
764    match (expected, actual) {
765        (None, None) => true,
766        (Some(expected), Some(actual)) => expected.matches(actual),
767        (None, Some(_)) | (Some(_), None) => false,
768    }
769}
770
771#[cfg(feature = "cluster")]
772struct GrpcState {
773    /// Latest gRPC-delivered announcement, fed in arrival order by the
774    /// relay task draining the incoming queue. Latest-wins semantics
775    /// (matching the gossip-KV fallback) so concurrent observers — the
776    /// pipeline's resume gate and the background durable tail — never
777    /// steal announcements from each other.
778    latest_rx: watch::Receiver<Option<BarrierAnnouncement>>,
779    /// Local half of the same ordered announcement stream used by the gRPC server. Cluster
780    /// leaders use it to observe their own reversible Aligned notification without polling the
781    /// durable fallback.
782    incoming_tx: crossfire::MAsyncTx<BarrierFlavor>,
783    merge_error: Arc<parking_lot::Mutex<Option<String>>>,
784    prepare_acks: Arc<parking_lot::Mutex<PrepareAckState>>,
785    /// The one clustered Prepare admitted after durable publication. It progresses from pending
786    /// fan-out through claimed quorum collection to an exact quorum-ready Aligned admission.
787    prepare_fanout: parking_lot::Mutex<Option<PrepareFanoutState>>,
788    clients: BarrierClientPool,
789    server_handle: Arc<parking_lot::Mutex<Option<tokio::task::JoinHandle<()>>>>,
790    relay_handle: Arc<parking_lot::Mutex<Option<tokio::task::JoinHandle<()>>>>,
791    advertise_addr: String,
792    local_process: Arc<std::sync::OnceLock<BarrierProcessIdentity>>,
793}
794
795#[cfg(feature = "cluster")]
796fn abort_grpc_tasks(state: &GrpcState) {
797    if let Some(handle) = state.server_handle.lock().take() {
798        handle.abort();
799    }
800    if let Some(handle) = state.relay_handle.lock().take() {
801        handle.abort();
802    }
803}
804
805#[cfg(feature = "cluster")]
806type ActiveLeaderState = Option<(NodeId, watch::Receiver<Vec<NodeInfo>>, Arc<AtomicBool>)>;
807
808#[cfg(feature = "cluster")]
809fn leader_proof_challenge_from_wire(bytes: &[u8]) -> Result<uuid::Uuid, tonic::Status> {
810    let challenge = uuid::Uuid::from_slice(bytes).map_err(|_| {
811        tonic::Status::invalid_argument("Leader proof challenge must contain exactly 16 bytes")
812    })?;
813    if challenge.is_nil() {
814        return Err(tonic::Status::invalid_argument(
815            "Leader proof challenge must be nonzero",
816        ));
817    }
818    Ok(challenge)
819}
820
821#[cfg(feature = "cluster")]
822fn leader_proof_ack_matches(challenge: uuid::Uuid, acknowledged: &[u8]) -> bool {
823    acknowledged == challenge.as_bytes()
824}
825
826#[cfg(feature = "cluster")]
827pub(crate) type LocalLeaderProofProvider =
828    Arc<dyn Fn() -> Option<super::LeaderProof> + Send + Sync>;
829
830#[cfg(feature = "cluster")]
831struct GrpcBarrierServer {
832    incoming_tx: crossfire::MAsyncTx<BarrierFlavor>,
833    prepare_acks: Arc<parking_lot::Mutex<PrepareAckState>>,
834    leader_lease_store: Arc<parking_lot::Mutex<Option<Arc<super::LeaderLeaseStore>>>>,
835    local_leader_proof: Arc<parking_lot::Mutex<Option<LocalLeaderProofProvider>>>,
836    local_process: Arc<std::sync::OnceLock<BarrierProcessIdentity>>,
837    process_lease_deadline: Arc<std::sync::OnceLock<Arc<super::LeaseDeadline>>>,
838}
839
840#[cfg(feature = "cluster")]
841#[derive(Clone, Default)]
842struct WireAssignmentFence {
843    version: u64,
844    vnode_count: u32,
845    map_digest: Vec<u8>,
846    participants: Vec<barrier_v1::CheckpointParticipant>,
847}
848
849#[cfg(feature = "cluster")]
850fn leader_proof_from_wire(
851    proof: Option<barrier_v1::LeaderProof>,
852) -> Result<Option<super::LeaderProof>, tonic::Status> {
853    proof
854        .map(|proof| {
855            let boot = uuid::Uuid::from_slice(&proof.boot_id).map_err(|_| {
856                tonic::Status::invalid_argument(
857                    "Leader proof boot identity must contain exactly 16 bytes",
858                )
859            })?;
860            if proof.node_id == 0
861                || boot.is_nil()
862                || proof.process_term == 0
863                || proof.fencing_token == 0
864            {
865                return Err(tonic::Status::invalid_argument(
866                    "Leader proof identity and fencing token must be nonzero",
867                ));
868            }
869            Ok(super::LeaderProof {
870                owner: crate::checkpoint::LeaderProofOwner {
871                    node_id: proof.node_id,
872                    boot_id: boot,
873                    process_term: proof.process_term,
874                },
875                fencing_token: proof.fencing_token,
876            })
877        })
878        .transpose()
879}
880
881#[cfg(feature = "cluster")]
882fn leader_proof_to_wire(proof: Option<&super::LeaderProof>) -> Option<barrier_v1::LeaderProof> {
883    proof.map(|proof| barrier_v1::LeaderProof {
884        node_id: proof.owner.node_id,
885        boot_id: proof.owner.boot_id.as_bytes().to_vec(),
886        process_term: proof.owner.process_term,
887        fencing_token: proof.fencing_token,
888    })
889}
890
891#[cfg(feature = "cluster")]
892fn assignment_fence_from_wire(
893    version: u64,
894    vnode_count: u32,
895    map_digest: Vec<u8>,
896    participants: Vec<barrier_v1::CheckpointParticipant>,
897) -> Result<Option<super::CheckpointAssignmentFence>, tonic::Status> {
898    if version == 0 && vnode_count == 0 && map_digest.is_empty() && participants.is_empty() {
899        return Ok(None);
900    }
901
902    let assignment_digest: [u8; 32] = map_digest.try_into().map_err(|_| {
903        tonic::Status::invalid_argument(
904            "Checkpoint assignment map digest must contain exactly 32 bytes",
905        )
906    })?;
907    let participants = participants
908        .into_iter()
909        .map(|participant| {
910            let boot_incarnation =
911                uuid::Uuid::from_slice(&participant.boot_incarnation).map_err(|_| {
912                    tonic::Status::invalid_argument(
913                        "Checkpoint participant incarnation must contain exactly 16 bytes",
914                    )
915                })?;
916            Ok(super::CheckpointParticipant {
917                node_id: participant.node_id,
918                boot_incarnation,
919            })
920        })
921        .collect::<Result<Vec<_>, tonic::Status>>()?;
922    let fence = super::CheckpointAssignmentFence {
923        assignment_version: version,
924        partitioning_abi_version: crate::state::PARTITIONING_ABI_VERSION,
925        vnode_count,
926        assignment_digest,
927        participants,
928    };
929    if !fence.is_canonical() {
930        return Err(tonic::Status::invalid_argument(
931            "Non-canonical checkpoint assignment certificate",
932        ));
933    }
934    Ok(Some(fence))
935}
936
937#[cfg(feature = "cluster")]
938fn assignment_fence_to_wire(
939    fence: Option<&super::CheckpointAssignmentFence>,
940) -> WireAssignmentFence {
941    fence.map_or_else(WireAssignmentFence::default, |fence| WireAssignmentFence {
942        version: fence.assignment_version,
943        vnode_count: fence.vnode_count,
944        map_digest: fence.assignment_digest.to_vec(),
945        participants: fence
946            .participants
947            .iter()
948            .map(|participant| barrier_v1::CheckpointParticipant {
949                node_id: participant.node_id,
950                boot_incarnation: participant.boot_incarnation.as_bytes().to_vec(),
951            })
952            .collect(),
953    })
954}
955
956#[cfg(feature = "cluster")]
957fn checkpoint_watermark_from_wire(
958    status: i32,
959    active_watermark_ms: Option<i64>,
960) -> Result<CheckpointWatermark, String> {
961    use barrier_v1::CheckpointWatermarkStatus as WireStatus;
962
963    let wire_status = WireStatus::try_from(status)
964        .map_err(|_| format!("unknown checkpoint watermark status {status}"))?;
965    let watermark = match (wire_status, active_watermark_ms) {
966        (WireStatus::CheckpointWatermarkUninitialized, None) => CheckpointWatermark::Uninitialized,
967        (WireStatus::CheckpointWatermarkIdle, None) => CheckpointWatermark::Idle,
968        (WireStatus::CheckpointWatermarkActive, Some(value)) => CheckpointWatermark::Active(value),
969        _ => {
970            return Err(format!(
971                "invalid checkpoint watermark status {status} with active value {active_watermark_ms:?}"
972            ));
973        }
974    };
975    watermark.validate()?;
976    Ok(watermark)
977}
978
979#[cfg(feature = "cluster")]
980fn grpc_ack(ack: BarrierAck) -> barrier_v1::Ack {
981    use barrier_v1::CheckpointWatermarkStatus as WireStatus;
982
983    let (watermark_status, local_watermark_ms) = match ack.watermark {
984        CheckpointWatermark::Uninitialized => {
985            (WireStatus::CheckpointWatermarkUninitialized as i32, None)
986        }
987        CheckpointWatermark::Idle => (WireStatus::CheckpointWatermarkIdle as i32, None),
988        CheckpointWatermark::Active(value) => {
989            (WireStatus::CheckpointWatermarkActive as i32, Some(value))
990        }
991    };
992    barrier_v1::Ack {
993        epoch: ack.epoch,
994        ok: ack.ok,
995        error: ack.error,
996        local_watermark_ms,
997        checkpoint_id: ack.checkpoint_id,
998        assignment_digest: ack
999            .assignment_digest
1000            .map_or_else(Vec::new, |digest| digest.to_vec()),
1001        watermark_status,
1002    }
1003}
1004
1005#[cfg(feature = "cluster")]
1006fn validate_phase_ack(ack: &barrier_v1::Ack, ann: &BarrierAnnouncement) -> Result<(), String> {
1007    if !CheckpointAttempt::new(ack.epoch, ack.checkpoint_id).is_canonical() {
1008        return Err("Barrier phase acknowledgement has a non-canonical checkpoint ID".into());
1009    }
1010    let expected_digest = ann
1011        .assignment_fence
1012        .as_ref()
1013        .map(super::CheckpointAssignmentFence::digest);
1014    if ack.epoch != ann.epoch
1015        || ack.checkpoint_id != ann.checkpoint_id
1016        || ack.assignment_digest.as_slice()
1017            != expected_digest
1018                .as_ref()
1019                .map_or(&[][..], <[u8; 32]>::as_slice)
1020    {
1021        return Err("Barrier phase acknowledgement identity mismatch".into());
1022    }
1023    if !ack.ok {
1024        return Err(ack
1025            .error
1026            .clone()
1027            .unwrap_or_else(|| "Barrier phase was rejected by follower".into()));
1028    }
1029    Ok(())
1030}
1031
1032#[cfg(feature = "cluster")]
1033impl GrpcBarrierServer {
1034    fn require_live_process_lease(&self) -> Result<Option<&super::LeaseDeadline>, tonic::Status> {
1035        if self.local_process.get().is_none() {
1036            return Ok(None);
1037        }
1038        let deadline = self.process_lease_deadline.get().ok_or_else(|| {
1039            tonic::Status::failed_precondition("Process lease deadline is not installed")
1040        })?;
1041        if !deadline.is_live() {
1042            return Err(tonic::Status::failed_precondition(
1043                "Process lease deadline has expired",
1044            ));
1045        }
1046        Ok(Some(deadline))
1047    }
1048
1049    fn require_exact_local_proof_process(
1050        &self,
1051        proof: &super::LeaderProof,
1052    ) -> Result<(), tonic::Status> {
1053        let local = self.local_process.get().copied().ok_or_else(|| {
1054            tonic::Status::failed_precondition(
1055                "Leader proof confirmation requires a process-bound endpoint",
1056            )
1057        })?;
1058        if local.node_id != proof.owner.node_id
1059            || local.boot_incarnation != proof.owner.boot_id
1060            || local.process_term != proof.owner.process_term
1061        {
1062            return Err(tonic::Status::failed_precondition(
1063                "Leader proof challenge does not target this exact process",
1064            ));
1065        }
1066        Ok(())
1067    }
1068
1069    async fn enqueue_while_process_live(
1070        &self,
1071        announcement: BarrierAnnouncement,
1072    ) -> Result<(), tonic::Status> {
1073        let deadline = self.require_live_process_lease()?;
1074        let send_result = if let Some(deadline) = deadline {
1075            tokio::select! {
1076                biased;
1077                () = deadline.wait_until_expired() => {
1078                    return Err(tonic::Status::failed_precondition(
1079                        "Process lease deadline expired before barrier delivery",
1080                    ));
1081                }
1082                result = self.incoming_tx.send(announcement) => result,
1083            }
1084        } else {
1085            self.incoming_tx.send(announcement).await
1086        };
1087        if send_result.is_err() {
1088            return Err(tonic::Status::aborted("Follower coordinator shutdown"));
1089        }
1090        self.require_live_process_lease()?;
1091        Ok(())
1092    }
1093
1094    async fn wait_for_prepare_ack(
1095        &self,
1096        rx: tokio::sync::oneshot::Receiver<BarrierAck>,
1097    ) -> Result<BarrierAck, tonic::Status> {
1098        let deadline = self.require_live_process_lease()?;
1099        let wait_for_ack = async {
1100            match tokio::time::timeout(PREPARE_RPC_TIMEOUT, rx).await {
1101                Ok(Ok(ack)) => Ok(ack),
1102                Ok(Err(_)) => Err(tonic::Status::internal("Ack sender dropped")),
1103                Err(_) => Err(tonic::Status::deadline_exceeded(
1104                    "Follower checkpoint prepare timed out",
1105                )),
1106            }
1107        };
1108        tokio::pin!(wait_for_ack);
1109        if let Some(deadline) = deadline {
1110            tokio::select! {
1111                biased;
1112                () = deadline.wait_until_expired() => Err(tonic::Status::failed_precondition(
1113                    "Process lease deadline expired while awaiting Prepare acknowledgement",
1114                )),
1115                result = &mut wait_for_ack => result,
1116            }
1117        } else {
1118            wait_for_ack.await
1119        }
1120    }
1121
1122    fn require_local_assignment_process(
1123        &self,
1124        fence: Option<&super::CheckpointAssignmentFence>,
1125    ) -> Result<(), tonic::Status> {
1126        match (self.local_process.get().copied(), fence) {
1127            (None, None) => Ok(()),
1128            (Some(_), None) => Err(tonic::Status::failed_precondition(
1129                "Assignment-less barrier cannot target a process-bound cluster endpoint",
1130            )),
1131            (None, Some(_)) => Err(tonic::Status::failed_precondition(
1132                "Certified barrier cannot target a process-unbound control endpoint",
1133            )),
1134            (Some(local), Some(fence))
1135                if fence.participant_incarnation(local.node_id) == Some(local.boot_incarnation) =>
1136            {
1137                Ok(())
1138            }
1139            (Some(_), Some(_)) => Err(tonic::Status::failed_precondition(
1140                "Certified barrier does not target this exact process",
1141            )),
1142        }
1143    }
1144
1145    async fn require_latest_proof(&self, proof: &super::LeaderProof) -> Result<(), tonic::Status> {
1146        let store = self.leader_lease_store.lock().clone().ok_or_else(|| {
1147            tonic::Status::failed_precondition("Durable leader lease store is not installed")
1148        })?;
1149        let lease = store
1150            .load()
1151            .await
1152            .map_err(|error| {
1153                tonic::Status::unavailable(format!("Leader lease read failed: {error}"))
1154            })?
1155            .ok_or_else(|| tonic::Status::permission_denied("No durable leader lease exists"))?;
1156        if !lease.matches_proof(proof) {
1157            return Err(tonic::Status::permission_denied(
1158                "Leader proof does not match the latest durable lease",
1159            ));
1160        }
1161        Ok(())
1162    }
1163
1164    async fn validate_reversible_leader(
1165        &self,
1166        proof: Option<barrier_v1::LeaderProof>,
1167    ) -> Result<super::LeaderProof, tonic::Status> {
1168        let proof = leader_proof_from_wire(proof)?
1169            .ok_or_else(|| tonic::Status::permission_denied("Missing durable leader proof"))?;
1170        self.require_latest_proof(&proof).await?;
1171        Ok(proof)
1172    }
1173}
1174
1175#[cfg(feature = "cluster")]
1176#[tonic::async_trait]
1177impl barrier_v1::barrier_sync_server::BarrierSync for GrpcBarrierServer {
1178    async fn confirm_leader_proof(
1179        &self,
1180        request: tonic::Request<barrier_v1::LeaderProofChallenge>,
1181    ) -> Result<tonic::Response<barrier_v1::LeaderProofAck>, tonic::Status> {
1182        self.require_live_process_lease()?;
1183        let requested = request.into_inner();
1184        let expected = leader_proof_from_wire(requested.expected_proof)?
1185            .ok_or_else(|| tonic::Status::invalid_argument("Leader proof challenge is missing"))?;
1186        leader_proof_challenge_from_wire(&requested.challenge_id)?;
1187        self.require_exact_local_proof_process(&expected)?;
1188        let provider = self.local_leader_proof.lock().clone().ok_or_else(|| {
1189            tonic::Status::failed_precondition("Local leader proof provider is not installed")
1190        })?;
1191        let local = provider().ok_or_else(|| {
1192            tonic::Status::failed_precondition("No process-local leader grant is live")
1193        })?;
1194        if local != expected {
1195            return Err(tonic::Status::failed_precondition(
1196                "Live process-local leader grant does not match the durable proof challenge",
1197            ));
1198        }
1199        self.require_live_process_lease()?;
1200        Ok(tonic::Response::new(barrier_v1::LeaderProofAck {
1201            challenge_id: requested.challenge_id,
1202        }))
1203    }
1204
1205    async fn prepare(
1206        &self,
1207        request: tonic::Request<barrier_v1::PrepareRequest>,
1208    ) -> Result<tonic::Response<barrier_v1::Ack>, tonic::Status> {
1209        self.require_live_process_lease()?;
1210        let req = request.into_inner();
1211        validate_wire_checkpoint_attempt(req.epoch, req.checkpoint_id)?;
1212        let assignment_fence = assignment_fence_from_wire(
1213            req.assignment_version,
1214            req.assignment_vnode_count,
1215            req.assignment_map_digest,
1216            req.assignment_participants,
1217        )?;
1218        self.require_local_assignment_process(assignment_fence.as_ref())?;
1219        let leader_proof = self
1220            .validate_reversible_leader(req.leader_proof.clone())
1221            .await?;
1222        let attempt = CheckpointAttempt::new(req.epoch, req.checkpoint_id);
1223        let assignment_digest = assignment_fence
1224            .as_ref()
1225            .map(super::CheckpointAssignmentFence::digest);
1226        let identity = BarrierIdentity {
1227            attempt,
1228            assignment_digest,
1229        };
1230
1231        self.require_live_process_lease()?;
1232        let (waiter_id, rx) = {
1233            let mut state = self.prepare_acks.lock();
1234            state.record_receipt(identity);
1235            if let Some(ack) = state.completed.get(&identity) {
1236                self.require_live_process_lease()?;
1237                return Ok(tonic::Response::new(grpc_ack(ack.clone())));
1238            }
1239            if state.pending.get(&identity).map_or(0, Vec::len) >= MAX_PREPARE_WAITERS_PER_IDENTITY
1240            {
1241                return Err(tonic::Status::resource_exhausted(
1242                    "Too many concurrent retries for one checkpoint Prepare",
1243                ));
1244            }
1245            if !state.pending.contains_key(&identity)
1246                && state.pending.len() >= MAX_RETAINED_BARRIER_IDENTITIES
1247            {
1248                return Err(tonic::Status::resource_exhausted(
1249                    "Too many concurrent checkpoint Prepare identities",
1250                ));
1251            }
1252
1253            let (tx, rx) = tokio::sync::oneshot::channel::<BarrierAck>();
1254            let waiter_id = state.next_waiter_id();
1255            state
1256                .pending
1257                .entry(identity)
1258                .or_default()
1259                .push(PendingPrepareWaiter {
1260                    id: waiter_id,
1261                    response: tx,
1262                });
1263            (waiter_id, rx)
1264        };
1265        let _registration = PrepareWaiterRegistration {
1266            state: Arc::clone(&self.prepare_acks),
1267            identity,
1268            waiter_id,
1269        };
1270
1271        let ann = BarrierAnnouncement {
1272            epoch: req.epoch,
1273            checkpoint_id: req.checkpoint_id,
1274            assignment_fence,
1275            leader_proof: Some(leader_proof.clone()),
1276            phase: Phase::Prepare,
1277            flags: req.flags,
1278        };
1279
1280        self.enqueue_while_process_live(ann).await?;
1281
1282        let ack = self.wait_for_prepare_ack(rx).await?;
1283        if ack.epoch != attempt.epoch
1284            || ack.checkpoint_id != attempt.checkpoint_id
1285            || ack.assignment_digest != assignment_digest
1286        {
1287            return Err(tonic::Status::failed_precondition(
1288                "Follower acknowledgement identity mismatch",
1289            ));
1290        }
1291        self.require_latest_proof(&leader_proof).await?;
1292        self.require_live_process_lease()?;
1293        Ok(tonic::Response::new(grpc_ack(ack)))
1294    }
1295
1296    async fn aligned(
1297        &self,
1298        request: tonic::Request<barrier_v1::AlignedRequest>,
1299    ) -> Result<tonic::Response<barrier_v1::Ack>, tonic::Status> {
1300        self.require_live_process_lease()?;
1301        let req = request.into_inner();
1302        validate_wire_checkpoint_attempt(req.epoch, req.checkpoint_id)?;
1303        let assignment_fence = assignment_fence_from_wire(
1304            req.assignment_version,
1305            req.assignment_vnode_count,
1306            req.assignment_map_digest,
1307            req.assignment_participants,
1308        )?;
1309        self.require_local_assignment_process(assignment_fence.as_ref())?;
1310        let leader_proof = self
1311            .validate_reversible_leader(req.leader_proof.clone())
1312            .await?;
1313
1314        // Unlike Commit/Abort, Aligned is mid-protocol: the epoch's ack
1315        // bookkeeping stays untouched — only the announcement is relayed
1316        // so the pipeline's resume gate can release.
1317        let ann = BarrierAnnouncement {
1318            epoch: req.epoch,
1319            checkpoint_id: req.checkpoint_id,
1320            assignment_fence: assignment_fence.clone(),
1321            leader_proof: Some(leader_proof),
1322            phase: Phase::Aligned,
1323            flags: req.flags,
1324        };
1325        self.enqueue_while_process_live(ann).await?;
1326        self.require_live_process_lease()?;
1327        Ok(tonic::Response::new(barrier_v1::Ack {
1328            epoch: req.epoch,
1329            ok: true,
1330            error: None,
1331            local_watermark_ms: None,
1332            checkpoint_id: req.checkpoint_id,
1333            assignment_digest: assignment_fence
1334                .as_ref()
1335                .map_or_else(Vec::new, |fence| fence.digest().to_vec()),
1336            watermark_status: 0,
1337        }))
1338    }
1339
1340    async fn commit(
1341        &self,
1342        request: tonic::Request<barrier_v1::CommitRequest>,
1343    ) -> Result<tonic::Response<barrier_v1::Ack>, tonic::Status> {
1344        self.require_live_process_lease()?;
1345        let req = request.into_inner();
1346        validate_wire_checkpoint_attempt(req.epoch, req.checkpoint_id)?;
1347        let leader_proof = leader_proof_from_wire(req.leader_proof.clone())?;
1348        let assignment_fence = assignment_fence_from_wire(
1349            req.assignment_version,
1350            req.assignment_vnode_count,
1351            req.assignment_map_digest,
1352            req.assignment_participants,
1353        )?;
1354        self.require_local_assignment_process(assignment_fence.as_ref())?;
1355
1356        let ann = BarrierAnnouncement {
1357            epoch: req.epoch,
1358            checkpoint_id: req.checkpoint_id,
1359            assignment_fence: assignment_fence.clone(),
1360            leader_proof,
1361            phase: Phase::Commit,
1362            flags: req.flags,
1363        };
1364        self.enqueue_while_process_live(ann).await?;
1365        self.require_live_process_lease()?;
1366        Ok(tonic::Response::new(barrier_v1::Ack {
1367            epoch: req.epoch,
1368            ok: true,
1369            error: None,
1370            local_watermark_ms: None,
1371            checkpoint_id: req.checkpoint_id,
1372            assignment_digest: assignment_fence
1373                .as_ref()
1374                .map_or_else(Vec::new, |fence| fence.digest().to_vec()),
1375            watermark_status: 0,
1376        }))
1377    }
1378
1379    async fn abort(
1380        &self,
1381        request: tonic::Request<barrier_v1::AbortRequest>,
1382    ) -> Result<tonic::Response<barrier_v1::Ack>, tonic::Status> {
1383        self.require_live_process_lease()?;
1384        let req = request.into_inner();
1385        validate_wire_checkpoint_attempt(req.epoch, req.checkpoint_id)?;
1386        let leader_proof = leader_proof_from_wire(req.leader_proof.clone())?;
1387        let assignment_fence = assignment_fence_from_wire(
1388            req.assignment_version,
1389            req.assignment_vnode_count,
1390            req.assignment_map_digest,
1391            req.assignment_participants,
1392        )?;
1393        self.require_local_assignment_process(assignment_fence.as_ref())?;
1394
1395        let ann = BarrierAnnouncement {
1396            epoch: req.epoch,
1397            checkpoint_id: req.checkpoint_id,
1398            assignment_fence: assignment_fence.clone(),
1399            leader_proof,
1400            phase: Phase::Abort,
1401            flags: req.flags,
1402        };
1403        self.enqueue_while_process_live(ann).await?;
1404        self.require_live_process_lease()?;
1405        Ok(tonic::Response::new(barrier_v1::Ack {
1406            epoch: req.epoch,
1407            ok: true,
1408            error: None,
1409            local_watermark_ms: None,
1410            checkpoint_id: req.checkpoint_id,
1411            assignment_digest: assignment_fence
1412                .as_ref()
1413                .map_or_else(Vec::new, |fence| fence.digest().to_vec()),
1414            watermark_status: 0,
1415        }))
1416    }
1417}
1418
1419#[cfg(feature = "cluster")]
1420async fn get_barrier_client(
1421    peer: NodeId,
1422    expected_process: Option<ExpectedBarrierProcess>,
1423    pool: &BarrierClientPool,
1424    kv: &Arc<dyn ClusterKv>,
1425) -> Result<
1426    Option<barrier_v1::barrier_sync_client::BarrierSyncClient<tonic::transport::Channel>>,
1427    BarrierClientResolutionError,
1428> {
1429    {
1430        let pool = pool.lock();
1431        if let Some(entry) = pool.get(&peer) {
1432            if barrier_client_process_matches(expected_process, entry.process) {
1433                return Ok(Some(entry.client.clone()));
1434            }
1435        }
1436    }
1437
1438    let Some(raw_endpoint) = kv.read_from(peer, BARRIER_ADDR_KEY).await else {
1439        return Ok(None);
1440    };
1441    let (address, published_process) =
1442        decode_barrier_endpoint(&raw_endpoint).map_err(BarrierClientResolutionError::Invalid)?;
1443    if let Some(expected) = expected_process {
1444        let actual = published_process.ok_or(BarrierClientResolutionError::ProcessMismatch)?;
1445        if actual.node_id != peer.0 {
1446            return Err(BarrierClientResolutionError::Invalid(format!(
1447                "cluster control endpoint slot {} advertises node {}",
1448                peer.0, actual.node_id
1449            )));
1450        }
1451        if !expected.matches(actual) {
1452            return Err(BarrierClientResolutionError::ProcessMismatch);
1453        }
1454    } else if let Some(process) = published_process {
1455        if process.node_id != peer.0 {
1456            return Err(BarrierClientResolutionError::Invalid(format!(
1457                "cluster control endpoint slot {} advertises a different node",
1458                peer.0
1459            )));
1460        }
1461        return Err(BarrierClientResolutionError::ProcessMismatch);
1462    }
1463    let endpoint =
1464        super::tls::client_endpoint(&address).map_err(BarrierClientResolutionError::Invalid)?;
1465    let channel = endpoint.connect_lazy();
1466    let client = barrier_v1::barrier_sync_client::BarrierSyncClient::new(channel);
1467
1468    let mut pool = pool.lock();
1469    if let Some(entry) = pool.get(&peer) {
1470        if barrier_client_process_matches(expected_process, entry.process) {
1471            return Ok(Some(entry.client.clone()));
1472        }
1473        if expected_process.is_none() && entry.process.is_some() {
1474            return Ok(Some(client));
1475        }
1476    }
1477    while pool.len() >= crate::checkpoint::MAX_CHECKPOINT_PARTICIPANTS && !pool.contains_key(&peer)
1478    {
1479        let victim = if expected_process.is_none() {
1480            pool.iter()
1481                .find_map(|(node, entry)| entry.process.is_none().then_some(*node))
1482        } else {
1483            pool.keys().next().copied()
1484        };
1485        let Some(victim) = victim else {
1486            break;
1487        };
1488        pool.remove(&victim);
1489    }
1490    if pool.len() >= crate::checkpoint::MAX_CHECKPOINT_PARTICIPANTS && !pool.contains_key(&peer) {
1491        return Ok(Some(client));
1492    }
1493    pool.insert(
1494        peer,
1495        BarrierClientEntry {
1496            process: published_process,
1497            client: client.clone(),
1498        },
1499    );
1500    Ok(Some(client))
1501}
1502
1503#[cfg(feature = "cluster")]
1504fn evict_barrier_client(
1505    pool: &BarrierClientPool,
1506    peer: NodeId,
1507    expected_process: Option<ExpectedBarrierProcess>,
1508) {
1509    let mut pool = pool.lock();
1510    let remove = pool
1511        .get(&peer)
1512        .is_some_and(|entry| barrier_client_process_matches(expected_process, entry.process));
1513    if remove {
1514        pool.remove(&peer);
1515    }
1516}
1517
1518#[cfg(feature = "cluster")]
1519async fn call_phase_rpc(
1520    client: &mut barrier_v1::barrier_sync_client::BarrierSyncClient<tonic::transport::Channel>,
1521    ann: &BarrierAnnouncement,
1522    request_timeout: Duration,
1523) -> Result<(), String> {
1524    let assignment = assignment_fence_to_wire(ann.assignment_fence.as_ref());
1525    match ann.phase {
1526        Phase::Aligned => {
1527            let mut request = tonic::Request::new(barrier_v1::AlignedRequest {
1528                epoch: ann.epoch,
1529                checkpoint_id: ann.checkpoint_id,
1530                flags: ann.flags,
1531                assignment_version: assignment.version,
1532                assignment_participants: assignment.participants,
1533                assignment_vnode_count: assignment.vnode_count,
1534                assignment_map_digest: assignment.map_digest,
1535                leader_proof: leader_proof_to_wire(ann.leader_proof.as_ref()),
1536            });
1537            request.set_timeout(request_timeout);
1538            client
1539                .aligned(request)
1540                .await
1541                .map_err(|error| error.to_string())
1542                .and_then(|response| validate_phase_ack(&response.into_inner(), ann))
1543        }
1544        Phase::Commit => {
1545            let mut request = tonic::Request::new(barrier_v1::CommitRequest {
1546                epoch: ann.epoch,
1547                checkpoint_id: ann.checkpoint_id,
1548                flags: ann.flags,
1549                assignment_version: assignment.version,
1550                assignment_participants: assignment.participants,
1551                assignment_vnode_count: assignment.vnode_count,
1552                assignment_map_digest: assignment.map_digest,
1553                leader_proof: leader_proof_to_wire(ann.leader_proof.as_ref()),
1554            });
1555            request.set_timeout(request_timeout);
1556            client
1557                .commit(request)
1558                .await
1559                .map_err(|error| error.to_string())
1560                .and_then(|response| validate_phase_ack(&response.into_inner(), ann))
1561        }
1562        Phase::Abort => {
1563            let mut request = tonic::Request::new(barrier_v1::AbortRequest {
1564                epoch: ann.epoch,
1565                checkpoint_id: ann.checkpoint_id,
1566                flags: ann.flags,
1567                assignment_version: assignment.version,
1568                assignment_participants: assignment.participants,
1569                assignment_vnode_count: assignment.vnode_count,
1570                assignment_map_digest: assignment.map_digest,
1571                leader_proof: leader_proof_to_wire(ann.leader_proof.as_ref()),
1572            });
1573            request.set_timeout(request_timeout);
1574            client
1575                .abort(request)
1576                .await
1577                .map_err(|error| error.to_string())
1578                .and_then(|response| validate_phase_ack(&response.into_inner(), ann))
1579        }
1580        Phase::Prepare => Err("Prepare cannot use the phase-notification RPC path".into()),
1581    }
1582}
1583
1584/// Fan a non-Prepare phase announcement to one peer over gRPC. A failed
1585/// RPC evicts the pooled client so the next round re-resolves the peer.
1586#[cfg(feature = "cluster")]
1587async fn send_phase_rpc(
1588    peer: NodeId,
1589    clients_pool: BarrierClientPool,
1590    kv: Arc<dyn ClusterKv>,
1591    ann: BarrierAnnouncement,
1592    deadline: tokio::time::Instant,
1593) -> Result<(), String> {
1594    let rpc = match ann.phase {
1595        Phase::Aligned => "aligned",
1596        Phase::Commit => "commit",
1597        Phase::Abort => "abort",
1598        Phase::Prepare => "prepare",
1599    };
1600
1601    let expected_process = ann
1602        .assignment_fence
1603        .as_ref()
1604        .map(|fence| {
1605            fence
1606                .participant_incarnation(peer.0)
1607                .map(|boot| ExpectedBarrierProcess::participant(peer.0, boot))
1608                .ok_or_else(|| {
1609                    format!("{rpc} RPC peer {} is outside the assignment roster", peer.0)
1610                })
1611        })
1612        .transpose()?;
1613    let result = tokio::time::timeout_at(deadline, async {
1614        let mut client = get_barrier_client(peer, expected_process, &clients_pool, &kv)
1615            .await
1616            .map_err(|error| format!("failed to resolve peer {}: {error}", peer.0))?
1617            .ok_or_else(|| format!("failed to get client for peer {}", peer.0))?;
1618        let request_timeout = deadline.saturating_duration_since(tokio::time::Instant::now());
1619        call_phase_rpc(&mut client, &ann, request_timeout).await
1620    })
1621    .await;
1622
1623    match result {
1624        Ok(Ok(())) => Ok(()),
1625        Ok(Err(error)) => {
1626            evict_barrier_client(&clients_pool, peer, expected_process);
1627            if tokio::time::Instant::now() >= deadline || error.contains("Timeout expired") {
1628                Err(format!(
1629                    "{rpc} RPC to peer {} exceeded its request deadline",
1630                    peer.0
1631                ))
1632            } else {
1633                Err(format!("{rpc} RPC to peer {} failed: {error}", peer.0))
1634            }
1635        }
1636        Err(_) => {
1637            evict_barrier_client(&clients_pool, peer, expected_process);
1638            Err(format!(
1639                "{rpc} RPC to peer {} exceeded its request deadline",
1640                peer.0
1641            ))
1642        }
1643    }
1644}
1645
1646/// Deliver a non-Prepare phase to remote participants over the low-latency notification path.
1647#[cfg(feature = "cluster")]
1648async fn send_phase_notifications(
1649    state: &GrpcState,
1650    kv: &Arc<dyn ClusterKv>,
1651    ann: &BarrierAnnouncement,
1652    expected: Vec<NodeId>,
1653) -> Vec<Result<(), String>> {
1654    let deadline = tokio::time::Instant::now() + PHASE_RPC_TIMEOUT;
1655    futures::future::join_all(expected.into_iter().map(|peer| {
1656        send_phase_rpc(
1657            peer,
1658            Arc::clone(&state.clients),
1659            Arc::clone(kv),
1660            ann.clone(),
1661            deadline,
1662        )
1663    }))
1664    .await
1665}
1666
1667#[cfg(feature = "cluster")]
1668fn send_local_phase_notification(
1669    state: &GrpcState,
1670    ann: &BarrierAnnouncement,
1671    process_lease: &super::LeaseDeadline,
1672) -> Result<(), String> {
1673    if !process_lease.is_live() {
1674        return Err("local process lease expired before barrier delivery".into());
1675    }
1676    match state.incoming_tx.try_send(ann.clone()) {
1677        Ok(()) => {}
1678        Err(crossfire::TrySendError::Full(_)) => {
1679            return Err("local barrier announcement relay is full".into());
1680        }
1681        Err(crossfire::TrySendError::Disconnected(_)) => {
1682            return Err("local barrier announcement relay is closed".into());
1683        }
1684    }
1685    if !process_lease.is_live() {
1686        return Err("local process lease expired during barrier delivery".into());
1687    }
1688    Ok(())
1689}
1690
1691/// Typed prepare-failure classification for the quorum wait:
1692/// `Unreachable` counts toward `TimedOut{missing}` (the peer cannot
1693/// participate), `Nack` toward `Failed` (a live follower answered
1694/// `ok = false`).
1695#[cfg(feature = "cluster")]
1696enum PeerFailure {
1697    Unreachable,
1698    Nack(String),
1699}
1700
1701/// Exact announcement and participant roster bound to one eager Prepare round.
1702#[cfg(feature = "cluster")]
1703struct PrepareFanoutBatch {
1704    announcement: BarrierAnnouncement,
1705    expected: Vec<NodeId>,
1706    // `JoinSet` aborts all remaining tasks when the batch or quorum future is dropped.
1707    tasks: tokio::task::JoinSet<Result<(NodeId, CheckpointWatermark), (NodeId, PeerFailure)>>,
1708}
1709
1710#[cfg(feature = "cluster")]
1711#[derive(Clone, Copy)]
1712struct PrepareFanoutBudget {
1713    total: Duration,
1714    per_attempt: Duration,
1715}
1716
1717#[cfg(feature = "cluster")]
1718fn prepare_fanout_budget(quorum_window: Duration) -> Result<PrepareFanoutBudget, String> {
1719    if quorum_window.is_zero() {
1720        return Err("Prepare quorum window must be greater than zero".into());
1721    }
1722    let per_attempt = quorum_window / 2;
1723    if per_attempt.is_zero() {
1724        return Err("Prepare quorum window is too small to divide into retry attempts".into());
1725    }
1726    Ok(PrepareFanoutBudget {
1727        total: PREPARE_RPC_TIMEOUT.max(quorum_window),
1728        per_attempt,
1729    })
1730}
1731
1732#[cfg(feature = "cluster")]
1733enum PrepareFanoutState {
1734    Pending(PrepareFanoutBatch),
1735    Claimed(BarrierAnnouncement),
1736    QuorumReached(BarrierAnnouncement),
1737}
1738
1739#[cfg(feature = "cluster")]
1740impl PrepareFanoutState {
1741    const fn announcement(&self) -> &BarrierAnnouncement {
1742        match self {
1743            Self::Pending(batch) => &batch.announcement,
1744            Self::Claimed(announcement) | Self::QuorumReached(announcement) => announcement,
1745        }
1746    }
1747}
1748
1749#[cfg(feature = "cluster")]
1750fn clustered_prepare_roster(prepare: &BarrierAnnouncement) -> Result<Option<Vec<NodeId>>, String> {
1751    if prepare.phase != Phase::Prepare {
1752        return Err("Prepare fan-out received a different barrier phase".into());
1753    }
1754    let Some(fence) = prepare.assignment_fence.as_ref() else {
1755        // A build with the cluster feature may still run embedded or single-node. Those modes
1756        // retain the existing KV/direct-on-wait behavior and do not install a clustered batch.
1757        return Ok(None);
1758    };
1759    if !fence.is_canonical() {
1760        return Err("clustered Prepare has a non-canonical assignment certificate".into());
1761    }
1762    let proof = prepare
1763        .leader_proof
1764        .as_ref()
1765        .filter(|proof| proof.is_canonical())
1766        .ok_or_else(|| "clustered Prepare has no canonical leader proof".to_string())?;
1767    if fence.participant_incarnation(proof.owner.node_id) != Some(proof.owner.boot_id) {
1768        return Err(
1769            "clustered Prepare leader proof is outside its exact assignment process roster".into(),
1770        );
1771    }
1772
1773    Ok(Some(
1774        fence
1775            .participants
1776            .iter()
1777            .filter(|participant| participant.node_id != proof.owner.node_id)
1778            .map(|participant| NodeId(participant.node_id))
1779            .collect(),
1780    ))
1781}
1782
1783#[cfg(feature = "cluster")]
1784fn clustered_phase_roster(
1785    announcement: &BarrierAnnouncement,
1786    local_process: Option<BarrierProcessIdentity>,
1787) -> Result<Option<Vec<NodeId>>, String> {
1788    if announcement.phase == Phase::Prepare {
1789        return Err("non-Prepare fan-out received a Prepare announcement".into());
1790    }
1791    let Some(fence) = announcement.assignment_fence.as_ref() else {
1792        return Ok(None);
1793    };
1794    if !fence.is_canonical() {
1795        return Err("clustered barrier has a non-canonical assignment certificate".into());
1796    }
1797
1798    let local_process = local_process
1799        .ok_or_else(|| "assignment-certified phase has no local process identity".to_string())?;
1800    if announcement.phase == Phase::Aligned {
1801        let proof = announcement
1802            .leader_proof
1803            .as_ref()
1804            .filter(|proof| proof.is_canonical())
1805            .ok_or_else(|| "clustered Aligned has no canonical leader proof".to_string())?;
1806        if fence.participant_incarnation(proof.owner.node_id) != Some(proof.owner.boot_id) {
1807            return Err(
1808                "clustered Aligned leader proof is outside its exact assignment process roster"
1809                    .into(),
1810            );
1811        }
1812        if local_process.node_id != proof.owner.node_id
1813            || local_process.boot_incarnation != proof.owner.boot_id
1814            || local_process.process_term != proof.owner.process_term
1815        {
1816            return Err("clustered Aligned sender does not own its live leader proof".into());
1817        }
1818    }
1819
1820    Ok(Some(
1821        fence
1822            .participants
1823            .iter()
1824            .filter(|participant| {
1825                if is_terminal_phase(announcement.phase) {
1826                    participant.node_id != local_process.node_id
1827                } else {
1828                    !local_process.matches_participant(participant)
1829                }
1830            })
1831            .map(|participant| NodeId(participant.node_id))
1832            .collect(),
1833    ))
1834}
1835
1836#[cfg(feature = "cluster")]
1837fn prepare_fanout_plan(
1838    announcement: &BarrierAnnouncement,
1839    quorum_window: Option<Duration>,
1840) -> Result<(Option<Vec<NodeId>>, Option<PrepareFanoutBudget>), String> {
1841    if announcement.phase != Phase::Prepare {
1842        return Ok((None, None));
1843    }
1844    let roster = clustered_prepare_roster(announcement)?;
1845    let budget = roster
1846        .as_ref()
1847        .map(|_| {
1848            quorum_window.ok_or_else(|| {
1849                "assignment-certified Prepare has no quorum retry window".to_string()
1850            })
1851        })
1852        .transpose()?
1853        .map(prepare_fanout_budget)
1854        .transpose()?;
1855    Ok((roster, budget))
1856}
1857
1858#[cfg(feature = "cluster")]
1859fn canonical_expected_roster(expected: &[NodeId]) -> Result<Vec<NodeId>, String> {
1860    let mut canonical = expected.to_vec();
1861    canonical.sort_unstable_by_key(|peer| peer.0);
1862    if canonical.iter().any(NodeId::is_unassigned)
1863        || canonical.windows(2).any(|pair| pair[0] == pair[1])
1864    {
1865        return Err("Prepare quorum participant roster is not canonical".into());
1866    }
1867    Ok(canonical)
1868}
1869
1870#[cfg(feature = "cluster")]
1871fn install_prepare_fanout(
1872    state: &GrpcState,
1873    kv: &Arc<dyn ClusterKv>,
1874    prepare: &BarrierAnnouncement,
1875    expected: Vec<NodeId>,
1876    budget: PrepareFanoutBudget,
1877) {
1878    let mut pending = state.prepare_fanout.lock();
1879    let rpc_deadline = tokio::time::Instant::now() + budget.total;
1880    let mut tasks = tokio::task::JoinSet::new();
1881    for &peer in &expected {
1882        let clients_pool = Arc::clone(&state.clients);
1883        let kv = Arc::clone(kv);
1884        let prepare = prepare.clone();
1885        tasks.spawn(async move {
1886            prepare_peer_until_deadline(
1887                peer,
1888                clients_pool,
1889                kv,
1890                prepare,
1891                rpc_deadline,
1892                budget.per_attempt,
1893            )
1894            .await
1895        });
1896    }
1897    let incoming = PrepareFanoutBatch {
1898        announcement: prepare.clone(),
1899        expected,
1900        tasks,
1901    };
1902    *pending = Some(PrepareFanoutState::Pending(incoming));
1903}
1904
1905#[cfg(feature = "cluster")]
1906fn preflight_prepare_fanout(
1907    state: &GrpcState,
1908    prepare: &BarrierAnnouncement,
1909) -> Result<bool, String> {
1910    let mut pending = state.prepare_fanout.lock();
1911    let Some(current) = pending.as_ref() else {
1912        return Ok(true);
1913    };
1914    validate_announcement_attempt(prepare)?;
1915    validate_announcement_attempt(current.announcement())?;
1916    match prepare
1917        .checkpoint_id
1918        .cmp(&current.announcement().checkpoint_id)
1919    {
1920        std::cmp::Ordering::Greater => {
1921            // Admission has already rejected attempt regression. Cancel the obsolete structured
1922            // fan-out before cancellable durable I/O so it cannot complete after being superseded.
1923            pending.take();
1924            Ok(true)
1925        }
1926        std::cmp::Ordering::Equal if current.announcement() == prepare => match current {
1927            PrepareFanoutState::Pending(_) => Ok(false),
1928            PrepareFanoutState::Claimed(_) => {
1929                Err("Prepare cannot be republished while its quorum is being collected".into())
1930            }
1931            PrepareFanoutState::QuorumReached(_) => {
1932                Err("Prepare cannot regress an exact quorum-ready checkpoint".into())
1933            }
1934        },
1935        std::cmp::Ordering::Less => {
1936            Err("stale Prepare cannot replace a newer in-flight fan-out".into())
1937        }
1938        std::cmp::Ordering::Equal => {
1939            Err("conflicting Prepare certificate cannot replace the in-flight fan-out".into())
1940        }
1941    }
1942}
1943
1944#[cfg(feature = "cluster")]
1945fn mark_prepare_quorum_reached(
1946    state: &GrpcState,
1947    prepare: &BarrierAnnouncement,
1948) -> Result<(), String> {
1949    let mut fanout = state.prepare_fanout.lock();
1950    match fanout.take() {
1951        Some(PrepareFanoutState::Claimed(claimed)) if claimed == *prepare => {
1952            *fanout = Some(PrepareFanoutState::QuorumReached(claimed));
1953            Ok(())
1954        }
1955        Some(other) => {
1956            *fanout = Some(other);
1957            Err("Prepare quorum completion lost its exact claimed fan-out".into())
1958        }
1959        None => Err("Prepare quorum completion has no claimed fan-out".into()),
1960    }
1961}
1962
1963#[cfg(feature = "cluster")]
1964fn require_aligned_quorum(state: &GrpcState, aligned: &BarrierAnnouncement) -> Result<(), String> {
1965    let fanout = state.prepare_fanout.lock();
1966    let Some(PrepareFanoutState::QuorumReached(prepare)) = fanout.as_ref() else {
1967        return Err("clustered Aligned requires a successful exact Prepare quorum".into());
1968    };
1969    if !same_announcement_identity(prepare, aligned) {
1970        return Err("clustered Aligned does not match the exact reached Prepare quorum".into());
1971    }
1972    Ok(())
1973}
1974
1975#[cfg(feature = "cluster")]
1976fn retire_prepare_fanout(state: &GrpcState) {
1977    state.prepare_fanout.lock().take();
1978}
1979
1980#[cfg(feature = "cluster")]
1981fn retryable_prepare_status(status: &tonic::Status) -> bool {
1982    matches!(
1983        status.code(),
1984        // Tonic maps a client-service readiness failure to `Unknown`; it is
1985        // transport state, not a response from the follower. The remaining
1986        // codes represent a connection that cannot currently complete the
1987        // request. Fence and validation failures use distinct semantic codes.
1988        tonic::Code::Unknown
1989            | tonic::Code::Unavailable
1990            | tonic::Code::DeadlineExceeded
1991            | tonic::Code::Cancelled
1992            | tonic::Code::Aborted
1993    )
1994}
1995
1996#[cfg(feature = "cluster")]
1997async fn wait_for_prepare_retry(deadline: tokio::time::Instant, backoff: &mut Duration) -> bool {
1998    let now = tokio::time::Instant::now();
1999    if now >= deadline {
2000        return false;
2001    }
2002
2003    tokio::time::sleep_until((now + *backoff).min(deadline)).await;
2004    *backoff = backoff.saturating_mul(2).min(PREPARE_RETRY_MAX_BACKOFF);
2005    tokio::time::Instant::now() < deadline
2006}
2007
2008#[cfg(feature = "cluster")]
2009fn prepare_rpc_request(
2010    prepare: &BarrierAnnouncement,
2011    assignment: &WireAssignmentFence,
2012    leader_proof: Option<&barrier_v1::LeaderProof>,
2013    timeout: Duration,
2014) -> tonic::Request<barrier_v1::PrepareRequest> {
2015    let mut request = tonic::Request::new(barrier_v1::PrepareRequest {
2016        epoch: prepare.epoch,
2017        checkpoint_id: prepare.checkpoint_id,
2018        flags: prepare.flags,
2019        assignment_version: assignment.version,
2020        assignment_participants: assignment.participants.clone(),
2021        assignment_vnode_count: assignment.vnode_count,
2022        assignment_map_digest: assignment.map_digest.clone(),
2023        leader_proof: leader_proof.cloned(),
2024    });
2025    request.set_timeout(timeout);
2026    request
2027}
2028
2029#[cfg(feature = "cluster")]
2030fn validate_prepare_ack(
2031    peer: NodeId,
2032    prepare: &BarrierAnnouncement,
2033    assignment_digest: Option<&[u8; 32]>,
2034    ack: barrier_v1::Ack,
2035) -> Result<(NodeId, CheckpointWatermark), (NodeId, PeerFailure)> {
2036    if ack.epoch != prepare.epoch
2037        || ack.checkpoint_id != prepare.checkpoint_id
2038        || ack.assignment_digest.as_slice()
2039            != assignment_digest.map_or(&[][..], <[u8; 32]>::as_slice)
2040    {
2041        return Err((
2042            peer,
2043            PeerFailure::Nack("Prepare acknowledgement identity mismatch".into()),
2044        ));
2045    }
2046    if !ack.ok {
2047        return Err((
2048            peer,
2049            PeerFailure::Nack(
2050                ack.error
2051                    .unwrap_or_else(|| "Unknown prepare failure".to_string()),
2052            ),
2053        ));
2054    }
2055    checkpoint_watermark_from_wire(ack.watermark_status, ack.local_watermark_ms)
2056        .map(|watermark| (peer, watermark))
2057        .map_err(|error| (peer, PeerFailure::Nack(error)))
2058}
2059
2060#[cfg(feature = "cluster")]
2061fn prepare_expected_process(
2062    peer: NodeId,
2063    prepare: &BarrierAnnouncement,
2064) -> Result<Option<ExpectedBarrierProcess>, (NodeId, PeerFailure)> {
2065    let Some(fence) = prepare.assignment_fence.as_ref() else {
2066        return Ok(None);
2067    };
2068    let boot = fence.participant_incarnation(peer.0).ok_or_else(|| {
2069        (
2070            peer,
2071            PeerFailure::Nack("Prepare peer is outside the assignment process roster".into()),
2072        )
2073    })?;
2074    Ok(Some(ExpectedBarrierProcess::participant(peer.0, boot)))
2075}
2076
2077#[cfg(feature = "cluster")]
2078async fn prepare_peer_until_deadline(
2079    peer: NodeId,
2080    clients_pool: BarrierClientPool,
2081    kv: Arc<dyn ClusterKv>,
2082    prepare: BarrierAnnouncement,
2083    deadline: tokio::time::Instant,
2084    max_attempt_duration: Duration,
2085) -> Result<(NodeId, CheckpointWatermark), (NodeId, PeerFailure)> {
2086    let assignment = assignment_fence_to_wire(prepare.assignment_fence.as_ref());
2087    let assignment_digest = prepare
2088        .assignment_fence
2089        .as_ref()
2090        .map(super::CheckpointAssignmentFence::digest);
2091    let leader_proof = leader_proof_to_wire(prepare.leader_proof.as_ref());
2092    let expected_process = prepare_expected_process(peer, &prepare)?;
2093    let mut backoff = PREPARE_RETRY_INITIAL_BACKOFF;
2094
2095    loop {
2096        if tokio::time::Instant::now() >= deadline {
2097            return Err((peer, PeerFailure::Unreachable));
2098        }
2099
2100        let Ok(client) = tokio::time::timeout_at(
2101            deadline,
2102            get_barrier_client(peer, expected_process, &clients_pool, &kv),
2103        )
2104        .await
2105        else {
2106            return Err((peer, PeerFailure::Unreachable));
2107        };
2108        let client = match client {
2109            Ok(client) => client,
2110            Err(BarrierClientResolutionError::ProcessMismatch) => {
2111                if wait_for_prepare_retry(deadline, &mut backoff).await {
2112                    continue;
2113                }
2114                return Err((peer, PeerFailure::Unreachable));
2115            }
2116            Err(BarrierClientResolutionError::Invalid(error)) => {
2117                return Err((peer, PeerFailure::Nack(error)));
2118            }
2119        };
2120        let Some(mut client) = client else {
2121            if wait_for_prepare_retry(deadline, &mut backoff).await {
2122                continue;
2123            }
2124            return Err((peer, PeerFailure::Unreachable));
2125        };
2126
2127        let now = tokio::time::Instant::now();
2128        let remaining = deadline.saturating_duration_since(now);
2129        if remaining.is_zero() {
2130            evict_barrier_client(&clients_pool, peer, expected_process);
2131            return Err((peer, PeerFailure::Unreachable));
2132        }
2133        // Leave part of the existing quorum budget available to evict and re-resolve a lazy
2134        // channel whose readiness check stalls. Prepare identities are idempotent, so a live
2135        // follower may safely complete an earlier attempt and serve its cached acknowledgement.
2136        let attempt_budget = (remaining / 2).min(max_attempt_duration);
2137        if attempt_budget.is_zero() {
2138            evict_barrier_client(&clients_pool, peer, expected_process);
2139            return Err((peer, PeerFailure::Unreachable));
2140        }
2141        let attempt_deadline = now + attempt_budget;
2142
2143        let request =
2144            prepare_rpc_request(&prepare, &assignment, leader_proof.as_ref(), attempt_budget);
2145
2146        match tokio::time::timeout_at(attempt_deadline, client.prepare(request)).await {
2147            Ok(Ok(response)) => {
2148                return validate_prepare_ack(
2149                    peer,
2150                    &prepare,
2151                    assignment_digest.as_ref(),
2152                    response.into_inner(),
2153                );
2154            }
2155            Ok(Err(status)) => {
2156                evict_barrier_client(&clients_pool, peer, expected_process);
2157                if retryable_prepare_status(&status) {
2158                    if wait_for_prepare_retry(deadline, &mut backoff).await {
2159                        continue;
2160                    }
2161                    return Err((peer, PeerFailure::Unreachable));
2162                }
2163                return Err((peer, PeerFailure::Nack(status.to_string())));
2164            }
2165            Err(_) => {
2166                evict_barrier_client(&clients_pool, peer, expected_process);
2167                if wait_for_prepare_retry(deadline, &mut backoff).await {
2168                    continue;
2169                }
2170                return Err((peer, PeerFailure::Unreachable));
2171            }
2172        }
2173    }
2174}
2175
2176#[derive(Default)]
2177struct AnnouncementPublicationState {
2178    initialized: bool,
2179    latest: Option<BarrierAnnouncement>,
2180}
2181
2182/// Cross-instance barrier coordination.
2183pub struct BarrierCoordinator {
2184    kv: Arc<dyn ClusterKv>,
2185    /// Serializes local publication across every runtime mode. The latest admitted value advances
2186    /// before cancellable I/O so an ambiguous write result cannot reopen an older phase.
2187    publication: tokio::sync::Mutex<AnnouncementPublicationState>,
2188    #[cfg(feature = "cluster")]
2189    grpc: Arc<parking_lot::Mutex<Option<Arc<GrpcState>>>>,
2190    #[cfg(feature = "cluster")]
2191    leader_election: Arc<parking_lot::Mutex<ActiveLeaderState>>,
2192    #[cfg(feature = "cluster")]
2193    leader_lease_store: Arc<parking_lot::Mutex<Option<Arc<super::LeaderLeaseStore>>>>,
2194    #[cfg(feature = "cluster")]
2195    local_leader_proof: Arc<parking_lot::Mutex<Option<LocalLeaderProofProvider>>>,
2196    #[cfg(feature = "cluster")]
2197    local_process: Arc<std::sync::OnceLock<BarrierProcessIdentity>>,
2198    #[cfg(feature = "cluster")]
2199    unbound_endpoint_started: parking_lot::Mutex<bool>,
2200    #[cfg(feature = "cluster")]
2201    process_lease_deadline: Arc<std::sync::OnceLock<Arc<super::LeaseDeadline>>>,
2202}
2203
2204impl std::fmt::Debug for BarrierCoordinator {
2205    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2206        f.debug_struct("BarrierCoordinator").finish_non_exhaustive()
2207    }
2208}
2209
2210impl Drop for BarrierCoordinator {
2211    fn drop(&mut self) {
2212        #[cfg(feature = "cluster")]
2213        {
2214            let grpc_opt = self.grpc.lock().take();
2215            if let Some(state) = grpc_opt {
2216                abort_grpc_tasks(&state);
2217            }
2218        }
2219    }
2220}
2221
2222impl BarrierCoordinator {
2223    /// Wrap a KV implementation.
2224    #[must_use]
2225    pub fn new(kv: Arc<dyn ClusterKv>) -> Self {
2226        Self {
2227            kv,
2228            publication: tokio::sync::Mutex::new(AnnouncementPublicationState::default()),
2229            #[cfg(feature = "cluster")]
2230            grpc: Arc::new(parking_lot::Mutex::new(None)),
2231            #[cfg(feature = "cluster")]
2232            leader_election: Arc::new(parking_lot::Mutex::new(None)),
2233            #[cfg(feature = "cluster")]
2234            leader_lease_store: Arc::new(parking_lot::Mutex::new(None)),
2235            #[cfg(feature = "cluster")]
2236            local_leader_proof: Arc::new(parking_lot::Mutex::new(None)),
2237            #[cfg(feature = "cluster")]
2238            local_process: Arc::new(std::sync::OnceLock::new()),
2239            #[cfg(feature = "cluster")]
2240            unbound_endpoint_started: parking_lot::Mutex::new(false),
2241            #[cfg(feature = "cluster")]
2242            process_lease_deadline: Arc::new(std::sync::OnceLock::new()),
2243        }
2244    }
2245
2246    #[cfg(feature = "cluster")]
2247    fn require_live_bound_process_lease(&self) -> Result<(), String> {
2248        if self.local_process.get().is_none() {
2249            return Ok(());
2250        }
2251        let deadline = self
2252            .process_lease_deadline
2253            .get()
2254            .ok_or_else(|| "process lease deadline is not installed".to_string())?;
2255        if !deadline.is_live() {
2256            return Err("process lease deadline has expired".into());
2257        }
2258        Ok(())
2259    }
2260
2261    #[cfg(feature = "cluster")]
2262    fn claim_endpoint_process(&self) -> Option<BarrierProcessIdentity> {
2263        let mut unbound_endpoint_started = self.unbound_endpoint_started.lock();
2264        let process = self.local_process.get().copied();
2265        if process.is_none() {
2266            *unbound_endpoint_started = true;
2267        }
2268        process
2269    }
2270
2271    #[cfg(feature = "cluster")]
2272    pub(crate) fn install_process_lease_deadline(
2273        &self,
2274        deadline: Arc<super::LeaseDeadline>,
2275    ) -> Result<(), String> {
2276        match self.process_lease_deadline.set(deadline) {
2277            Ok(()) => Ok(()),
2278            Err(deadline)
2279                if self
2280                    .process_lease_deadline
2281                    .get()
2282                    .is_some_and(|current| Arc::ptr_eq(current, &deadline)) =>
2283            {
2284                Ok(())
2285            }
2286            Err(_) => Err("process lease deadline is already installed".into()),
2287        }
2288    }
2289
2290    #[cfg(feature = "cluster")]
2291    pub(crate) fn install_local_process_lease(
2292        &self,
2293        lease: &super::ProcessLease,
2294    ) -> Result<(), String> {
2295        let process = BarrierProcessIdentity::from_process_lease(lease)?;
2296        let unbound_endpoint_started = self.unbound_endpoint_started.lock();
2297        if *unbound_endpoint_started {
2298            return Err(
2299                "an assignment-less cluster control endpoint cannot be promoted in place".into(),
2300            );
2301        }
2302        match self.local_process.set(process) {
2303            Ok(()) => Ok(()),
2304            Err(_) if self.local_process.get() == Some(&process) => Ok(()),
2305            Err(_) => Err("cluster control endpoint process identity is already installed".into()),
2306        }
2307    }
2308
2309    #[cfg(feature = "cluster")]
2310    pub(crate) fn set_local_leader_proof_provider(&self, provider: LocalLeaderProofProvider) {
2311        *self.local_leader_proof.lock() = Some(provider);
2312    }
2313
2314    /// Install the durable authority used to validate clustered reversible barrier phases.
2315    /// Without it, clustered `Prepare` and `Aligned` traffic fails closed.
2316    #[cfg(feature = "cluster")]
2317    pub fn set_leader_lease_store(&self, store: Arc<super::LeaderLeaseStore>) {
2318        *self.leader_lease_store.lock() = Some(store);
2319    }
2320
2321    /// Exact durable authority installed for clustered barriers and checkpoint decisions.
2322    ///
2323    /// Embedded and single-node runtimes do not call this path. A cluster runtime that omitted
2324    /// authority wiring fails closed instead of falling back to standalone outcome objects.
2325    ///
2326    /// # Errors
2327    /// Returns `NotConfigured` when durable cluster checkpoint authority is not installed.
2328    #[cfg(feature = "cluster")]
2329    pub fn checkpoint_authority(
2330        &self,
2331    ) -> Result<Arc<super::LeaderLeaseStore>, super::ClusterCheckpointAuthorityError> {
2332        self.leader_lease_store
2333            .lock()
2334            .clone()
2335            .ok_or(super::ClusterCheckpointAuthorityError::NotConfigured)
2336    }
2337
2338    #[cfg(feature = "cluster")]
2339    async fn validate_reversible_announcement(
2340        &self,
2341        ann: &BarrierAnnouncement,
2342    ) -> Result<(), String> {
2343        validate_announcement_attempt(ann)?;
2344        if !matches!(ann.phase, Phase::Prepare | Phase::Aligned) {
2345            return Ok(());
2346        }
2347        // An assignment certificate is the barrier layer's cluster-runtime marker. Embedded and
2348        // single-node coordinators can be built with the cluster feature enabled, but do not have
2349        // a remote leader lease and must retain their local KV path.
2350        if ann.assignment_fence.is_none() {
2351            return Ok(());
2352        }
2353        let proof = ann.leader_proof.as_ref().ok_or_else(|| {
2354            format!(
2355                "clustered {:?} for checkpoint {}/{} is missing a durable leader proof",
2356                ann.phase, ann.epoch, ann.checkpoint_id
2357            )
2358        })?;
2359        self.validate_announcement_leader(ann, proof).await
2360    }
2361
2362    #[cfg(feature = "cluster")]
2363    fn require_exact_local_leader_proof(&self, ann: &BarrierAnnouncement) -> Result<(), String> {
2364        let expected = ann
2365            .leader_proof
2366            .as_ref()
2367            .ok_or_else(|| "assignment-certified barrier has no exact leader proof".to_string())?;
2368        let provider = self
2369            .local_leader_proof
2370            .lock()
2371            .clone()
2372            .ok_or_else(|| "local leader proof provider is not installed".to_string())?;
2373        if provider().as_ref() != Some(expected) {
2374            return Err(
2375                "assignment-certified barrier sender no longer owns its exact leader proof".into(),
2376            );
2377        }
2378        Ok(())
2379    }
2380
2381    #[cfg(feature = "cluster")]
2382    async fn validate_announcement_leader(
2383        &self,
2384        ann: &BarrierAnnouncement,
2385        proof: &crate::checkpoint::LeaderProof,
2386    ) -> Result<(), String> {
2387        let local_proof = self
2388            .local_leader_proof
2389            .lock()
2390            .clone()
2391            .and_then(|provider| provider());
2392        if local_proof.as_ref() == Some(proof) {
2393            return Ok(());
2394        }
2395        let store = self
2396            .leader_lease_store
2397            .lock()
2398            .clone()
2399            .ok_or_else(|| "durable leader lease store is not installed".to_string())?;
2400        let lease = store
2401            .load()
2402            .await
2403            .map_err(|error| format!("leader lease read failed: {error}"))?
2404            .ok_or_else(|| "no durable leader lease exists".to_string())?;
2405        if !lease.matches_proof(proof) {
2406            return Err(format!(
2407                "clustered {:?} for checkpoint {}/{} does not match the latest durable leader lease",
2408                ann.phase, ann.epoch, ann.checkpoint_id
2409            ));
2410        }
2411        Ok(())
2412    }
2413
2414    #[cfg(feature = "cluster")]
2415    pub(super) async fn validate_checkpoint_prepare(
2416        &self,
2417        announcement: &BarrierAnnouncement,
2418    ) -> Result<(), String> {
2419        validate_announcement_attempt(announcement)?;
2420        if announcement.phase != Phase::Prepare {
2421            return Err("checkpoint Prepare validation received a different barrier phase".into());
2422        }
2423        let proof = announcement.leader_proof.as_ref().ok_or_else(|| {
2424            format!(
2425                "clustered Prepare for checkpoint {}/{} is missing a durable leader proof",
2426                announcement.epoch, announcement.checkpoint_id
2427            )
2428        })?;
2429        self.validate_announcement_leader(announcement, proof).await
2430    }
2431
2432    /// Configure membership used to target active barrier peers.
2433    /// Gossip election is not a barrier authority boundary.
2434    #[cfg(feature = "cluster")]
2435    pub fn set_leader_election(
2436        &mut self,
2437        instance_id: NodeId,
2438        members_rx: watch::Receiver<Vec<NodeInfo>>,
2439        leader_eligible: Arc<AtomicBool>,
2440    ) {
2441        *self.leader_election.lock() = Some((instance_id, members_rx, leader_eligible));
2442    }
2443
2444    /// Local monotonic receipt time for this exact gRPC Prepare.
2445    #[cfg(feature = "cluster")]
2446    #[must_use]
2447    pub fn prepare_received_at(&self, prepare: &BarrierAnnouncement) -> Option<std::time::Instant> {
2448        if prepare.phase != Phase::Prepare {
2449            return None;
2450        }
2451        let identity = BarrierIdentity::from_announcement(prepare);
2452        self.grpc.lock().as_ref().and_then(|state| {
2453            state
2454                .prepare_acks
2455                .lock()
2456                .received_at
2457                .get(&identity)
2458                .copied()
2459        })
2460    }
2461
2462    /// Bind and run the follower's direct gRPC barrier sync server.
2463    ///
2464    /// # Errors
2465    /// Returns an error string on bind or socket address retrieval failures.
2466    #[cfg(feature = "cluster")]
2467    pub async fn start_server(
2468        &self,
2469        bind_addr: std::net::SocketAddr,
2470        advertise_host: Option<String>,
2471    ) -> Result<std::net::SocketAddr, String> {
2472        use barrier_v1::barrier_sync_server::BarrierSyncServer;
2473        use std::net::TcpListener;
2474        use tonic::transport::Server;
2475
2476        self.require_live_bound_process_lease()?;
2477        let advertised_process = self.claim_endpoint_process();
2478        let listener = TcpListener::bind(bind_addr).map_err(|e| e.to_string())?;
2479        let local_addr = listener.local_addr().map_err(|e| e.to_string())?;
2480        listener.set_nonblocking(true).map_err(|e| e.to_string())?;
2481        let tokio_listener =
2482            tokio::net::TcpListener::from_std(listener).map_err(|e| e.to_string())?;
2483        let advertise_addr = if let Some(ref host) = advertise_host {
2484            format!("{host}:{}", local_addr.port())
2485        } else if local_addr.ip().is_unspecified() {
2486            let hostname = gethostname::gethostname();
2487            let hostname = hostname.to_string_lossy();
2488            if hostname.is_empty() {
2489                local_addr.to_string()
2490            } else {
2491                format!("{hostname}:{}", local_addr.port())
2492            }
2493        } else {
2494            local_addr.to_string()
2495        };
2496        let advertisement = encode_barrier_endpoint(&advertise_addr, advertised_process)?;
2497
2498        let (incoming_tx, incoming_rx) = crossfire::mpsc::bounded_async::<BarrierAnnouncement>(128);
2499        let prepare_acks = Arc::new(parking_lot::Mutex::new(PrepareAckState::default()));
2500        let clients = Arc::new(parking_lot::Mutex::new(FxHashMap::default()));
2501        let local_process = Arc::clone(&self.local_process);
2502
2503        let server_impl = GrpcBarrierServer {
2504            incoming_tx: incoming_tx.clone(),
2505            prepare_acks: Arc::clone(&prepare_acks),
2506            leader_lease_store: Arc::clone(&self.leader_lease_store),
2507            local_leader_proof: Arc::clone(&self.local_leader_proof),
2508            local_process: Arc::clone(&local_process),
2509            process_lease_deadline: Arc::clone(&self.process_lease_deadline),
2510        };
2511
2512        // Apply TLS synchronously so a bad cert fails start_server (before
2513        // publishing BARRIER_ADDR_KEY) rather than silently never serving.
2514        let mut builder = Server::builder();
2515        if let Some(tls) = super::tls::server_tls() {
2516            builder = builder
2517                .tls_config(tls.clone())
2518                .map_err(|e| format!("cluster control-plane TLS config: {e}"))?;
2519        }
2520        let router = builder.add_service(BarrierSyncServer::new(server_impl));
2521        let server_task = tokio::spawn(async move {
2522            let incoming_stream = tokio_stream::wrappers::TcpListenerStream::new(tokio_listener);
2523            let _ = router.serve_with_incoming(incoming_stream).await;
2524        });
2525
2526        // Relay every gRPC-delivered announcement into a relation-validated
2527        // watch in arrival order. Observation is then non-destructive,
2528        // so the pipeline's resume gate and the background durable
2529        // tail can watch concurrently (matching the gossip-KV
2530        // fallback's read-latest semantics).
2531        let (latest_tx, latest_rx) = watch::channel::<Option<BarrierAnnouncement>>(None);
2532        let merge_error = Arc::new(parking_lot::Mutex::new(None));
2533        let relay_merge_error = Arc::clone(&merge_error);
2534        let relay_task = tokio::spawn(async move {
2535            while let Ok(ann) = incoming_rx.recv().await {
2536                let merged = match latest_tx.borrow().clone() {
2537                    Some(current) => merge_direct_announcement(current, ann),
2538                    None => Ok(ann),
2539                };
2540                match merged {
2541                    Ok(merged) => {
2542                        let changed = latest_tx.borrow().as_ref() != Some(&merged);
2543                        if changed {
2544                            let _ = latest_tx.send(Some(merged));
2545                        }
2546                    }
2547                    Err(error) => {
2548                        tracing::error!(%error, "rejecting conflicting direct barrier history");
2549                        let mut retained = relay_merge_error.lock();
2550                        if retained.is_none() {
2551                            *retained = Some(error);
2552                        }
2553                    }
2554                }
2555            }
2556        });
2557
2558        let grpc_state = Arc::new(GrpcState {
2559            latest_rx,
2560            incoming_tx,
2561            merge_error,
2562            prepare_acks,
2563            prepare_fanout: parking_lot::Mutex::new(None),
2564            clients,
2565            server_handle: Arc::new(parking_lot::Mutex::new(Some(server_task))),
2566            relay_handle: Arc::new(parking_lot::Mutex::new(Some(relay_task))),
2567            advertise_addr: advertise_addr.clone(),
2568            local_process: Arc::clone(&local_process),
2569        });
2570
2571        if let Err(error) = self.require_live_bound_process_lease() {
2572            abort_grpc_tasks(&grpc_state);
2573            return Err(error);
2574        }
2575        if let Err(error) = self.kv.write_checked(BARRIER_ADDR_KEY, advertisement).await {
2576            abort_grpc_tasks(&grpc_state);
2577            return Err(format!(
2578                "publish cluster control endpoint advertisement: {error}"
2579            ));
2580        }
2581
2582        *self.grpc.lock() = Some(grpc_state);
2583
2584        Ok(local_addr)
2585    }
2586
2587    /// Ask one exact remote process to confirm a proof already read from durable authority.
2588    ///
2589    /// The response echoes only a fresh challenge id. It never returns a process-local or durable
2590    /// fencing token.
2591    ///
2592    /// # Errors
2593    /// Fails when the proof, peer address, RPC, acknowledgement, or deadline is invalid.
2594    #[cfg(feature = "cluster")]
2595    pub async fn confirm_remote_leader_proof(
2596        &self,
2597        proof: &super::LeaderProof,
2598        deadline: tokio::time::Instant,
2599    ) -> Result<bool, String> {
2600        if !proof.is_canonical() {
2601            return Err("remote leader proof challenge is not canonical".into());
2602        }
2603        let peer = NodeId(proof.owner.node_id);
2604        let state = self
2605            .grpc
2606            .lock()
2607            .clone()
2608            .ok_or_else(|| "cluster control RPC server is not started".to_string())?;
2609        let clients = Arc::clone(&state.clients);
2610        let request_timeout = deadline.saturating_duration_since(tokio::time::Instant::now());
2611        let challenge = uuid::Uuid::new_v4();
2612        let expected_process = Some(ExpectedBarrierProcess::exact(&proof.owner));
2613        let result = tokio::time::timeout_at(deadline, async {
2614            let mut client =
2615                match get_barrier_client(peer, expected_process, &clients, &self.kv).await {
2616                    Ok(Some(client)) => client,
2617                    Ok(None) => {
2618                        return Err(format!(
2619                            "cluster control address for peer {} is unavailable",
2620                            peer.0
2621                        ));
2622                    }
2623                    Err(BarrierClientResolutionError::ProcessMismatch) => return Ok(false),
2624                    Err(BarrierClientResolutionError::Invalid(error)) => return Err(error),
2625                };
2626            let mut request = tonic::Request::new(barrier_v1::LeaderProofChallenge {
2627                expected_proof: leader_proof_to_wire(Some(proof)),
2628                challenge_id: challenge.as_bytes().to_vec(),
2629            });
2630            request.set_timeout(request_timeout);
2631            match client.confirm_leader_proof(request).await {
2632                Ok(response) => {
2633                    let acknowledged = response.into_inner().challenge_id;
2634                    if !leader_proof_ack_matches(challenge, &acknowledged) {
2635                        return Err("remote leader proof acknowledgement challenge mismatch".into());
2636                    }
2637                    Ok(true)
2638                }
2639                Err(status) if status.code() == tonic::Code::FailedPrecondition => {
2640                    // The stable node id may now advertise a replacement process. Do not pin
2641                    // subsequent proof attempts to a still-responsive channel for the old boot.
2642                    evict_barrier_client(&clients, peer, expected_process);
2643                    Ok(false)
2644                }
2645                Err(status) => Err(status.to_string()),
2646            }
2647        })
2648        .await;
2649        match result {
2650            Ok(Ok(confirmed)) => Ok(confirmed),
2651            Ok(Err(error)) => {
2652                evict_barrier_client(&clients, peer, expected_process);
2653                Err(error)
2654            }
2655            Err(_) => {
2656                evict_barrier_client(&clients, peer, expected_process);
2657                Err(format!(
2658                    "remote leader proof request for peer {} timed out",
2659                    peer.0
2660                ))
2661            }
2662        }
2663    }
2664
2665    /// Leader-side announcement for terminal, aligned, and assignment-less local/KV phases.
2666    ///
2667    /// # Errors
2668    /// Assignment-certified Prepare must use [`Self::announce_prepare`]. Assignment-certified
2669    /// reversible phases require a started leased barrier server, and Aligned requires the exact
2670    /// Prepare to have completed [`Self::wait_for_quorum`]. Other errors propagate validation,
2671    /// encoding, and publication failures.
2672    pub async fn announce(&self, ann: &BarrierAnnouncement) -> Result<(), String> {
2673        #[cfg(feature = "cluster")]
2674        if ann.phase == Phase::Prepare && ann.assignment_fence.is_some() {
2675            return Err(
2676                "assignment-certified Prepare requires an explicit quorum retry window".into(),
2677            );
2678        }
2679        self.announce_inner(ann, None).await
2680    }
2681
2682    /// Durably publish one assignment-certified Prepare and immediately start its direct fan-out.
2683    ///
2684    /// # Errors
2685    /// Rejects a different phase, an assignment-less announcement, a zero/indivisible quorum
2686    /// window, malformed authority, conflicting in-flight Prepare state, or publication failure.
2687    #[cfg(feature = "cluster")]
2688    pub async fn announce_prepare(
2689        &self,
2690        ann: &BarrierAnnouncement,
2691        quorum_window: Duration,
2692    ) -> Result<(), String> {
2693        if ann.phase != Phase::Prepare || ann.assignment_fence.is_none() {
2694            return Err("explicit Prepare fan-out requires an assignment certificate".into());
2695        }
2696        self.announce_inner(ann, Some(quorum_window)).await
2697    }
2698
2699    #[cfg(feature = "cluster")]
2700    async fn discover_assignment_less_phase_peers(&self, local_address: &str) -> Vec<NodeId> {
2701        let live: Option<FxHashSet<NodeId>> =
2702            self.leader_election
2703                .lock()
2704                .clone()
2705                .map(|(_, members_rx, _)| {
2706                    members_rx
2707                        .borrow()
2708                        .iter()
2709                        .filter(|member| {
2710                            matches!(member.state, NodeState::Active | NodeState::Draining)
2711                        })
2712                        .map(|member| member.id)
2713                        .collect()
2714                });
2715        let mut discovered = Vec::new();
2716        for (node_id, raw_endpoint) in self.kv.scan(BARRIER_ADDR_KEY).await {
2717            let address = match decode_barrier_endpoint(&raw_endpoint) {
2718                Ok((address, None)) => address,
2719                Ok((_, Some(_))) => continue,
2720                Err(error) => {
2721                    tracing::warn!(peer = node_id.0, %error, "ignoring invalid cluster control endpoint");
2722                    continue;
2723                }
2724            };
2725            if address == local_address {
2726                continue;
2727            }
2728            if live.as_ref().is_some_and(|live| !live.contains(&node_id)) {
2729                continue;
2730            }
2731            discovered.push(node_id);
2732        }
2733        discovered
2734    }
2735
2736    async fn announce_inner(
2737        &self,
2738        ann: &BarrierAnnouncement,
2739        prepare_quorum_window: Option<Duration>,
2740    ) -> Result<(), String> {
2741        validate_announcement_attempt(ann)?;
2742        #[cfg(feature = "cluster")]
2743        {
2744            self.validate_reversible_announcement(ann).await?;
2745            let (prepare_roster, prepare_budget) = prepare_fanout_plan(ann, prepare_quorum_window)?;
2746            let grpc_opt = self.grpc.lock().clone();
2747            let process_bound = self.local_process.get().is_some();
2748            match (process_bound, ann.assignment_fence.is_some()) {
2749                (true, false) => {
2750                    return Err(format!(
2751                        "process-bound cluster {:?} requires an assignment certificate",
2752                        ann.phase
2753                    ));
2754                }
2755                (false, true) => {
2756                    return Err(format!(
2757                        "assignment-certified {:?} requires a process-bound leased endpoint",
2758                        ann.phase
2759                    ));
2760                }
2761                (true, true) | (false, false) => {}
2762            }
2763            if ann.assignment_fence.is_some()
2764                && matches!(ann.phase, Phase::Prepare | Phase::Aligned)
2765                && grpc_opt.is_none()
2766            {
2767                return Err(format!(
2768                    "assignment-certified {:?} requires a started leased barrier server",
2769                    ann.phase
2770                ));
2771            }
2772            let phase_roster = if let Some(state) = grpc_opt.as_ref() {
2773                let local_process = state.local_process.get().copied();
2774                if ann.phase == Phase::Prepare {
2775                    None
2776                } else {
2777                    clustered_phase_roster(ann, local_process)?
2778                }
2779            } else {
2780                None
2781            };
2782            let json = serde_json::to_string(ann).map_err(|e| e.to_string())?;
2783            let mut publication = self.publication.lock().await;
2784            if !publication.initialized {
2785                publication.latest = self.scan_latest_announcement().await?;
2786                publication.initialized = true;
2787            }
2788            if let Some(current) = publication.latest.as_ref() {
2789                validate_publication_order(current, ann)?;
2790            }
2791
2792            if let Some(state) = grpc_opt {
2793                let certified_prepare = ann.phase == Phase::Prepare && prepare_roster.is_some();
2794                if certified_prepare {
2795                    self.require_live_bound_process_lease()?;
2796                    self.require_exact_local_leader_proof(ann)?;
2797                }
2798                let replace_prepare_fanout = if prepare_roster.is_some() {
2799                    // Reject stale/equivocating retries before they can overwrite the durable
2800                    // gossip slot. The publication lock keeps this check and durable mutation
2801                    // atomic with respect to every local phase transition.
2802                    preflight_prepare_fanout(&state, ann)?
2803                } else {
2804                    false
2805                };
2806                let certified_aligned = ann.phase == Phase::Aligned && phase_roster.is_some();
2807                if certified_aligned {
2808                    // Recheck authority after lock contention and require the exact Prepare quorum
2809                    // before releasing any participant. Aligned remains reversible: it cannot
2810                    // advance sources, commit sinks, authorize recovery, or collect old state.
2811                    self.require_live_bound_process_lease()?;
2812                    self.require_exact_local_leader_proof(ann)?;
2813                    require_aligned_quorum(&state, ann)?;
2814                    let process_lease = self
2815                        .process_lease_deadline
2816                        .get()
2817                        .cloned()
2818                        .ok_or_else(|| "process lease deadline is not installed".to_string())?;
2819                    // Admission advances before cancellable I/O. An ambiguous durable error must
2820                    // not reopen Prepare or allow a terminal attempt to regress locally.
2821                    publication.latest = Some(ann.clone());
2822
2823                    let expected =
2824                        phase_roster.expect("certified Aligned has an assignment roster");
2825                    let remote = send_phase_notifications(&state, &self.kv, ann, expected);
2826                    let local_result = send_local_phase_notification(&state, ann, &process_lease);
2827                    let durable = self.kv.write_checked(ANNOUNCEMENT_KEY, json);
2828                    tokio::pin!(remote);
2829                    tokio::pin!(durable);
2830                    let mut completed_remote = None;
2831                    let durable_result = tokio::select! {
2832                        result = &mut durable => result,
2833                        results = &mut remote => {
2834                            completed_remote = Some(results);
2835                            durable.await
2836                        }
2837                    };
2838                    let authority_result = self
2839                        .require_live_bound_process_lease()
2840                        .and_then(|()| self.require_exact_local_leader_proof(ann));
2841                    drop(publication);
2842                    authority_result?;
2843
2844                    let direct_results = match completed_remote {
2845                        Some(results) => results,
2846                        None => remote.await,
2847                    };
2848                    if let Err(error) = local_result {
2849                        tracing::warn!(
2850                            epoch = ann.epoch,
2851                            error = %error,
2852                            "local aligned announcement delivery failed; resume falls back to durable observation or Commit"
2853                        );
2854                    }
2855                    for error in direct_results.into_iter().filter_map(Result::err) {
2856                        tracing::warn!(
2857                            epoch = ann.epoch,
2858                            error = %error,
2859                            "aligned announcement delivery failed; resume falls back to durable observation or Commit"
2860                        );
2861                    }
2862                    self.require_live_bound_process_lease()?;
2863                    self.require_exact_local_leader_proof(ann)?;
2864                    durable_result
2865                        .map_err(|error| format!("publish barrier announcement: {error}"))?;
2866                    return Ok(());
2867                }
2868                // Prepare, terminal, and assignment-less phases remain durable-first. Recovery
2869                // and irreversible decisions never depend on best-effort notification delivery.
2870                if is_terminal_phase(ann.phase) {
2871                    // Cancel superseded transport work at admission, before a cancellable durable
2872                    // write. Publication ordering retains the terminal identity if I/O is
2873                    // ambiguous, so the old Prepare must not remain actionable.
2874                    retire_prepare_fanout(&state);
2875                }
2876                publication.latest = Some(ann.clone());
2877                self.kv
2878                    .write_checked(ANNOUNCEMENT_KEY, json)
2879                    .await
2880                    .map_err(|error| format!("publish barrier announcement: {error}"))?;
2881                if ann.phase == Phase::Prepare {
2882                    if certified_prepare {
2883                        self.require_live_bound_process_lease()?;
2884                        self.require_exact_local_leader_proof(ann)?;
2885                    }
2886                    if let Some(expected) = prepare_roster.filter(|_| replace_prepare_fanout) {
2887                        // Start the exact assignment-complete batch only after Prepare is durable.
2888                        // Local source fencing, shuffle alignment, and state capture can now run
2889                        // concurrently with follower capture instead of delaying RPC delivery.
2890                        install_prepare_fanout(
2891                            &state,
2892                            &self.kv,
2893                            ann,
2894                            expected,
2895                            prepare_budget.expect("certified Prepare budget was validated"),
2896                        );
2897                    }
2898                    drop(publication);
2899                } else {
2900                    drop(publication);
2901                    let expected = if let Some(roster) = phase_roster {
2902                        // The checkpoint certificate, not mutable membership, is the phase roster.
2903                        // This also removes a discovery/object-store scan from the clustered hot
2904                        // path and excludes Active processes outside the frozen cut.
2905                        roster
2906                    } else {
2907                        // Feature-enabled embedded and single-node use remains assignment-less.
2908                        // Only that compatibility path discovers direct peers from membership.
2909                        self.discover_assignment_less_phase_peers(&state.advertise_addr)
2910                            .await
2911                    };
2912
2913                    let results = send_phase_notifications(&state, &self.kv, ann, expected).await;
2914                    for res in results {
2915                        match res {
2916                            Ok(()) => {}
2917                            // Aligned is best-effort per peer: a missed
2918                            // delivery only delays that peer's pipeline
2919                            // resume until Commit (or its gate timeout) —
2920                            // never fail the announce, and never skip the
2921                            // KV write below.
2922                            Err(e) if ann.phase == Phase::Aligned => {
2923                                tracing::warn!(
2924                                    epoch = ann.epoch,
2925                                    error = %e,
2926                                    "aligned announcement RPC failed; peer resumes on Commit"
2927                                );
2928                            }
2929                            Err(e) => return Err(e),
2930                        }
2931                    }
2932                }
2933
2934                return Ok(());
2935            }
2936
2937            publication.latest = Some(ann.clone());
2938            self.kv
2939                .write_checked(ANNOUNCEMENT_KEY, json)
2940                .await
2941                .map_err(|error| format!("publish barrier announcement: {error}"))?;
2942            Ok(())
2943        }
2944
2945        #[cfg(not(feature = "cluster"))]
2946        {
2947            let _ = prepare_quorum_window;
2948            let json = serde_json::to_string(ann).map_err(|e| e.to_string())?;
2949            let mut publication = self.publication.lock().await;
2950            if !publication.initialized {
2951                publication.latest = self.scan_latest_announcement().await?;
2952                publication.initialized = true;
2953            }
2954            if let Some(current) = publication.latest.as_ref() {
2955                validate_publication_order(current, ann)?;
2956            }
2957            publication.latest = Some(ann.clone());
2958            self.kv
2959                .write_checked(ANNOUNCEMENT_KEY, json)
2960                .await
2961                .map_err(|error| format!("publish barrier announcement: {error}"))?;
2962            Ok(())
2963        }
2964    }
2965
2966    /// Watch over gRPC-delivered announcements, for push-driven waits
2967    /// (the decision wait and the Aligned resume gate). `None` until
2968    /// the gRPC server is started — gossip-KV-only deployments fall
2969    /// back to polling the merged gossip history.
2970    #[cfg(feature = "cluster")]
2971    #[must_use]
2972    pub fn announcement_watch(&self) -> Option<watch::Receiver<Option<BarrierAnnouncement>>> {
2973        self.grpc.lock().as_ref().map(|s| s.latest_rx.clone())
2974    }
2975
2976    /// Merge the latest direct and gossip announcements without consulting remote authority.
2977    /// Callers may inspect the result, but must validate a matching reversible phase before use.
2978    /// Observation is non-destructive, and direct plus gossip histories must remain related by
2979    /// canonical checkpoint ID. Terminal durable KV values remain the decision authority.
2980    ///
2981    /// # Errors
2982    /// Returns a string on transport, decode, or conflicting-history failure.
2983    pub(super) async fn observe_hint(
2984        &self,
2985        leader: NodeId,
2986    ) -> Result<Option<BarrierAnnouncement>, String> {
2987        #[cfg(feature = "cluster")]
2988        let grpc_latest: Option<BarrierAnnouncement> = {
2989            let grpc_opt = self.grpc.lock().clone();
2990            if let Some(error) = grpc_opt
2991                .as_ref()
2992                .and_then(|state| state.merge_error.lock().clone())
2993            {
2994                return Err(error);
2995            }
2996            grpc_opt.and_then(|state| state.latest_rx.borrow().clone())
2997        };
2998        #[cfg(not(feature = "cluster"))]
2999        let grpc_latest: Option<BarrierAnnouncement> = None;
3000
3001        let kv_latest: Option<BarrierAnnouncement> =
3002            match self.kv.read_from_checked(leader, ANNOUNCEMENT_KEY).await? {
3003                Some(json) => Some(serde_json::from_str(&json).map_err(|error| {
3004                    format!("malformed durable barrier announcement from {leader}: {error}")
3005                })?),
3006                None => None,
3007            };
3008
3009        let observed = match (grpc_latest, kv_latest) {
3010            (Some(g), Some(k)) => Some(merge_observed_announcement(g, k)?),
3011            (Some(g), None) => Some(g),
3012            (None, k) => k,
3013        };
3014        if let Some(announcement) = observed.as_ref() {
3015            validate_announcement_attempt(announcement)?;
3016        }
3017        Ok(observed)
3018    }
3019
3020    #[cfg(feature = "cluster")]
3021    /// Validate one merged announcement immediately before a caller uses it.
3022    pub(super) async fn validate_observed(
3023        &self,
3024        announcement: &BarrierAnnouncement,
3025    ) -> Result<(), String> {
3026        self.validate_reversible_announcement(announcement).await
3027    }
3028
3029    async fn scan_latest_announcement(&self) -> Result<Option<BarrierAnnouncement>, String> {
3030        let mut announcements = Vec::new();
3031        for (node, json) in self.kv.scan_checked(ANNOUNCEMENT_KEY).await? {
3032            let announcement: BarrierAnnouncement = serde_json::from_str(&json)
3033                .map_err(|error| format!("malformed barrier announcement from {node}: {error}"))?;
3034            announcements.push(announcement);
3035        }
3036        validate_scanned_announcements(announcements)
3037    }
3038
3039    /// Follower-side ack.
3040    ///
3041    /// # Errors
3042    /// Returns a string on JSON encode failure.
3043    pub async fn ack(&self, ack: &BarrierAck) -> Result<(), String> {
3044        validate_ack_attempt(ack)?;
3045        #[cfg(feature = "cluster")]
3046        {
3047            let grpc_opt = self.grpc.lock().clone();
3048            if let Some(state) = grpc_opt {
3049                let identity = BarrierIdentity::from_ack(ack);
3050                let (cached, waiters) = {
3051                    let mut prepare = state.prepare_acks.lock();
3052                    let cached = prepare.record_ack(identity, ack);
3053                    let waiters = prepare.pending.remove(&identity).unwrap_or_default();
3054                    (cached, waiters)
3055                };
3056                for waiter in waiters {
3057                    let _ = waiter.response.send(cached.clone());
3058                }
3059                return Ok(());
3060            }
3061        }
3062
3063        let json = serde_json::to_string(ack).map_err(|e| e.to_string())?;
3064        self.kv.write(ACK_KEY, json).await;
3065        Ok(())
3066    }
3067
3068    /// Leader-side: wait until quorum or `deadline`.
3069    #[allow(clippy::too_many_lines)]
3070    // `PeerFailure` (module level, below) classifies each peer's
3071    // prepare outcome.
3072    pub async fn wait_for_quorum(
3073        &self,
3074        prepare: &BarrierAnnouncement,
3075        expected: &[NodeId],
3076        deadline: Duration,
3077    ) -> QuorumOutcome {
3078        if let Err(error) = validate_announcement_attempt(prepare) {
3079            return QuorumOutcome::Failed {
3080                failures: vec![(
3081                    expected.first().copied().unwrap_or(NodeId::UNASSIGNED),
3082                    error,
3083                )],
3084            };
3085        }
3086        let epoch = prepare.epoch;
3087        let checkpoint_id = prepare.checkpoint_id;
3088        let assignment_digest = prepare
3089            .assignment_fence
3090            .as_ref()
3091            .map(super::CheckpointAssignmentFence::digest);
3092        #[cfg(feature = "cluster")]
3093        {
3094            let grpc_opt = self.grpc.lock().clone();
3095            if let Some(state) = grpc_opt {
3096                let expected_roster = match canonical_expected_roster(expected) {
3097                    Ok(roster) => roster,
3098                    Err(error) => {
3099                        return QuorumOutcome::Failed {
3100                            failures: vec![(
3101                                expected.first().copied().unwrap_or(NodeId::UNASSIGNED),
3102                                error,
3103                            )],
3104                        };
3105                    }
3106                };
3107                let eager_batch = if prepare.assignment_fence.is_some() {
3108                    let mut pending = state.prepare_fanout.lock();
3109                    match pending.take() {
3110                        Some(PrepareFanoutState::Pending(batch))
3111                            if batch.announcement == *prepare
3112                                && batch.expected == expected_roster =>
3113                        {
3114                            *pending = Some(PrepareFanoutState::Claimed(prepare.clone()));
3115                            Some(batch)
3116                        }
3117                        Some(state @ PrepareFanoutState::Pending(_))
3118                            if state.announcement() != prepare =>
3119                        {
3120                            let claimed = state.announcement().clone();
3121                            *pending = Some(PrepareFanoutState::Claimed(claimed));
3122                            return QuorumOutcome::Failed {
3123                                failures: vec![(
3124                                    expected_roster
3125                                        .first()
3126                                        .copied()
3127                                        .unwrap_or(NodeId::UNASSIGNED),
3128                                    "Prepare quorum does not match the exact announced fan-out"
3129                                        .into(),
3130                                )],
3131                            };
3132                        }
3133                        Some(state @ PrepareFanoutState::Pending(_)) => {
3134                            let claimed = state.announcement().clone();
3135                            *pending = Some(PrepareFanoutState::Claimed(claimed));
3136                            return QuorumOutcome::Failed {
3137                                failures: vec![(
3138                                    expected_roster
3139                                        .first()
3140                                        .copied()
3141                                        .unwrap_or(NodeId::UNASSIGNED),
3142                                    "Prepare quorum roster does not match the announced assignment"
3143                                        .into(),
3144                                )],
3145                            };
3146                        }
3147                        Some(
3148                            state @ (PrepareFanoutState::Claimed(_)
3149                            | PrepareFanoutState::QuorumReached(_)),
3150                        ) => {
3151                            let exact = state.announcement() == prepare;
3152                            *pending = Some(state);
3153                            return QuorumOutcome::Failed {
3154                                failures: vec![(
3155                                    expected_roster
3156                                        .first()
3157                                        .copied()
3158                                        .unwrap_or(NodeId::UNASSIGNED),
3159                                    if exact {
3160                                        "clustered Prepare fan-out was already claimed or completed"
3161                                    } else {
3162                                        "Prepare quorum does not match the claimed fan-out"
3163                                    }
3164                                    .into(),
3165                                )],
3166                            };
3167                        }
3168                        None => {
3169                            return QuorumOutcome::Failed {
3170                                failures: vec![(
3171                                    expected_roster
3172                                        .first()
3173                                        .copied()
3174                                        .unwrap_or(NodeId::UNASSIGNED),
3175                                    "clustered Prepare has no in-flight announced fan-out".into(),
3176                                )],
3177                            };
3178                        }
3179                    }
3180                } else {
3181                    None
3182                };
3183
3184                let prepare_deadline = tokio::time::Instant::now() + deadline;
3185                let results = if let Some(mut batch) = eager_batch {
3186                    debug_assert_eq!(batch.tasks.len(), expected_roster.len());
3187                    let mut results = Vec::with_capacity(expected_roster.len());
3188                    loop {
3189                        match tokio::time::timeout_at(prepare_deadline, batch.tasks.join_next())
3190                            .await
3191                        {
3192                            Ok(Some(Ok(result))) => results.push(result),
3193                            Ok(Some(Err(error))) => results.push(Err((
3194                                NodeId::UNASSIGNED,
3195                                PeerFailure::Nack(format!("Prepare RPC task failed: {error}")),
3196                            ))),
3197                            Ok(None) | Err(_) => break,
3198                        }
3199                    }
3200                    results
3201                } else {
3202                    // Embedded/single-node use with the cluster feature and legacy direct tests
3203                    // have no assignment certificate. Keep their on-demand direct path intact.
3204                    let futures = expected.iter().map(|&peer| {
3205                        prepare_peer_until_deadline(
3206                            peer,
3207                            Arc::clone(&state.clients),
3208                            Arc::clone(&self.kv),
3209                            prepare.clone(),
3210                            prepare_deadline,
3211                            deadline / 2,
3212                        )
3213                    });
3214                    futures::future::join_all(futures).await
3215                };
3216
3217                let mut successful = Vec::new();
3218                let mut failures = Vec::new();
3219                let mut follower_watermark = None;
3220                let mut timed_out = Vec::new();
3221
3222                for res in results {
3223                    match res {
3224                        Ok((peer, wm)) => {
3225                            successful.push(peer);
3226                            follower_watermark = Some(follower_watermark.map_or(wm, |current| {
3227                                CheckpointWatermark::cluster_min(current, wm)
3228                            }));
3229                        }
3230                        Err((peer, PeerFailure::Unreachable)) => timed_out.push(peer),
3231                        Err((peer, PeerFailure::Nack(msg))) => failures.push((peer, msg)),
3232                    }
3233                }
3234                successful.sort_unstable_by_key(|peer| peer.0);
3235                failures.sort_unstable_by_key(|(peer, _)| peer.0);
3236
3237                let completed: FxHashSet<NodeId> = successful
3238                    .iter()
3239                    .copied()
3240                    .chain(timed_out.iter().copied())
3241                    .chain(failures.iter().map(|(peer, _)| *peer))
3242                    .collect();
3243                for &peer in &expected_roster {
3244                    if !completed.contains(&peer) {
3245                        timed_out.push(peer);
3246                    }
3247                }
3248                timed_out.sort_unstable_by_key(|peer| peer.0);
3249
3250                if !failures.is_empty() {
3251                    return QuorumOutcome::Failed { failures };
3252                }
3253
3254                if !timed_out.is_empty() || successful.len() < expected.len() {
3255                    let got = successful;
3256                    let mut missing = timed_out;
3257                    for &peer in expected {
3258                        if !got.contains(&peer) && !missing.contains(&peer) {
3259                            missing.push(peer);
3260                        }
3261                    }
3262                    missing.sort_unstable_by_key(|peer| peer.0);
3263                    return QuorumOutcome::TimedOut { got, missing };
3264                }
3265
3266                if prepare.assignment_fence.is_some() {
3267                    if let Err(error) = mark_prepare_quorum_reached(&state, prepare) {
3268                        return QuorumOutcome::Failed {
3269                            failures: vec![(
3270                                expected_roster
3271                                    .first()
3272                                    .copied()
3273                                    .unwrap_or(NodeId::UNASSIGNED),
3274                                error,
3275                            )],
3276                        };
3277                    }
3278                }
3279
3280                return QuorumOutcome::Reached {
3281                    acks: successful,
3282                    follower_watermark: follower_watermark
3283                        .unwrap_or(CheckpointWatermark::Uninitialized),
3284                };
3285            }
3286        }
3287
3288        let start = Instant::now();
3289        let expected_set: FxHashSet<NodeId> = expected.iter().copied().collect();
3290        let mut successful: Vec<NodeId> = Vec::new();
3291        let mut failures: Vec<(NodeId, String)> = Vec::new();
3292        let mut follower_watermark: Option<CheckpointWatermark>;
3293
3294        loop {
3295            successful.clear();
3296            failures.clear();
3297            follower_watermark = None;
3298
3299            for (from, json) in self.kv.scan(ACK_KEY).await {
3300                if !expected_set.contains(&from) {
3301                    continue;
3302                }
3303                let Ok(ack) = serde_json::from_str::<BarrierAck>(&json) else {
3304                    continue;
3305                };
3306                if ack.epoch != epoch
3307                    || ack.checkpoint_id != checkpoint_id
3308                    || ack.assignment_digest != assignment_digest
3309                {
3310                    continue;
3311                }
3312                if ack.ok {
3313                    if let Err(error) = ack.watermark.validate() {
3314                        failures.push((from, error));
3315                    } else {
3316                        successful.push(from);
3317                        follower_watermark =
3318                            Some(follower_watermark.map_or(ack.watermark, |current| {
3319                                current.cluster_min(ack.watermark)
3320                            }));
3321                    }
3322                } else {
3323                    failures.push((from, ack.error.unwrap_or_default()));
3324                }
3325            }
3326
3327            if !failures.is_empty() {
3328                return QuorumOutcome::Failed { failures };
3329            }
3330            if successful.len() == expected.len() {
3331                return QuorumOutcome::Reached {
3332                    acks: successful,
3333                    follower_watermark: follower_watermark
3334                        .unwrap_or(CheckpointWatermark::Uninitialized),
3335                };
3336            }
3337            if start.elapsed() >= deadline {
3338                let got: FxHashSet<NodeId> = successful.iter().copied().collect();
3339                let missing: Vec<NodeId> = expected
3340                    .iter()
3341                    .copied()
3342                    .filter(|n| !got.contains(n))
3343                    .collect();
3344                return QuorumOutcome::TimedOut {
3345                    got: successful,
3346                    missing,
3347                };
3348            }
3349            tokio::time::sleep(Duration::from_millis(50)).await;
3350        }
3351    }
3352}
3353
3354#[cfg(test)]
3355mod tests;