Skip to main content

laminar_core/cluster/control/process_lease/
store.rs

1//! Append-only lease storage, history retention, and indexed takeover certificates.
2
3use 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
26/// Append-only object-store authority for one stable node identity.
27pub 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    /// Create an authority for `node`.
49    #[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        // Seal every newly observed term transition under the predecessor boot identity before
224        // deleting either side. The indexed certificate keeps recovery lookup O(1) and prevents
225        // routine renewals or many later terms from erasing takeover evidence.
226        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            // A concurrent pruner may have sealed the direct index and removed the history pair
502            // after our first point read but before the scan completed.
503            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    /// Load the highest durable sequence.
547    ///
548    /// # Errors
549    /// Fails on object-store I/O, malformed JSON, or a record in the wrong node namespace.
550    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                    // RECOVERY: A transient LIST failure says nothing about retention health.
557                    // Starting a prune would duplicate authority I/O during the outage and could
558                    // force the next renewal through synchronous repair. Invalid inventory still
559                    // triggers the bounded repair path.
560                    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    /// Acquire or renew this stable node identity for `owner` at `now_ms`.
591    ///
592    /// # Errors
593    /// Fails closed on object-store I/O or an invalid durable record.
594    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    /// Start a candidate-local monotonic observation of a rival lease.
670    ///
671    /// # Errors
672    /// Rejects a malformed record or one from another node namespace.
673    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    /// Attempt takeover after the same sequence and owner have been observed unchanged for a
685    /// full TTL on this candidate's monotonic clock.
686    ///
687    /// # Errors
688    /// Fails closed on early observation, object-store I/O, or malformed durable state.
689    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}