Skip to main content

rightkit_control/
lease.rs

1//! PTY-002 input lease: epoch fencing plus monotonic, acknowledged input
2//! sequence so a stale controller can never inject input.
3//!
4//! Contract (CodeRight canon PTY-002): "Input ownership uses a fenced lease
5//! epoch plus monotonic input sequence and acknowledgement; stale-lease writes
6//! are rejected, and reconnect explicitly reconciles or rejects
7//! unacknowledged bytes rather than replaying them ambiguously." Shape follows
8//! CodeRight `ProcessLeaseSnapshot` / `UnacknowledgedInputDecision`
9//! (`engine/crates/tools/src/background_shell/process_model.rs`): attach bumps
10//! the epoch; a stale epoch is rejected; unacknowledged input is reconciled
11//! before new input.
12//!
13//! Sequence model (one global sequence per lease, continuing across epochs):
14//! - `next_input_sequence` is the only sequence accepted as new input.
15//! - every sequence `< acked_input_sequence` was delivered and acknowledged.
16//! - `[acked_input_sequence, next_input_sequence)` is the unacknowledged
17//!   range: input whose delivery is ambiguous (executor failed or panicked
18//!   after it may have produced a side effect). While it is non-empty, new
19//!   input and plain `acquire` are rejected until the holder or a new
20//!   controller explicitly reconciles it.
21//!
22//! Input rules for a submit `(epoch, sequence)`:
23//! 1. lease not attached → rejected (`NotAttached`).
24//! 2. `epoch < current` → rejected (`StaleEpoch`): the holder was fenced.
25//!    `epoch > current` → rejected (`FutureEpoch`).
26//! 3. `sequence < acked_input_sequence` → deduplicated: acknowledged again,
27//!    not executed (`InputOutcome::Duplicate`).
28//! 4. `sequence` inside the unacknowledged range, or any new sequence while
29//!    that range is non-empty → rejected (`UnacknowledgedInput`).
30//! 5. `sequence > next_input_sequence` → rejected (`OutOfOrder`); no gap
31//!    buffering.
32//! 6. `sequence == next_input_sequence` → executed exactly once while the
33//!    lease is held, then acknowledged.
34//!
35//! Submissions are serialized: the lease lock is held across the executor, so
36//! an acquisition that fences a holder waits for its in-flight input to
37//! settle. Executors must not call back into the same lease.
38//!
39//! Execution state: every lease carries an [`ExecutionIdentity`] and an
40//! [`ExecutionState`]; new input is accepted only while the execution is
41//! `Running` (in-memory leases start `Running`, so behaviour is unchanged).
42//!
43//! Durability (optional, [`InputLease::open`]): every state change is saved
44//! to a [`LeaseStore`] before it takes effect. An epoch is only granted after
45//! it is durable, and an input sequence is reserved durably before its
46//! executor runs, so a crash at any point leaves either the old state or an
47//! honest unacknowledged range. Recovery semantics are in
48//! [`crate::lease_store`]. If saving the post-execute acknowledgement fails,
49//! the live lease keeps the acknowledgement and the durable record stays
50//! conservative (that input recovers as `Unknown`);
51//! [`InputLease::store_degraded`] reports it until the next successful save.
52use std::panic::{catch_unwind, AssertUnwindSafe};
53use std::sync::atomic::{AtomicBool, Ordering};
54use std::sync::{Arc, Mutex};
55
56use crate::admission::{lock, CancelToken, EffectGate, EffectRequest, ExecError, SettleOutcome};
57use crate::events::{emit, ControlEvent, EventSink};
58pub use crate::lease_store::{
59    ExecutionIdentity, ExecutionState, FileLeaseStore, LeaseRecord, LeaseStore, MemoryLeaseStore,
60    StoreError,
61};
62
63/// Point-in-time lease state.
64#[derive(Debug, Clone, Copy, PartialEq, Eq)]
65pub struct LeaseSnapshot {
66    /// Current epoch; `0` means never acquired.
67    pub epoch: u64,
68    pub attached: bool,
69    pub next_input_sequence: u64,
70    /// Exclusive: every sequence below this is acknowledged.
71    pub acked_input_sequence: u64,
72    /// `Some((from, to))` when `[from, to)` is unacknowledged.
73    pub unacknowledged_input: Option<(u64, u64)>,
74}
75
76impl LeaseSnapshot {
77    /// CodeRight-shaped `last_acked_input_sequence` (highest acked sequence).
78    pub fn last_acked_input_sequence(&self) -> Option<u64> {
79        self.acked_input_sequence.checked_sub(1)
80    }
81}
82
83/// What a controller receives when it acquires the lease.
84#[derive(Debug, Clone, Copy, PartialEq, Eq)]
85pub struct LeaseGrant {
86    pub epoch: u64,
87    /// First sequence this holder must send.
88    pub next_input_sequence: u64,
89    /// Epoch of the attached holder this grant fenced, if any.
90    pub fenced_epoch: Option<u64>,
91}
92
93/// Explicit decision about an unacknowledged range (CodeRight
94/// `UnacknowledgedInputDecision`). There is no implicit option.
95#[derive(Debug, Clone, Copy, PartialEq, Eq)]
96pub enum UnacknowledgedInputDecision {
97    /// Caller confirmed delivery (e.g. observed the effect): acknowledge it.
98    ReconcileAsDelivered,
99    /// Caller treats it as lost: roll the cursor back so the input is resent
100    /// under fresh, unambiguous sequence numbers.
101    ReconcileAsLost,
102}
103
104/// Acknowledgement for one input sequence.
105#[derive(Debug, Clone, Copy, PartialEq, Eq)]
106pub struct InputAck {
107    pub epoch: u64,
108    pub sequence: u64,
109}
110
111#[derive(Debug, PartialEq, Eq)]
112pub enum InputOutcome<T> {
113    /// Executed once and acknowledged.
114    Applied { ack: InputAck, value: T },
115    /// Already acknowledged; not executed again.
116    Duplicate { ack: InputAck },
117}
118
119impl<T> InputOutcome<T> {
120    pub fn ack(&self) -> InputAck {
121        match self {
122            Self::Applied { ack, .. } | Self::Duplicate { ack } => *ack,
123        }
124    }
125}
126
127/// Fencing / sequencing refusal. Nothing executed.
128#[derive(Debug, Clone, PartialEq, Eq)]
129pub enum LeaseError {
130    NotAttached,
131    StaleEpoch {
132        presented: u64,
133        current: u64,
134    },
135    FutureEpoch {
136        presented: u64,
137        current: u64,
138    },
139    UnacknowledgedInput {
140        from: u64,
141        to: u64,
142    },
143    OutOfOrder {
144        expected: u64,
145        got: u64,
146    },
147    /// New input while the execution is not `Running`.
148    ExecutionNotRunning {
149        state: ExecutionState,
150    },
151    /// `acknowledge_input` named a sequence that was never sent.
152    UnsentSequence {
153        sequence: u64,
154        next: u64,
155    },
156    /// `acknowledge_input` went backwards past an acknowledged sequence.
157    NonMonotonicAck {
158        sequence: u64,
159        acked: u64,
160    },
161    /// `set_execution_state` attempted a transition out of `Exited`.
162    InvalidExecutionTransition {
163        from: ExecutionState,
164        to: ExecutionState,
165    },
166    /// The durable store refused the change; nothing changed.
167    Store {
168        reason: String,
169    },
170}
171
172impl From<StoreError> for LeaseError {
173    fn from(e: StoreError) -> Self {
174        Self::Store {
175            reason: e.to_string(),
176        }
177    }
178}
179
180impl LeaseError {
181    pub fn label(&self) -> &'static str {
182        match self {
183            Self::NotAttached => "not_attached",
184            Self::StaleEpoch { .. } => "stale_epoch",
185            Self::FutureEpoch { .. } => "future_epoch",
186            Self::UnacknowledgedInput { .. } => "unacknowledged_input",
187            Self::OutOfOrder { .. } => "out_of_order",
188            Self::ExecutionNotRunning { .. } => "execution_not_running",
189            Self::UnsentSequence { .. } => "unsent_sequence",
190            Self::NonMonotonicAck { .. } => "non_monotonic_ack",
191            Self::InvalidExecutionTransition { .. } => "invalid_execution_transition",
192            Self::Store { .. } => "store",
193        }
194    }
195}
196
197impl std::fmt::Display for LeaseError {
198    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
199        match self {
200            Self::NotAttached => write!(f, "input lease is not held; acquire before input"),
201            Self::StaleEpoch { presented, current } => {
202                write!(f, "stale input lease epoch {presented} (current {current})")
203            }
204            Self::FutureEpoch { presented, current } => {
205                write!(f, "unknown input lease epoch {presented} (current {current})")
206            }
207            Self::UnacknowledgedInput { from, to } => write!(
208                f,
209                "input lease has unacknowledged input in sequence range [{from}, {to}); reconcile it explicitly before new input"
210            ),
211            Self::OutOfOrder { expected, got } => {
212                write!(f, "input sequence must be {expected}, got {got}")
213            }
214            Self::ExecutionNotRunning { state } => {
215                write!(f, "execution is {}; input requires running", state.label())
216            }
217            Self::UnsentSequence { sequence, next } => write!(
218                f,
219                "cannot acknowledge unsent input sequence {sequence} (next {next})"
220            ),
221            Self::NonMonotonicAck { sequence, acked } => write!(
222                f,
223                "input acknowledgement must be monotonic: {sequence} is below acknowledged {acked}"
224            ),
225            Self::InvalidExecutionTransition { from, to } => write!(
226                f,
227                "execution state cannot change from {} to {}",
228                from.label(),
229                to.label()
230            ),
231            Self::Store { reason } => write!(f, "{reason}"),
232        }
233    }
234}
235
236impl std::error::Error for LeaseError {}
237
238/// How an executor failed, from the lease's point of view.
239#[derive(Debug, Clone, PartialEq, Eq)]
240pub enum InputFailure {
241    /// Definitely no side effect (e.g. denied by admission): the sequence is
242    /// released and may be reused.
243    NotDelivered(String),
244    /// A side effect may have happened: the sequence stays unacknowledged
245    /// until explicitly reconciled.
246    Ambiguous(String),
247}
248
249#[derive(Debug, Clone, PartialEq, Eq)]
250pub enum InputError {
251    Lease(LeaseError),
252    NotDelivered {
253        reason: String,
254    },
255    Ambiguous {
256        reason: String,
257        unacknowledged: (u64, u64),
258    },
259}
260
261impl std::fmt::Display for InputError {
262    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
263        match self {
264            Self::Lease(e) => e.fmt(f),
265            Self::NotDelivered { reason } => write!(f, "input not delivered: {reason}"),
266            Self::Ambiguous {
267                reason,
268                unacknowledged: (a, b),
269            } => write!(
270                f,
271                "input delivery ambiguous ({reason}); sequence range [{a}, {b}) is unacknowledged"
272            ),
273        }
274    }
275}
276
277impl std::error::Error for InputError {}
278
279impl From<LeaseError> for InputError {
280    fn from(e: LeaseError) -> Self {
281        Self::Lease(e)
282    }
283}
284
285/// Delivery status of one input sequence.
286#[derive(Debug, Clone, Copy, PartialEq, Eq)]
287pub enum InputStatus {
288    Acknowledged,
289    /// Delivery ambiguous in this host run; reconcile explicitly.
290    Unacknowledged,
291    /// Was unacknowledged when the host crashed/restarted; resolve explicitly.
292    Unknown,
293    /// Never sent (at or beyond `next_input_sequence`).
294    Unsent,
295}
296
297#[derive(Debug, Clone)]
298struct State {
299    epoch: u64,
300    attached: bool,
301    next: u64,
302    acked: u64,
303    exec_id: ExecutionIdentity,
304    exec_state: ExecutionState,
305    /// The unacknowledged range was inherited from a crash/restart.
306    unknown: bool,
307}
308
309impl State {
310    fn fresh(exec_id: ExecutionIdentity) -> Self {
311        Self {
312            epoch: 0,
313            attached: false,
314            next: 0,
315            acked: 0,
316            exec_id,
317            exec_state: ExecutionState::Running,
318            unknown: false,
319        }
320    }
321    fn recovered(r: LeaseRecord) -> Self {
322        let exec_state = match r.execution_state {
323            ExecutionState::Running => ExecutionState::Unknown,
324            other => other,
325        };
326        Self {
327            // Recovery fence: every pre-restart holder is now stale.
328            epoch: r.epoch.saturating_add(1),
329            attached: false,
330            next: r.next_input_sequence,
331            acked: r.acked_input_sequence,
332            exec_id: r.execution_id,
333            exec_state,
334            unknown: r.next_input_sequence > r.acked_input_sequence,
335        }
336    }
337    fn normalize(&mut self) {
338        if self.unacked().is_none() {
339            self.unknown = false;
340        }
341    }
342    fn record(&self, name: &str) -> LeaseRecord {
343        LeaseRecord {
344            name: name.to_string(),
345            execution_id: self.exec_id.clone(),
346            execution_state: self.exec_state,
347            epoch: self.epoch,
348            holder_epoch: self.attached.then_some(self.epoch),
349            next_input_sequence: self.next,
350            acked_input_sequence: self.acked,
351            unknown_input: self.unacked().filter(|_| self.unknown),
352        }
353    }
354    fn check_epoch(&self, epoch: u64) -> Result<(), LeaseError> {
355        if epoch < self.epoch {
356            return Err(LeaseError::StaleEpoch {
357                presented: epoch,
358                current: self.epoch,
359            });
360        }
361        if epoch > self.epoch {
362            return Err(LeaseError::FutureEpoch {
363                presented: epoch,
364                current: self.epoch,
365            });
366        }
367        Ok(())
368    }
369    fn unacked(&self) -> Option<(u64, u64)> {
370        (self.next > self.acked).then_some((self.acked, self.next))
371    }
372    fn snapshot(&self) -> LeaseSnapshot {
373        LeaseSnapshot {
374            epoch: self.epoch,
375            attached: self.attached,
376            next_input_sequence: self.next,
377            acked_input_sequence: self.acked,
378            unacknowledged_input: self.unacked(),
379        }
380    }
381    fn check_holder(&self, epoch: u64) -> Result<(), LeaseError> {
382        if !self.attached {
383            return Err(LeaseError::NotAttached);
384        }
385        self.check_epoch(epoch)
386    }
387}
388
389/// One exclusive input lease (per PTY, window, or control target).
390pub struct InputLease {
391    name: String,
392    state: Mutex<State>,
393    events: Option<EventSink>,
394    store: Option<Arc<dyn LeaseStore>>,
395    degraded: AtomicBool,
396}
397
398impl std::fmt::Debug for InputLease {
399    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
400        f.debug_struct("InputLease")
401            .field("name", &self.name)
402            .field("state", &self.snapshot())
403            .field("durable", &self.store.is_some())
404            .finish()
405    }
406}
407
408impl InputLease {
409    /// In-memory lease (the default): nothing survives the process.
410    pub fn new(name: impl Into<String>) -> Self {
411        Self::from_state(
412            name.into(),
413            State::fresh(ExecutionIdentity::generate()),
414            None,
415        )
416    }
417
418    /// Durable lease backed by `store`. A first open creates a fresh record
419    /// (new [`ExecutionIdentity`], `Running`, epoch 0); a later open recovers
420    /// it (see [`crate::lease_store`]): epoch fenced forward, detached,
421    /// `Running` becomes `Unknown`, unacknowledged input becomes `Unknown`
422    /// input. Fails closed on an unreadable/corrupt record, a record for a
423    /// different lease name, or a failed save.
424    pub fn open(name: impl Into<String>, store: Arc<dyn LeaseStore>) -> Result<Self, LeaseError> {
425        Self::open_inner(name.into(), store, None)
426    }
427
428    /// As [`Self::open`], using `identity` when no record exists yet. An
429    /// existing record keeps its own identity; a different one is refused.
430    pub fn open_with_identity(
431        name: impl Into<String>,
432        store: Arc<dyn LeaseStore>,
433        identity: ExecutionIdentity,
434    ) -> Result<Self, LeaseError> {
435        Self::open_inner(name.into(), store, Some(identity))
436    }
437
438    fn open_inner(
439        name: String,
440        store: Arc<dyn LeaseStore>,
441        identity: Option<ExecutionIdentity>,
442    ) -> Result<Self, LeaseError> {
443        let state = match store.load()? {
444            None => State::fresh(identity.unwrap_or_else(ExecutionIdentity::generate)),
445            Some(r) => {
446                if r.name != name {
447                    return Err(StoreError::Invalid(format!(
448                        "record belongs to lease {:?}, not {name:?}",
449                        r.name
450                    ))
451                    .into());
452                }
453                if let Some(id) = identity.filter(|id| *id != r.execution_id) {
454                    return Err(StoreError::Invalid(format!(
455                        "record holds execution {}, not {id}",
456                        r.execution_id
457                    ))
458                    .into());
459                }
460                State::recovered(r)
461            }
462        };
463        store.save(&state.record(&name))?;
464        Ok(Self::from_state(name, state, Some(store)))
465    }
466
467    fn from_state(name: String, state: State, store: Option<Arc<dyn LeaseStore>>) -> Self {
468        Self {
469            name,
470            state: Mutex::new(state),
471            events: None,
472            store,
473            degraded: AtomicBool::new(false),
474        }
475    }
476
477    pub fn with_events(mut self, sink: EventSink) -> Self {
478        self.events = Some(sink);
479        self
480    }
481    pub fn name(&self) -> &str {
482        &self.name
483    }
484    pub fn snapshot(&self) -> LeaseSnapshot {
485        lock(&self.state).snapshot()
486    }
487    /// True when backed by a [`LeaseStore`].
488    pub fn is_durable(&self) -> bool {
489        self.store.is_some()
490    }
491    /// True when the last best-effort save (after an executor ran) failed;
492    /// cleared by the next successful save.
493    pub fn store_degraded(&self) -> bool {
494        self.degraded.load(Ordering::SeqCst)
495    }
496    /// The durable projection (identity, execution state, epoch, holder,
497    /// sequences, unknown input), whether or not a store backs this lease.
498    pub fn record(&self) -> LeaseRecord {
499        lock(&self.state).record(&self.name)
500    }
501    pub fn execution_identity(&self) -> ExecutionIdentity {
502        lock(&self.state).exec_id.clone()
503    }
504    pub fn execution_state(&self) -> ExecutionState {
505        lock(&self.state).exec_state
506    }
507    /// `[from, to)` inherited unacknowledged from a crash/restart, if any.
508    pub fn unknown_input(&self) -> Option<(u64, u64)> {
509        let s = lock(&self.state);
510        s.unacked().filter(|_| s.unknown)
511    }
512    pub fn input_status(&self, sequence: u64) -> InputStatus {
513        let s = lock(&self.state);
514        if sequence < s.acked {
515            InputStatus::Acknowledged
516        } else if sequence >= s.next {
517            InputStatus::Unsent
518        } else if s.unknown {
519            InputStatus::Unknown
520        } else {
521            InputStatus::Unacknowledged
522        }
523    }
524
525    /// Host-side, explicit execution-state resolution (e.g. after a restart
526    /// the host verified the target is alive → `Running`, or observed it
527    /// ended → `Exited`). `Exited` is terminal.
528    pub fn set_execution_state(&self, state: ExecutionState) -> Result<LeaseSnapshot, LeaseError> {
529        self.mutate(|s| {
530            if s.exec_state == ExecutionState::Exited && state != ExecutionState::Exited {
531                return Err(LeaseError::InvalidExecutionTransition {
532                    from: s.exec_state,
533                    to: state,
534                });
535            }
536            s.exec_state = state;
537            Ok(s.snapshot())
538        })
539    }
540
541    /// Apply `f` to a copy of the state, save it, then commit. A failed `f`
542    /// or save leaves the lease unchanged.
543    fn mutate<R>(
544        &self,
545        f: impl FnOnce(&mut State) -> Result<R, LeaseError>,
546    ) -> Result<R, LeaseError> {
547        let mut guard = lock(&self.state);
548        let mut next = guard.clone();
549        let out = f(&mut next)?;
550        next.normalize();
551        self.persist(&next)?;
552        *guard = next;
553        Ok(out)
554    }
555
556    fn persist(&self, s: &State) -> Result<(), LeaseError> {
557        if let Some(store) = &self.store {
558            store.save(&s.record(&self.name))?;
559            self.degraded.store(false, Ordering::SeqCst);
560        }
561        Ok(())
562    }
563
564    fn persist_best_effort(&self, s: &State) {
565        if self.persist(s).is_err() {
566            self.degraded.store(true, Ordering::SeqCst);
567        }
568    }
569
570    /// Take the lease: bump the epoch (fencing any attached holder). Rejected
571    /// while unacknowledged input is outstanding; use
572    /// [`Self::acquire_reconciling`]. On a durable lease the new epoch is
573    /// saved before it is granted.
574    pub fn acquire(&self) -> Result<LeaseGrant, LeaseError> {
575        let (grant, fenced) = self.mutate(|s| {
576            if let Some((from, to)) = s.unacked() {
577                return Err(LeaseError::UnacknowledgedInput { from, to });
578            }
579            Ok(take(s))
580        })?;
581        self.announce(grant, fenced);
582        Ok(grant)
583    }
584
585    /// Resolve any unacknowledged (or crash-`Unknown`) range with an explicit
586    /// decision, then take the lease (bumping the epoch).
587    ///
588    /// # Panics
589    /// On a durable lease whose store refuses the save (an epoch that is not
590    /// durable is never granted). Durable callers should use
591    /// [`Self::try_acquire_reconciling`]; an in-memory lease never panics.
592    pub fn acquire_reconciling(&self, decision: UnacknowledgedInputDecision) -> LeaseGrant {
593        match self.try_acquire_reconciling(decision) {
594            Ok(g) => g,
595            Err(e) => panic!("input lease {:?}: {e}", self.name),
596        }
597    }
598
599    /// Fallible [`Self::acquire_reconciling`]: errs only when the store
600    /// refuses the save, in which case nothing changed.
601    pub fn try_acquire_reconciling(
602        &self,
603        decision: UnacknowledgedInputDecision,
604    ) -> Result<LeaseGrant, LeaseError> {
605        let (grant, fenced) = self.mutate(|s| {
606            reconcile_state(s, decision);
607            Ok(take(s))
608        })?;
609        self.announce(grant, fenced);
610        Ok(grant)
611    }
612
613    fn announce(&self, grant: LeaseGrant, fenced: Option<u64>) {
614        if let Some(old) = fenced {
615            emit(
616                &self.events,
617                ControlEvent::LeaseFenced {
618                    lease: self.name.clone(),
619                    old_epoch: old,
620                    new_epoch: grant.epoch,
621                },
622            );
623        }
624        emit(
625            &self.events,
626            ControlEvent::LeaseAcquired {
627                lease: self.name.clone(),
628                epoch: grant.epoch,
629            },
630        );
631    }
632
633    /// The current holder resolves its own unacknowledged range in place.
634    pub fn reconcile(
635        &self,
636        epoch: u64,
637        decision: UnacknowledgedInputDecision,
638    ) -> Result<LeaseSnapshot, LeaseError> {
639        self.mutate(|s| {
640            s.check_holder(epoch)?;
641            reconcile_state(s, decision);
642            Ok(s.snapshot())
643        })
644    }
645
646    /// Explicitly acknowledge delivery of every sequence up to and including
647    /// `sequence` (CodeRight `acknowledge_input`). This is how a client that
648    /// observed the effect resolves unacknowledged or crash-`Unknown` input
649    /// one prefix at a time. `epoch` must be the current epoch; the lease
650    /// need not be attached (after a restart it is detached at the recovery
651    /// epoch, which only the host can read). Re-acknowledging the highest
652    /// acknowledged sequence is idempotent.
653    pub fn acknowledge_input(
654        &self,
655        epoch: u64,
656        sequence: u64,
657    ) -> Result<LeaseSnapshot, LeaseError> {
658        self.mutate(|s| {
659            s.check_epoch(epoch)?;
660            if sequence >= s.next {
661                return Err(LeaseError::UnsentSequence {
662                    sequence,
663                    next: s.next,
664                });
665            }
666            if sequence.saturating_add(1) < s.acked {
667                return Err(LeaseError::NonMonotonicAck {
668                    sequence,
669                    acked: s.acked,
670                });
671            }
672            s.acked = s.acked.max(sequence + 1);
673            Ok(s.snapshot())
674        })
675    }
676
677    /// Release by the current holder only; a fenced holder cannot release.
678    pub fn release(&self, epoch: u64) -> Result<LeaseSnapshot, LeaseError> {
679        let snap = self.mutate(|s| {
680            s.check_holder(epoch)?;
681            s.attached = false;
682            Ok(s.snapshot())
683        })?;
684        emit(
685            &self.events,
686            ControlEvent::LeaseReleased {
687                lease: self.name.clone(),
688                epoch,
689            },
690        );
691        Ok(snap)
692    }
693
694    /// Check fencing and sequencing without executing. Returns `Ok(true)` for
695    /// new input at `next_input_sequence`, `Ok(false)` for a duplicate.
696    pub fn check(&self, epoch: u64, sequence: u64) -> Result<bool, LeaseError> {
697        let s = lock(&self.state);
698        classify(&s, epoch, sequence).map(|c| c == Class::New)
699    }
700
701    /// Submit input `(epoch, sequence)`. `execute` runs at most once, only for
702    /// new in-order input from the current holder while the execution is
703    /// `Running`. On a durable lease the sequence is reserved durably before
704    /// `execute` runs; if that save fails nothing runs (`NotDelivered`).
705    pub fn submit<T>(
706        &self,
707        epoch: u64,
708        sequence: u64,
709        execute: impl FnOnce() -> Result<T, InputFailure>,
710    ) -> Result<InputOutcome<T>, InputError> {
711        let mut s = lock(&self.state);
712        match classify(&s, epoch, sequence) {
713            Err(e) => {
714                emit(
715                    &self.events,
716                    ControlEvent::InputRejected {
717                        lease: self.name.clone(),
718                        epoch,
719                        sequence,
720                        reason: e.label(),
721                    },
722                );
723                return Err(e.into());
724            }
725            Ok(Class::Duplicate) => {
726                return Ok(InputOutcome::Duplicate {
727                    ack: InputAck { epoch, sequence },
728                })
729            }
730            Ok(Class::New) => {}
731        }
732        // Reserve the sequence before the side effect: if we crash mid-way the
733        // range is visibly unacknowledged, never silently reusable.
734        s.next = sequence.saturating_add(1);
735        if let Err(e) = self.persist(&s) {
736            s.next = sequence;
737            return Err(InputError::NotDelivered {
738                reason: e.to_string(),
739            });
740        }
741        let result = match catch_unwind(AssertUnwindSafe(execute)) {
742            Ok(Ok(value)) => {
743                s.acked = s.next;
744                Ok(InputOutcome::Applied {
745                    ack: InputAck { epoch, sequence },
746                    value,
747                })
748            }
749            Ok(Err(InputFailure::NotDelivered(reason))) => {
750                s.next = sequence;
751                Err(InputError::NotDelivered { reason })
752            }
753            Ok(Err(InputFailure::Ambiguous(reason))) => Err(InputError::Ambiguous {
754                reason,
755                unacknowledged: (s.acked, s.next),
756            }),
757            Err(_) => Err(InputError::Ambiguous {
758                reason: "input executor panicked".into(),
759                unacknowledged: (s.acked, s.next),
760            }),
761        };
762        s.normalize();
763        self.persist_best_effort(&s);
764        result
765    }
766}
767
768fn take(s: &mut State) -> (LeaseGrant, Option<u64>) {
769    let fenced = s.attached.then_some(s.epoch);
770    s.epoch = s.epoch.saturating_add(1);
771    s.attached = true;
772    (
773        LeaseGrant {
774            epoch: s.epoch,
775            next_input_sequence: s.next,
776            fenced_epoch: fenced,
777        },
778        fenced,
779    )
780}
781
782#[derive(PartialEq, Eq)]
783enum Class {
784    New,
785    Duplicate,
786}
787
788fn classify(s: &State, epoch: u64, sequence: u64) -> Result<Class, LeaseError> {
789    s.check_holder(epoch)?;
790    if sequence < s.acked {
791        return Ok(Class::Duplicate);
792    }
793    if s.exec_state != ExecutionState::Running {
794        return Err(LeaseError::ExecutionNotRunning {
795            state: s.exec_state,
796        });
797    }
798    if let Some((from, to)) = s.unacked() {
799        return Err(LeaseError::UnacknowledgedInput { from, to });
800    }
801    if sequence != s.next {
802        return Err(LeaseError::OutOfOrder {
803            expected: s.next,
804            got: sequence,
805        });
806    }
807    Ok(Class::New)
808}
809
810fn reconcile_state(s: &mut State, decision: UnacknowledgedInputDecision) {
811    if s.unacked().is_some() {
812        match decision {
813            UnacknowledgedInputDecision::ReconcileAsDelivered => s.acked = s.next,
814            UnacknowledgedInputDecision::ReconcileAsLost => s.next = s.acked,
815        }
816    }
817    s.normalize();
818}
819
820/// A lease-fenced, admission-gated control session: input must first pass
821/// PTY-002 fencing/sequencing (a stale controller never even reaches the
822/// hook), then EFF-001 admission, then executes and settles.
823#[derive(Debug, Clone)]
824pub struct ControlSession {
825    pub lease: Arc<InputLease>,
826    pub gate: Arc<EffectGate>,
827}
828
829impl ControlSession {
830    pub fn new(lease: Arc<InputLease>, gate: Arc<EffectGate>) -> Self {
831        Self { lease, gate }
832    }
833
834    /// Lease-checked, admission-gated input. Denials and pre-execute
835    /// cancellation release the sequence (`NotDelivered`); a failure, panic
836    /// or cancellation after the executor started leaves it unacknowledged
837    /// (`Ambiguous`).
838    pub fn input<T>(
839        &self,
840        epoch: u64,
841        sequence: u64,
842        request: EffectRequest,
843        cancel: &CancelToken,
844        execute: impl FnOnce(&CancelToken) -> Result<T, ExecError>,
845    ) -> Result<InputOutcome<T>, InputError> {
846        self.lease.submit(epoch, sequence, || {
847            let effect = self.gate.admit(request, cancel, execute);
848            let st = &effect.settlement;
849            let reason = format!(
850                "{}: {}",
851                st.outcome.label(),
852                st.reason.clone().unwrap_or_default()
853            );
854            match (st.outcome, st.executed) {
855                (SettleOutcome::Ok, _) => effect.value.ok_or(InputFailure::Ambiguous(reason)),
856                (_, false) => Err(InputFailure::NotDelivered(reason)),
857                (_, true) => Err(InputFailure::Ambiguous(reason)),
858            }
859        })
860    }
861}