1use std::collections::BinaryHeap;
4use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
5use std::sync::Arc;
6use std::time::Duration;
7
8use bytes::Bytes;
9use futures::StreamExt;
10use object_store::path::Path as OsPath;
11use object_store::{ObjectStore, ObjectStoreExt, PutMode, PutOptions, PutPayload};
12use uuid::Uuid;
13
14use crate::cluster::discovery::NodeId;
15
16use super::{
17 fence_path, lease_path, lease_prefix, sequence_from_path, successor_fence_path, ProcessLease,
18 ProcessLeaseError, ProcessLeaseFence, ProcessLeaseObservation, ProcessLeaseOutcome,
19 MAX_PROCESS_LEASE_FENCE_BYTES, MAX_PROCESS_LEASE_RECORD_BYTES,
20 PROCESS_LEASE_HEAD_READ_ATTEMPTS, PROCESS_LEASE_HISTORY_TO_RETAIN,
21 PROCESS_LEASE_MAX_LIST_RECORDS, PROCESS_LEASE_MAX_PRUNE_BATCHES,
22 PROCESS_LEASE_PRUNE_BATCH_RECORDS, PROCESS_LEASE_PRUNE_IO_TIMEOUT,
23 PROCESS_LEASE_PRUNE_READ_CONCURRENCY, PROCESS_LEASE_WRITES_PER_PRUNE,
24};
25
26pub struct ProcessLeaseStore {
28 pub(super) store: Arc<dyn ObjectStore>,
29 pub(super) node: NodeId,
30 pub(super) ttl_ms: i64,
31 prune_running: Arc<AtomicBool>,
32 pub(super) prune_healthy: Arc<AtomicBool>,
33 writes_since_prune: AtomicU64,
34 sealed_term: AtomicU64,
35}
36
37impl std::fmt::Debug for ProcessLeaseStore {
38 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
39 formatter
40 .debug_struct("ProcessLeaseStore")
41 .field("node", &self.node)
42 .field("ttl_ms", &self.ttl_ms)
43 .finish_non_exhaustive()
44 }
45}
46
47impl ProcessLeaseStore {
48 #[must_use]
50 pub fn new(store: Arc<dyn ObjectStore>, node: NodeId, ttl_ms: i64) -> Self {
51 Self {
52 store,
53 node,
54 ttl_ms,
55 prune_running: Arc::new(AtomicBool::new(false)),
56 prune_healthy: Arc::new(AtomicBool::new(true)),
57 writes_since_prune: AtomicU64::new(0),
58 sealed_term: AtomicU64::new(0),
59 }
60 }
61
62 async fn list_seqs_from(
63 store: &Arc<dyn ObjectStore>,
64 node: NodeId,
65 ) -> Result<Vec<u64>, ProcessLeaseError> {
66 let prefix_string = lease_prefix(node);
67 let prefix = OsPath::from(prefix_string.clone());
68 let mut entries = store.list(Some(&prefix));
69 let mut sequences = Vec::new();
70 while let Some(entry) = entries.next().await {
71 let entry = entry.map_err(|error| ProcessLeaseError::Io(error.to_string()))?;
72 if sequences.len() == PROCESS_LEASE_MAX_LIST_RECORDS {
73 return Err(ProcessLeaseError::Invalid(format!(
74 "process lease history exceeds the fixed {PROCESS_LEASE_MAX_LIST_RECORDS}-record scan bound"
75 )));
76 }
77 sequences.push(sequence_from_path(node, &entry.location)?);
78 }
79 sequences.sort_unstable();
80 sequences.dedup();
81 Ok(sequences)
82 }
83
84 pub(super) async fn list_seqs(&self) -> Result<Vec<u64>, ProcessLeaseError> {
85 Self::list_seqs_from(&self.store, self.node).await
86 }
87
88 async fn oldest_prune_window(
89 store: &Arc<dyn ObjectStore>,
90 node: NodeId,
91 ) -> Result<(u64, Vec<u64>), ProcessLeaseError> {
92 let prefix = OsPath::from(lease_prefix(node));
93 let mut entries = store.list(Some(&prefix));
94 let window_len = PROCESS_LEASE_PRUNE_BATCH_RECORDS
95 .checked_add(1)
96 .ok_or_else(|| {
97 ProcessLeaseError::Invalid("process lease prune window overflow".into())
98 })?;
99 let mut oldest = BinaryHeap::with_capacity(window_len);
100 let mut count = 0_u64;
101 while let Some(entry) = entries.next().await {
102 let entry = entry.map_err(|error| ProcessLeaseError::Io(error.to_string()))?;
103 let sequence = sequence_from_path(node, &entry.location)?;
104 count = count.checked_add(1).ok_or_else(|| {
105 ProcessLeaseError::Invalid("process lease history count exhausted".into())
106 })?;
107 if oldest.len() < window_len {
108 oldest.push(sequence);
109 } else if oldest.peek().is_some_and(|largest| sequence < *largest) {
110 oldest.pop();
111 oldest.push(sequence);
112 }
113 }
114 Ok((count, oldest.into_sorted_vec()))
115 }
116
117 fn schedule_history_prune(&self, force: bool) {
118 if !force
119 && self.writes_since_prune.fetch_add(1, Ordering::AcqRel) + 1
120 < PROCESS_LEASE_WRITES_PER_PRUNE
121 {
122 return;
123 }
124 if self
125 .prune_running
126 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
127 .is_err()
128 {
129 return;
130 }
131 self.writes_since_prune.store(0, Ordering::Release);
132 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
133 self.prune_running.store(false, Ordering::Release);
134 self.prune_healthy.store(false, Ordering::Release);
135 return;
136 };
137 let store = Arc::clone(&self.store);
138 let node = self.node;
139 let prune_running = Arc::clone(&self.prune_running);
140 let prune_healthy = Arc::clone(&self.prune_healthy);
141 runtime.spawn(async move {
142 let prune = Self::prune_history(&store, node).await;
143 if let Err(error) = prune {
144 prune_healthy.store(false, Ordering::Release);
145 tracing::warn!(node = node.0, %error, "process lease history prune failed");
146 } else {
147 prune_healthy.store(true, Ordering::Release);
148 }
149 prune_running.store(false, Ordering::Release);
150 });
151 }
152
153 pub(super) async fn prune_history(
154 store: &Arc<dyn ObjectStore>,
155 node: NodeId,
156 ) -> Result<(), ProcessLeaseError> {
157 for _ in 0..PROCESS_LEASE_MAX_PRUNE_BATCHES {
158 let done = tokio::time::timeout(
159 PROCESS_LEASE_PRUNE_IO_TIMEOUT,
160 Self::prune_history_batch(store, node),
161 )
162 .await
163 .map_err(|_| ProcessLeaseError::Io("process lease history prune timed out".into()))??;
164 if done {
165 return Ok(());
166 }
167 tokio::task::yield_now().await;
168 }
169 Err(ProcessLeaseError::Io(
170 "process lease history still exceeds the bounded prune budget".into(),
171 ))
172 }
173
174 async fn repair_unhealthy_prune(&self) -> Result<(), ProcessLeaseError> {
175 if self.prune_healthy.load(Ordering::Acquire) {
176 return Ok(());
177 }
178 Self::prune_history(&self.store, self.node).await?;
179 self.prune_healthy.store(true, Ordering::Release);
180 self.writes_since_prune.store(0, Ordering::Release);
181 Ok(())
182 }
183
184 pub(super) async fn prune_history_batch(
185 store: &Arc<dyn ObjectStore>,
186 node: NodeId,
187 ) -> Result<bool, ProcessLeaseError> {
188 let (count, sequences) = Self::oldest_prune_window(store, node).await?;
189 let retain = u64::try_from(PROCESS_LEASE_HISTORY_TO_RETAIN)
190 .map_err(|_| ProcessLeaseError::Invalid("process lease retention overflow".into()))?;
191 if count <= retain {
192 return Ok(true);
193 }
194 let delete_count_u64 =
195 (count - retain).min(u64::try_from(PROCESS_LEASE_PRUNE_BATCH_RECORDS).map_err(
196 |_| ProcessLeaseError::Invalid("process lease prune batch overflow".into()),
197 )?);
198 let delete_count = usize::try_from(delete_count_u64).map_err(|_| {
199 ProcessLeaseError::Invalid("process lease delete count overflow".into())
200 })?;
201 if sequences.len() <= delete_count {
202 return Err(ProcessLeaseError::Invalid(
203 "process lease prune window is missing its retained boundary".into(),
204 ));
205 }
206 let mut reads = futures::stream::iter(sequences.iter().copied().map(|sequence| {
207 let store = Arc::clone(store);
208 async move {
209 Self::read_record_from(&store, node, sequence)
210 .await
211 .map(|record| record.map(|record| (sequence, record)))
212 }
213 }))
214 .buffer_unordered(PROCESS_LEASE_PRUNE_READ_CONCURRENCY);
215 let mut records = Vec::with_capacity(sequences.len());
216 while let Some(record) = reads.next().await {
217 let Some(record) = record? else {
218 return Ok(false);
219 };
220 records.push(record);
221 }
222 records.sort_unstable_by_key(|(sequence, _)| *sequence);
223 for pair in records.windows(2) {
227 let (left_sequence, left) = &pair[0];
228 let (right_sequence, right) = &pair[1];
229 if left_sequence.checked_add(1) != Some(*right_sequence) {
230 return Err(ProcessLeaseError::Invalid(
231 "process lease history contains a noncontiguous sequence".into(),
232 ));
233 }
234 if left.owner != right.owner || left.term != right.term {
235 if left.term.checked_add(1) != Some(right.term) || left.owner == right.owner {
236 return Err(ProcessLeaseError::Invalid(
237 "process lease history contains a noncanonical term transition".into(),
238 ));
239 }
240 Self::seal_fence(store, &ProcessLeaseFence::new(left.clone(), right.clone())?)
241 .await?;
242 }
243 }
244 let deletable = records
245 .iter()
246 .take(delete_count)
247 .map(|(sequence, _)| *sequence)
248 .collect::<Vec<_>>();
249 let deletions = futures::stream::iter(
250 deletable
251 .into_iter()
252 .map(move |sequence| Ok::<_, object_store::Error>(lease_path(node, sequence))),
253 )
254 .boxed();
255 let mut results = store.delete_stream(deletions);
256 while let Some(result) = results.next().await {
257 if let Err(error) = result {
258 if !matches!(error, object_store::Error::NotFound { .. }) {
259 return Err(ProcessLeaseError::Io(error.to_string()));
260 }
261 }
262 }
263 Ok(count.saturating_sub(delete_count_u64) <= retain)
264 }
265
266 async fn read_record_from(
267 store: &Arc<dyn ObjectStore>,
268 node: NodeId,
269 sequence: u64,
270 ) -> Result<Option<ProcessLease>, ProcessLeaseError> {
271 let result = match store.get(&lease_path(node, sequence)).await {
272 Ok(result) => result,
273 Err(object_store::Error::NotFound { .. }) => return Ok(None),
274 Err(error) => return Err(ProcessLeaseError::Io(error.to_string())),
275 };
276 if result.meta.size == 0 || result.meta.size > MAX_PROCESS_LEASE_RECORD_BYTES {
277 return Err(ProcessLeaseError::Invalid(format!(
278 "process lease record is {} bytes; expected 1..={MAX_PROCESS_LEASE_RECORD_BYTES}",
279 result.meta.size
280 )));
281 }
282 let bytes = result
283 .bytes()
284 .await
285 .map_err(|error| ProcessLeaseError::Io(error.to_string()))?;
286 let lease: ProcessLease = serde_json::from_slice(&bytes)?;
287 lease.validate(node)?;
288 if lease.seq != sequence {
289 return Err(ProcessLeaseError::Invalid(
290 "record sequence does not match its object name".into(),
291 ));
292 }
293 let canonical = serde_json::to_vec(&lease)?;
294 if canonical.as_slice() != bytes.as_ref() {
295 return Err(ProcessLeaseError::Invalid(format!(
296 "process lease record {sequence} does not use its canonical body"
297 )));
298 }
299 Ok(Some(lease))
300 }
301
302 async fn read_history(
303 store: &Arc<dyn ObjectStore>,
304 node: NodeId,
305 ) -> Result<Vec<(u64, ProcessLease)>, ProcessLeaseError> {
306 for attempt in 0..PROCESS_LEASE_HEAD_READ_ATTEMPTS {
307 let sequences = Self::list_seqs_from(store, node).await?;
308 let mut reads = futures::stream::iter(sequences.iter().copied().map(|sequence| {
309 let store = Arc::clone(store);
310 async move {
311 Self::read_record_from(&store, node, sequence)
312 .await
313 .map(|record| record.map(|record| (sequence, record)))
314 }
315 }))
316 .buffer_unordered(PROCESS_LEASE_PRUNE_READ_CONCURRENCY);
317 let mut records = Vec::with_capacity(sequences.len());
318 let mut changed = false;
319 while let Some(record) = reads.next().await {
320 match record? {
321 Some(record) => records.push(record),
322 None => changed = true,
323 }
324 }
325 if !changed {
326 records.sort_unstable_by_key(|(sequence, _)| *sequence);
327 return Ok(records);
328 }
329 if attempt + 1 < PROCESS_LEASE_HEAD_READ_ATTEMPTS {
330 tokio::task::yield_now().await;
331 }
332 }
333 Err(ProcessLeaseError::Io(format!(
334 "process lease history changed during {PROCESS_LEASE_HEAD_READ_ATTEMPTS} read attempts"
335 )))
336 }
337
338 async fn load_fence_at(
339 store: &Arc<dyn ObjectStore>,
340 path: &OsPath,
341 ) -> Result<Option<ProcessLeaseFence>, ProcessLeaseError> {
342 let result = match store.get(path).await {
343 Ok(result) => result,
344 Err(object_store::Error::NotFound { .. }) => return Ok(None),
345 Err(error) => return Err(ProcessLeaseError::Io(error.to_string())),
346 };
347 if result.meta.size == 0 || result.meta.size > MAX_PROCESS_LEASE_FENCE_BYTES {
348 return Err(ProcessLeaseError::Invalid(format!(
349 "process lease fence is {} bytes; expected 1..={MAX_PROCESS_LEASE_FENCE_BYTES}",
350 result.meta.size
351 )));
352 }
353 let bytes = result
354 .bytes()
355 .await
356 .map_err(|error| ProcessLeaseError::Io(error.to_string()))?;
357 let fence: ProcessLeaseFence = serde_json::from_slice(&bytes)?;
358 if !fence.is_canonical() {
359 return Err(ProcessLeaseError::Invalid(
360 "indexed process lease fence is not canonical".into(),
361 ));
362 }
363 let canonical = serde_json::to_vec(&fence)?;
364 if canonical.as_slice() != bytes.as_ref() {
365 return Err(ProcessLeaseError::Invalid(
366 "indexed process lease fence does not use its canonical body".into(),
367 ));
368 }
369 Ok(Some(fence))
370 }
371
372 pub(super) async fn load_fence(
373 store: &Arc<dyn ObjectStore>,
374 node: NodeId,
375 predecessor: Uuid,
376 ) -> Result<Option<ProcessLeaseFence>, ProcessLeaseError> {
377 let fence = Self::load_fence_at(store, &fence_path(node, predecessor)).await?;
378 if fence.as_ref().is_some_and(|fence| {
379 fence.predecessor.node != node || fence.predecessor.owner != predecessor
380 }) {
381 return Err(ProcessLeaseError::Invalid(
382 "predecessor-indexed process lease fence has the wrong identity".into(),
383 ));
384 }
385 Ok(fence)
386 }
387
388 async fn load_successor_fence(
389 store: &Arc<dyn ObjectStore>,
390 node: NodeId,
391 successor: Uuid,
392 term: u64,
393 ) -> Result<Option<ProcessLeaseFence>, ProcessLeaseError> {
394 let fence =
395 Self::load_fence_at(store, &successor_fence_path(node, successor, term)).await?;
396 if fence.as_ref().is_some_and(|fence| {
397 fence.successor.node != node
398 || fence.successor.owner != successor
399 || fence.successor.term != term
400 }) {
401 return Err(ProcessLeaseError::Invalid(
402 "successor-indexed process lease fence has the wrong identity".into(),
403 ));
404 }
405 Ok(fence)
406 }
407
408 async fn seal_fence_at(
409 store: &Arc<dyn ObjectStore>,
410 path: &OsPath,
411 fence: &ProcessLeaseFence,
412 ) -> Result<(), ProcessLeaseError> {
413 let options = PutOptions {
414 mode: PutMode::Create,
415 ..PutOptions::default()
416 };
417 let payload = PutPayload::from(Bytes::from(serde_json::to_vec(fence)?));
418 let put_error = store.put_opts(path, payload, options).await.err();
419 match Self::load_fence_at(store, path).await {
420 Ok(Some(stored)) if stored == *fence => Ok(()),
421 Ok(Some(_)) => Err(ProcessLeaseError::Invalid(
422 "process fence index maps to conflicting takeover evidence".into(),
423 )),
424 Ok(None) => Err(ProcessLeaseError::Io(
425 "process lease fence write was not durably visible".into(),
426 )),
427 Err(reconcile_error) => {
428 if let Some(put_error) = put_error {
429 Err(ProcessLeaseError::Io(format!(
430 "process lease fence write failed ({put_error}); reconciliation failed ({reconcile_error})"
431 )))
432 } else {
433 Err(reconcile_error)
434 }
435 }
436 }
437 }
438
439 async fn seal_fence(
440 store: &Arc<dyn ObjectStore>,
441 fence: &ProcessLeaseFence,
442 ) -> Result<(), ProcessLeaseError> {
443 if !fence.is_canonical() {
444 return Err(ProcessLeaseError::Invalid(
445 "cannot seal a noncanonical process lease fence".into(),
446 ));
447 }
448 Self::seal_fence_at(
449 store,
450 &fence_path(fence.predecessor.node, fence.predecessor.owner),
451 fence,
452 )
453 .await?;
454 Self::seal_fence_at(
455 store,
456 &successor_fence_path(
457 fence.successor.node,
458 fence.successor.owner,
459 fence.successor.term,
460 ),
461 fence,
462 )
463 .await
464 }
465
466 async fn read_record(&self, sequence: u64) -> Result<Option<ProcessLease>, ProcessLeaseError> {
467 Self::read_record_from(&self.store, self.node, sequence).await
468 }
469
470 pub(super) async fn find_takeover_from(
471 &self,
472 owner: Uuid,
473 ) -> Result<Option<ProcessLeaseFence>, ProcessLeaseError> {
474 if let Some(fence) = Self::load_fence(&self.store, self.node, owner).await? {
475 return Ok(Some(fence));
476 }
477 let records = Self::read_history(&self.store, self.node).await?;
478 let mut found = None;
479 for pair in records.windows(2) {
480 let (predecessor_sequence, predecessor) = &pair[0];
481 let (successor_sequence, successor) = &pair[1];
482 if predecessor.owner != owner {
483 continue;
484 }
485 if predecessor_sequence.checked_add(1) != Some(*successor_sequence) {
486 continue;
487 }
488 let Ok(fence) = ProcessLeaseFence::new(predecessor.clone(), successor.clone()) else {
489 continue;
490 };
491 if found.replace(fence).is_some() {
492 return Err(ProcessLeaseError::Invalid(
493 "process owner appears in more than one durable takeover transition".into(),
494 ));
495 }
496 }
497 if let Some(fence) = found {
498 Self::seal_fence(&self.store, &fence).await?;
499 Ok(Some(fence))
500 } else {
501 Self::load_fence(&self.store, self.node, owner).await
504 }
505 }
506
507 pub(super) async fn ensure_current_term_fence(
508 &self,
509 current: &ProcessLease,
510 ) -> Result<(), ProcessLeaseError> {
511 if current.term == 1 || self.sealed_term.load(Ordering::Acquire) == current.term {
512 self.sealed_term.store(current.term, Ordering::Release);
513 return Ok(());
514 }
515 if let Some(fence) =
516 Self::load_successor_fence(&self.store, self.node, current.owner, current.term).await?
517 {
518 if fence.successor.seq > current.seq {
519 return Err(ProcessLeaseError::Invalid(
520 "process term fence starts after the current lease head".into(),
521 ));
522 }
523 self.sealed_term.store(current.term, Ordering::Release);
524 return Ok(());
525 }
526
527 let records = Self::read_history(&self.store, self.node).await?;
528 for pair in records.windows(2) {
529 let (left_sequence, left) = &pair[0];
530 let (right_sequence, right) = &pair[1];
531 if left_sequence.checked_add(1) == Some(*right_sequence)
532 && right.owner == current.owner
533 && right.term == current.term
534 {
535 let fence = ProcessLeaseFence::new(left.clone(), right.clone())?;
536 Self::seal_fence(&self.store, &fence).await?;
537 self.sealed_term.store(current.term, Ordering::Release);
538 return Ok(());
539 }
540 }
541 Err(ProcessLeaseError::Invalid(
542 "current process term has no durable takeover evidence".into(),
543 ))
544 }
545
546 pub async fn load(&self) -> Result<Option<ProcessLease>, ProcessLeaseError> {
551 let mut observed_head = false;
552 for attempt in 0..PROCESS_LEASE_HEAD_READ_ATTEMPTS {
553 let sequences = match self.list_seqs().await {
554 Ok(sequences) => sequences,
555 Err(error) => {
556 if matches!(&error, ProcessLeaseError::Invalid(_)) {
561 self.schedule_history_prune(true);
562 }
563 return Err(error);
564 }
565 };
566 let Some(sequence) = sequences.last().copied() else {
567 if !observed_head {
568 return Ok(None);
569 }
570 if attempt + 1 < PROCESS_LEASE_HEAD_READ_ATTEMPTS {
571 tokio::task::yield_now().await;
572 continue;
573 }
574 break;
575 };
576 observed_head = true;
577 match self.read_record(sequence).await? {
578 Some(lease) => return Ok(Some(lease)),
579 None if attempt + 1 < PROCESS_LEASE_HEAD_READ_ATTEMPTS => {
580 tokio::task::yield_now().await;
581 }
582 None => break,
583 }
584 }
585 Err(ProcessLeaseError::Io(format!(
586 "process lease head changed during {PROCESS_LEASE_HEAD_READ_ATTEMPTS} read attempts"
587 )))
588 }
589
590 pub async fn try_acquire(
595 &self,
596 owner: Uuid,
597 now_ms: i64,
598 ) -> Result<ProcessLeaseOutcome, ProcessLeaseError> {
599 if owner.is_nil() || self.node.is_unassigned() || self.ttl_ms <= 0 {
600 return Err(ProcessLeaseError::Invalid(
601 "node, owner, and lease TTL must be nonzero".into(),
602 ));
603 }
604 self.repair_unhealthy_prune().await?;
605 let current = self.load().await?;
606 if let Some(lease) = current.as_ref().filter(|lease| lease.owner == owner) {
607 self.ensure_current_term_fence(lease).await?;
608 }
609 let candidate = match current {
610 None => ProcessLease {
611 node: self.node,
612 owner,
613 term: 1,
614 seq: 1,
615 expires_at_ms: now_ms.saturating_add(self.ttl_ms),
616 },
617 Some(ref lease) if lease.owner == owner => ProcessLease {
618 node: self.node,
619 owner,
620 term: lease.term,
621 seq: lease
622 .seq
623 .checked_add(1)
624 .ok_or_else(|| ProcessLeaseError::Invalid("lease sequence exhausted".into()))?,
625 expires_at_ms: now_ms.saturating_add(self.ttl_ms),
626 },
627 Some(lease) => return Ok(ProcessLeaseOutcome::Held(lease)),
628 };
629 if candidate.term == 0 || candidate.seq == 0 {
630 return Err(ProcessLeaseError::Invalid(
631 "lease term or sequence exhausted".into(),
632 ));
633 }
634
635 let options = PutOptions {
636 mode: PutMode::Create,
637 ..PutOptions::default()
638 };
639 let payload = PutPayload::from(Bytes::from(serde_json::to_vec(&candidate)?));
640 match self
641 .store
642 .put_opts(&lease_path(self.node, candidate.seq), payload, options)
643 .await
644 {
645 Ok(_) => {
646 self.sealed_term.store(candidate.term, Ordering::Release);
647 self.schedule_history_prune(false);
648 Ok(ProcessLeaseOutcome::Acquired(candidate))
649 }
650 Err(
651 object_store::Error::AlreadyExists { .. }
652 | object_store::Error::Precondition { .. },
653 ) => {
654 let winner = self.load().await?.ok_or_else(|| {
655 ProcessLeaseError::Io("CAS conflict but the winner was not readable".into())
656 })?;
657 self.schedule_history_prune(false);
658 if winner.owner == owner {
659 self.ensure_current_term_fence(&winner).await?;
660 Ok(ProcessLeaseOutcome::Acquired(winner))
661 } else {
662 Ok(ProcessLeaseOutcome::Held(winner))
663 }
664 }
665 Err(error) => Err(ProcessLeaseError::Io(error.to_string())),
666 }
667 }
668
669 pub fn observe_rival(
674 &self,
675 lease: &ProcessLease,
676 ) -> Result<ProcessLeaseObservation, ProcessLeaseError> {
677 lease.validate(self.node)?;
678 Ok(ProcessLeaseObservation {
679 lease: lease.clone(),
680 started: std::time::Instant::now(),
681 })
682 }
683
684 pub async fn try_takeover(
690 &self,
691 owner: Uuid,
692 observation: &ProcessLeaseObservation,
693 now_ms: i64,
694 ) -> Result<ProcessLeaseOutcome, ProcessLeaseError> {
695 if owner.is_nil() || self.ttl_ms <= 0 {
696 return Err(ProcessLeaseError::Invalid(
697 "takeover owner and lease TTL must be nonzero".into(),
698 ));
699 }
700 self.repair_unhealthy_prune().await?;
701 observation.lease.validate(self.node)?;
702 if observation.started.elapsed()
703 < Duration::from_millis(u64::try_from(self.ttl_ms).unwrap_or(u64::MAX))
704 {
705 return Ok(ProcessLeaseOutcome::Held(observation.lease.clone()));
706 }
707 let current = self.load().await?.ok_or_else(|| {
708 ProcessLeaseError::Invalid("observed process lease disappeared".into())
709 })?;
710 if current != observation.lease {
711 return Ok(ProcessLeaseOutcome::Held(current));
712 }
713 let candidate = ProcessLease {
714 node: self.node,
715 owner,
716 term: current
717 .term
718 .checked_add(1)
719 .ok_or_else(|| ProcessLeaseError::Invalid("process term exhausted".into()))?,
720 seq: current
721 .seq
722 .checked_add(1)
723 .ok_or_else(|| ProcessLeaseError::Invalid("lease sequence exhausted".into()))?,
724 expires_at_ms: now_ms.saturating_add(self.ttl_ms),
725 };
726 let options = PutOptions {
727 mode: PutMode::Create,
728 ..PutOptions::default()
729 };
730 let payload = PutPayload::from(Bytes::from(serde_json::to_vec(&candidate)?));
731 match self
732 .store
733 .put_opts(&lease_path(self.node, candidate.seq), payload, options)
734 .await
735 {
736 Ok(_) => {
737 let fence = ProcessLeaseFence::new(current, candidate.clone())?;
738 Self::seal_fence(&self.store, &fence).await?;
739 self.sealed_term.store(candidate.term, Ordering::Release);
740 self.schedule_history_prune(true);
741 Ok(ProcessLeaseOutcome::Acquired(candidate))
742 }
743 Err(
744 object_store::Error::AlreadyExists { .. }
745 | object_store::Error::Precondition { .. },
746 ) => {
747 let winner = self.load().await?.ok_or_else(|| {
748 ProcessLeaseError::Io("takeover CAS winner was not readable".into())
749 })?;
750 self.schedule_history_prune(true);
751 Ok(ProcessLeaseOutcome::Held(winner))
752 }
753 Err(error) => Err(ProcessLeaseError::Io(error.to_string())),
754 }
755 }
756}