1#[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
22pub const ANNOUNCEMENT_KEY: &str = "control:barrier";
24
25pub const ACK_KEY: &str = "control:barrier-ack";
27
28#[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#[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#[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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
197pub enum Phase {
198 Prepare,
201 Aligned,
205 Commit,
207 Abort,
209}
210
211const fn is_terminal_phase(phase: Phase) -> bool {
212 matches!(phase, Phase::Commit | Phase::Abort)
213}
214
215#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
217pub struct BarrierAnnouncement {
218 pub epoch: u64,
220 pub checkpoint_id: u64,
222 #[serde(default)]
225 pub assignment_fence: Option<super::CheckpointAssignmentFence>,
226 #[serde(default)]
231 pub leader_proof: Option<super::LeaderProof>,
232 pub phase: Phase,
234 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
282fn 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(¤t, &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#[cfg(feature = "cluster")]
332fn merge_direct_announcement(
333 current: BarrierAnnouncement,
334 incoming: BarrierAnnouncement,
335) -> Result<BarrierAnnouncement, String> {
336 validate_announcement_attempt(¤t)?;
337 validate_announcement_attempt(&incoming)?;
338 match incoming.checkpoint_id.cmp(¤t.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
345fn 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
381fn merge_scanned_announcement(
384 current: BarrierAnnouncement,
385 incoming: BarrierAnnouncement,
386) -> Result<BarrierAnnouncement, String> {
387 validate_announcement_attempt(¤t)?;
388 validate_announcement_attempt(&incoming)?;
389 match incoming.checkpoint_id.cmp(¤t.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(¤t.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 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
450pub struct BarrierAck {
451 pub epoch: u64,
453 #[serde(default)]
455 pub checkpoint_id: u64,
456 #[serde(default)]
458 pub assignment_digest: Option<[u8; 32]>,
459 pub ok: bool,
461 pub error: Option<String>,
463 #[serde(default)]
466 pub watermark: CheckpointWatermark,
467}
468
469#[derive(Debug, Clone, PartialEq, Eq)]
471pub enum QuorumOutcome {
472 Reached {
474 acks: Vec<NodeId>,
476 follower_watermark: CheckpointWatermark,
478 },
479 TimedOut {
481 got: Vec<NodeId>,
483 missing: Vec<NodeId>,
485 },
486 Failed {
488 failures: Vec<(NodeId, String)>,
490 },
491}
492
493#[async_trait]
495pub trait ClusterKv: Send + Sync + 'static {
496 async fn write(&self, key: &str, value: String);
498 async fn write_checked(&self, key: &str, value: String) -> Result<(), String> {
508 self.write(key, value).await;
509 Ok(())
510 }
511 async fn read_from(&self, who: NodeId, key: &str) -> Option<String>;
513 async fn read_from_checked(&self, who: NodeId, key: &str) -> Result<Option<String>, String> {
519 Ok(self.read_from(who, key).await)
520 }
521 async fn scan(&self, key: &str) -> Vec<(NodeId, String)>;
523 async fn scan_checked(&self, key: &str) -> Result<Vec<(NodeId, String)>, String> {
528 Ok(self.scan(key).await)
529 }
530}
531
532#[derive(Debug)]
534pub struct InMemoryKv {
535 local_id: NodeId,
536 state: Mutex<FxHashMap<(NodeId, String), String>>,
537}
538
539impl InMemoryKv {
540 #[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 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#[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#[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#[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#[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_rx: watch::Receiver<Option<BarrierAnnouncement>>,
779 incoming_tx: crossfire::MAsyncTx<BarrierFlavor>,
783 merge_error: Arc<parking_lot::Mutex<Option<String>>>,
784 prepare_acks: Arc<parking_lot::Mutex<PrepareAckState>>,
785 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 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#[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#[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#[cfg(feature = "cluster")]
1696enum PeerFailure {
1697 Unreachable,
1698 Nack(String),
1699}
1700
1701#[cfg(feature = "cluster")]
1703struct PrepareFanoutBatch {
1704 announcement: BarrierAnnouncement,
1705 expected: Vec<NodeId>,
1706 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 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(¤t.announcement().checkpoint_id)
1919 {
1920 std::cmp::Ordering::Greater => {
1921 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::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 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
2182pub struct BarrierCoordinator {
2184 kv: Arc<dyn ClusterKv>,
2185 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 #[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 #[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 #[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 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 #[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 #[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 #[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 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 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 #[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 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 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 #[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 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 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 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 if is_terminal_phase(ann.phase) {
2871 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 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 roster
2906 } else {
2907 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 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 #[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 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 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 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 #[allow(clippy::too_many_lines)]
3070 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 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;