1use std::panic::{catch_unwind, AssertUnwindSafe};
39use std::sync::{Arc, Mutex};
40
41use crate::admission::{lock, CancelToken, EffectGate, EffectRequest, ExecError, SettleOutcome};
42use crate::events::{emit, ControlEvent, EventSink};
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub struct LeaseSnapshot {
47 pub epoch: u64,
49 pub attached: bool,
50 pub next_input_sequence: u64,
51 pub acked_input_sequence: u64,
53 pub unacknowledged_input: Option<(u64, u64)>,
55}
56
57impl LeaseSnapshot {
58 pub fn last_acked_input_sequence(&self) -> Option<u64> {
60 self.acked_input_sequence.checked_sub(1)
61 }
62}
63
64#[derive(Debug, Clone, Copy, PartialEq, Eq)]
66pub struct LeaseGrant {
67 pub epoch: u64,
68 pub next_input_sequence: u64,
70 pub fenced_epoch: Option<u64>,
72}
73
74#[derive(Debug, Clone, Copy, PartialEq, Eq)]
77pub enum UnacknowledgedInputDecision {
78 ReconcileAsDelivered,
80 ReconcileAsLost,
83}
84
85#[derive(Debug, Clone, Copy, PartialEq, Eq)]
87pub struct InputAck {
88 pub epoch: u64,
89 pub sequence: u64,
90}
91
92#[derive(Debug, PartialEq, Eq)]
93pub enum InputOutcome<T> {
94 Applied { ack: InputAck, value: T },
96 Duplicate { ack: InputAck },
98}
99
100impl<T> InputOutcome<T> {
101 pub fn ack(&self) -> InputAck {
102 match self {
103 Self::Applied { ack, .. } | Self::Duplicate { ack } => *ack,
104 }
105 }
106}
107
108#[derive(Debug, Clone, PartialEq, Eq)]
110pub enum LeaseError {
111 NotAttached,
112 StaleEpoch { presented: u64, current: u64 },
113 FutureEpoch { presented: u64, current: u64 },
114 UnacknowledgedInput { from: u64, to: u64 },
115 OutOfOrder { expected: u64, got: u64 },
116}
117
118impl LeaseError {
119 pub fn label(&self) -> &'static str {
120 match self {
121 Self::NotAttached => "not_attached",
122 Self::StaleEpoch { .. } => "stale_epoch",
123 Self::FutureEpoch { .. } => "future_epoch",
124 Self::UnacknowledgedInput { .. } => "unacknowledged_input",
125 Self::OutOfOrder { .. } => "out_of_order",
126 }
127 }
128}
129
130impl std::fmt::Display for LeaseError {
131 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
132 match self {
133 Self::NotAttached => write!(f, "input lease is not held; acquire before input"),
134 Self::StaleEpoch { presented, current } => {
135 write!(f, "stale input lease epoch {presented} (current {current})")
136 }
137 Self::FutureEpoch { presented, current } => {
138 write!(f, "unknown input lease epoch {presented} (current {current})")
139 }
140 Self::UnacknowledgedInput { from, to } => write!(
141 f,
142 "input lease has unacknowledged input in sequence range [{from}, {to}); reconcile it explicitly before new input"
143 ),
144 Self::OutOfOrder { expected, got } => {
145 write!(f, "input sequence must be {expected}, got {got}")
146 }
147 }
148 }
149}
150
151impl std::error::Error for LeaseError {}
152
153#[derive(Debug, Clone, PartialEq, Eq)]
155pub enum InputFailure {
156 NotDelivered(String),
159 Ambiguous(String),
162}
163
164#[derive(Debug, Clone, PartialEq, Eq)]
165pub enum InputError {
166 Lease(LeaseError),
167 NotDelivered {
168 reason: String,
169 },
170 Ambiguous {
171 reason: String,
172 unacknowledged: (u64, u64),
173 },
174}
175
176impl std::fmt::Display for InputError {
177 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
178 match self {
179 Self::Lease(e) => e.fmt(f),
180 Self::NotDelivered { reason } => write!(f, "input not delivered: {reason}"),
181 Self::Ambiguous {
182 reason,
183 unacknowledged: (a, b),
184 } => write!(
185 f,
186 "input delivery ambiguous ({reason}); sequence range [{a}, {b}) is unacknowledged"
187 ),
188 }
189 }
190}
191
192impl std::error::Error for InputError {}
193
194impl From<LeaseError> for InputError {
195 fn from(e: LeaseError) -> Self {
196 Self::Lease(e)
197 }
198}
199
200#[derive(Debug, Default)]
201struct State {
202 epoch: u64,
203 attached: bool,
204 next: u64,
205 acked: u64,
206}
207
208impl State {
209 fn unacked(&self) -> Option<(u64, u64)> {
210 (self.next > self.acked).then_some((self.acked, self.next))
211 }
212 fn snapshot(&self) -> LeaseSnapshot {
213 LeaseSnapshot {
214 epoch: self.epoch,
215 attached: self.attached,
216 next_input_sequence: self.next,
217 acked_input_sequence: self.acked,
218 unacknowledged_input: self.unacked(),
219 }
220 }
221 fn check_holder(&self, epoch: u64) -> Result<(), LeaseError> {
222 if !self.attached {
223 return Err(LeaseError::NotAttached);
224 }
225 if epoch < self.epoch {
226 return Err(LeaseError::StaleEpoch {
227 presented: epoch,
228 current: self.epoch,
229 });
230 }
231 if epoch > self.epoch {
232 return Err(LeaseError::FutureEpoch {
233 presented: epoch,
234 current: self.epoch,
235 });
236 }
237 Ok(())
238 }
239}
240
241pub struct InputLease {
243 name: String,
244 state: Mutex<State>,
245 events: Option<EventSink>,
246}
247
248impl std::fmt::Debug for InputLease {
249 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
250 f.debug_struct("InputLease")
251 .field("name", &self.name)
252 .field("state", &self.snapshot())
253 .finish()
254 }
255}
256
257impl InputLease {
258 pub fn new(name: impl Into<String>) -> Self {
259 Self {
260 name: name.into(),
261 state: Mutex::new(State::default()),
262 events: None,
263 }
264 }
265 pub fn with_events(mut self, sink: EventSink) -> Self {
266 self.events = Some(sink);
267 self
268 }
269 pub fn name(&self) -> &str {
270 &self.name
271 }
272 pub fn snapshot(&self) -> LeaseSnapshot {
273 lock(&self.state).snapshot()
274 }
275
276 pub fn acquire(&self) -> Result<LeaseGrant, LeaseError> {
280 let mut s = lock(&self.state);
281 if let Some((from, to)) = s.unacked() {
282 return Err(LeaseError::UnacknowledgedInput { from, to });
283 }
284 Ok(self.take(&mut s))
285 }
286
287 pub fn acquire_reconciling(&self, decision: UnacknowledgedInputDecision) -> LeaseGrant {
290 let mut s = lock(&self.state);
291 reconcile_state(&mut s, decision);
292 self.take(&mut s)
293 }
294
295 fn take(&self, s: &mut State) -> LeaseGrant {
296 let fenced = s.attached.then_some(s.epoch);
297 s.epoch = s.epoch.saturating_add(1);
298 s.attached = true;
299 if let Some(old) = fenced {
300 emit(
301 &self.events,
302 ControlEvent::LeaseFenced {
303 lease: self.name.clone(),
304 old_epoch: old,
305 new_epoch: s.epoch,
306 },
307 );
308 }
309 emit(
310 &self.events,
311 ControlEvent::LeaseAcquired {
312 lease: self.name.clone(),
313 epoch: s.epoch,
314 },
315 );
316 LeaseGrant {
317 epoch: s.epoch,
318 next_input_sequence: s.next,
319 fenced_epoch: fenced,
320 }
321 }
322
323 pub fn reconcile(
325 &self,
326 epoch: u64,
327 decision: UnacknowledgedInputDecision,
328 ) -> Result<LeaseSnapshot, LeaseError> {
329 let mut s = lock(&self.state);
330 s.check_holder(epoch)?;
331 reconcile_state(&mut s, decision);
332 Ok(s.snapshot())
333 }
334
335 pub fn release(&self, epoch: u64) -> Result<LeaseSnapshot, LeaseError> {
337 let mut s = lock(&self.state);
338 s.check_holder(epoch)?;
339 s.attached = false;
340 emit(
341 &self.events,
342 ControlEvent::LeaseReleased {
343 lease: self.name.clone(),
344 epoch,
345 },
346 );
347 Ok(s.snapshot())
348 }
349
350 pub fn check(&self, epoch: u64, sequence: u64) -> Result<bool, LeaseError> {
353 let s = lock(&self.state);
354 classify(&s, epoch, sequence).map(|c| c == Class::New)
355 }
356
357 pub fn submit<T>(
360 &self,
361 epoch: u64,
362 sequence: u64,
363 execute: impl FnOnce() -> Result<T, InputFailure>,
364 ) -> Result<InputOutcome<T>, InputError> {
365 let mut s = lock(&self.state);
366 match classify(&s, epoch, sequence) {
367 Err(e) => {
368 emit(
369 &self.events,
370 ControlEvent::InputRejected {
371 lease: self.name.clone(),
372 epoch,
373 sequence,
374 reason: e.label(),
375 },
376 );
377 return Err(e.into());
378 }
379 Ok(Class::Duplicate) => {
380 return Ok(InputOutcome::Duplicate {
381 ack: InputAck { epoch, sequence },
382 })
383 }
384 Ok(Class::New) => {}
385 }
386 s.next = sequence.saturating_add(1);
389 match catch_unwind(AssertUnwindSafe(execute)) {
390 Ok(Ok(value)) => {
391 s.acked = s.next;
392 Ok(InputOutcome::Applied {
393 ack: InputAck { epoch, sequence },
394 value,
395 })
396 }
397 Ok(Err(InputFailure::NotDelivered(reason))) => {
398 s.next = sequence;
399 Err(InputError::NotDelivered { reason })
400 }
401 Ok(Err(InputFailure::Ambiguous(reason))) => Err(InputError::Ambiguous {
402 reason,
403 unacknowledged: (s.acked, s.next),
404 }),
405 Err(_) => Err(InputError::Ambiguous {
406 reason: "input executor panicked".into(),
407 unacknowledged: (s.acked, s.next),
408 }),
409 }
410 }
411}
412
413#[derive(PartialEq, Eq)]
414enum Class {
415 New,
416 Duplicate,
417}
418
419fn classify(s: &State, epoch: u64, sequence: u64) -> Result<Class, LeaseError> {
420 s.check_holder(epoch)?;
421 if sequence < s.acked {
422 return Ok(Class::Duplicate);
423 }
424 if let Some((from, to)) = s.unacked() {
425 return Err(LeaseError::UnacknowledgedInput { from, to });
426 }
427 if sequence != s.next {
428 return Err(LeaseError::OutOfOrder {
429 expected: s.next,
430 got: sequence,
431 });
432 }
433 Ok(Class::New)
434}
435
436fn reconcile_state(s: &mut State, decision: UnacknowledgedInputDecision) {
437 if s.unacked().is_some() {
438 match decision {
439 UnacknowledgedInputDecision::ReconcileAsDelivered => s.acked = s.next,
440 UnacknowledgedInputDecision::ReconcileAsLost => s.next = s.acked,
441 }
442 }
443}
444
445#[derive(Debug, Clone)]
449pub struct ControlSession {
450 pub lease: Arc<InputLease>,
451 pub gate: Arc<EffectGate>,
452}
453
454impl ControlSession {
455 pub fn new(lease: Arc<InputLease>, gate: Arc<EffectGate>) -> Self {
456 Self { lease, gate }
457 }
458
459 pub fn input<T>(
464 &self,
465 epoch: u64,
466 sequence: u64,
467 request: EffectRequest,
468 cancel: &CancelToken,
469 execute: impl FnOnce(&CancelToken) -> Result<T, ExecError>,
470 ) -> Result<InputOutcome<T>, InputError> {
471 self.lease.submit(epoch, sequence, || {
472 let effect = self.gate.admit(request, cancel, execute);
473 let st = &effect.settlement;
474 let reason = format!(
475 "{}: {}",
476 st.outcome.label(),
477 st.reason.clone().unwrap_or_default()
478 );
479 match (st.outcome, st.executed) {
480 (SettleOutcome::Ok, _) => effect.value.ok_or(InputFailure::Ambiguous(reason)),
481 (_, false) => Err(InputFailure::NotDelivered(reason)),
482 (_, true) => Err(InputFailure::Ambiguous(reason)),
483 }
484 })
485 }
486}