1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
65pub struct LeaseSnapshot {
66 pub epoch: u64,
68 pub attached: bool,
69 pub next_input_sequence: u64,
70 pub acked_input_sequence: u64,
72 pub unacknowledged_input: Option<(u64, u64)>,
74}
75
76impl LeaseSnapshot {
77 pub fn last_acked_input_sequence(&self) -> Option<u64> {
79 self.acked_input_sequence.checked_sub(1)
80 }
81}
82
83#[derive(Debug, Clone, Copy, PartialEq, Eq)]
85pub struct LeaseGrant {
86 pub epoch: u64,
87 pub next_input_sequence: u64,
89 pub fenced_epoch: Option<u64>,
91}
92
93#[derive(Debug, Clone, Copy, PartialEq, Eq)]
96pub enum UnacknowledgedInputDecision {
97 ReconcileAsDelivered,
99 ReconcileAsLost,
102}
103
104#[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 Applied { ack: InputAck, value: T },
115 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#[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 ExecutionNotRunning {
149 state: ExecutionState,
150 },
151 UnsentSequence {
153 sequence: u64,
154 next: u64,
155 },
156 NonMonotonicAck {
158 sequence: u64,
159 acked: u64,
160 },
161 InvalidExecutionTransition {
163 from: ExecutionState,
164 to: ExecutionState,
165 },
166 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#[derive(Debug, Clone, PartialEq, Eq)]
240pub enum InputFailure {
241 NotDelivered(String),
244 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
287pub enum InputStatus {
288 Acknowledged,
289 Unacknowledged,
291 Unknown,
293 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 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 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
389pub 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 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 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 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 pub fn is_durable(&self) -> bool {
489 self.store.is_some()
490 }
491 pub fn store_degraded(&self) -> bool {
494 self.degraded.load(Ordering::SeqCst)
495 }
496 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 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 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 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 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 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 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 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 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 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 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 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 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#[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 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}