1use std::time::Duration;
4
5#[cfg(test)]
6use behavior::Recipient;
7use behavior::{
8 Actions, Address, Behavior, BehaviorActed, BehaviorBase, EventLayer, InterpreterRequests,
9 Never, NoBirths, SendEffects, User,
10};
11use thiserror::Error;
12
13use crate::{DeliveryRoute, ScheduleAfter, TimedEvent, TimerGeneration, TimerId};
14
15pub enum LeaseState<K, Route> {
17 Vacant {
19 next: TimerGeneration,
21 },
22 Held {
24 holder: K,
26 generation: TimerGeneration,
28 notify: Route,
30 next: Option<TimerGeneration>,
33 },
34 Exhausted,
36}
37
38#[derive(Debug, Error, Clone, PartialEq, Eq)]
40pub enum LeaseRejection<K> {
41 #[error("lease is occupied")]
43 Occupied {
44 current: K,
46 },
47 #[error("lease holder does not match")]
49 WrongHolder {
50 current: K,
52 },
53 #[error("lease generation is stale")]
55 StaleGeneration {
56 observed: TimerGeneration,
58 current: TimerGeneration,
60 },
61 #[error("lease is vacant")]
63 Vacant,
64 #[error("lease generation domain is exhausted")]
66 GenerationExhausted,
67}
68
69#[derive(Debug, Clone, PartialEq, Eq)]
71pub enum LeaseRequest<K> {
72 Acquire {
74 holder: K,
76 duration: Duration,
78 },
79 Renew {
81 holder: K,
83 generation: TimerGeneration,
85 duration: Duration,
87 },
88 Release {
90 holder: K,
92 generation: TimerGeneration,
94 },
95}
96
97#[derive(Debug, Clone, PartialEq, Eq)]
99pub enum LeaseOutcome<K> {
100 Acquired {
102 holder: K,
104 generation: TimerGeneration,
106 },
107 Renewed {
109 holder: K,
111 generation: TimerGeneration,
113 },
114 Released {
116 holder: K,
118 generation: TimerGeneration,
120 },
121 Expired {
123 holder: K,
125 generation: TimerGeneration,
127 },
128 Rejected {
130 request: LeaseRequest<K>,
132 reason: LeaseRejection<K>,
134 },
135}
136
137pub enum LeaseMessage<K, Route> {
139 Acquire {
141 holder: K,
143 duration: Duration,
145 reply_to: Route,
147 },
148 Renew {
150 holder: K,
152 generation: TimerGeneration,
154 duration: Duration,
156 reply_to: Route,
158 },
159 Release {
161 holder: K,
163 generation: TimerGeneration,
165 reply_to: Route,
167 },
168}
169
170#[derive(behavior_macros::SendProduct)]
172pub struct LeaseSends<OutcomeSends, Schedules> {
173 pub outcomes: OutcomeSends,
175 pub schedules: Schedules,
177}
178
179pub struct Lease<
194 A: Address,
195 K,
196 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = LeaseOutcome<K>>>,
197> {
198 id: TimerId,
199 state: LeaseState<K, Route>,
200 marker: core::marker::PhantomData<fn() -> A>,
201}
202type LeaseActions<A, OutcomeSends> =
203 Actions<A, Never, LeaseSends<OutcomeSends, InterpreterRequests<ScheduleAfter>>, NoBirths>;
204impl<A, K, Route> Lease<A, K, Route>
205where
206 A: Address,
207 K: Clone + Eq,
208 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = LeaseOutcome<K>>> + Clone,
209{
210 #[must_use]
212 pub const fn new(id: TimerId) -> Self {
213 Self {
214 id,
215 state: LeaseState::Vacant {
216 next: TimerGeneration(0),
217 },
218 marker: core::marker::PhantomData,
219 }
220 }
221 #[must_use]
223 pub const fn state(&self) -> &LeaseState<K, Route> {
224 &self.state
225 }
226 fn result(reply_to: Route, outcome: LeaseOutcome<K>) -> LeaseActions<A, Route::Sends> {
227 Actions::send(LeaseSends {
228 outcomes: reply_to.deliver(outcome),
229 schedules: InterpreterRequests::empty(),
230 })
231 }
232 fn acquire(
233 &mut self,
234 holder: K,
235 duration: Duration,
236 reply_to: Route,
237 ) -> LeaseActions<A, Route::Sends> {
238 let next = match &self.state {
239 LeaseState::Held {
240 holder: current, ..
241 } => {
242 return Self::result(
243 reply_to,
244 LeaseOutcome::Rejected {
245 request: LeaseRequest::Acquire { holder, duration },
246 reason: LeaseRejection::Occupied {
247 current: current.clone(),
248 },
249 },
250 );
251 }
252 LeaseState::Exhausted => {
253 return Self::result(
254 reply_to,
255 LeaseOutcome::Rejected {
256 request: LeaseRequest::Acquire { holder, duration },
257 reason: LeaseRejection::GenerationExhausted,
258 },
259 );
260 }
261 LeaseState::Vacant { next } => *next,
262 };
263 let successor = next.0.checked_add(1).map(TimerGeneration);
264 self.state = LeaseState::Held {
265 holder: holder.clone(),
266 generation: next,
267 notify: reply_to.clone(),
268 next: successor,
269 };
270 Actions::send(LeaseSends {
271 outcomes: reply_to.deliver(LeaseOutcome::Acquired {
272 holder,
273 generation: next,
274 }),
275 schedules: InterpreterRequests::one(ScheduleAfter::new(self.id, next, duration)),
276 })
277 }
278 fn successor(
279 &self,
280 holder: &K,
281 generation: TimerGeneration,
282 ) -> Result<Option<TimerGeneration>, LeaseRejection<K>> {
283 match &self.state {
284 LeaseState::Vacant { .. } => Err(LeaseRejection::Vacant),
285 LeaseState::Exhausted => Err(LeaseRejection::GenerationExhausted),
286 LeaseState::Held {
287 holder: current, ..
288 } if current != holder => Err(LeaseRejection::WrongHolder {
289 current: current.clone(),
290 }),
291 LeaseState::Held {
292 generation: active, ..
293 } if *active != generation => Err(LeaseRejection::StaleGeneration {
294 observed: generation,
295 current: *active,
296 }),
297 LeaseState::Held { next, .. } => Ok(*next),
298 }
299 }
300}
301impl<A, K, Route> BehaviorBase for Lease<A, K, Route>
302where
303 A: Address,
304 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = LeaseOutcome<K>>>,
305{
306 type Base = Self;
307 fn base(&self) -> &Self {
308 self
309 }
310}
311impl<A, K, Route> behavior::Protocol for Lease<A, K, Route>
312where
313 A: Address,
314 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = LeaseOutcome<K>>>,
315{
316 type Addr = A;
317 type Msg = LeaseMessage<K, Route>;
318}
319
320impl<A, K, Route> Behavior for Lease<A, K, Route>
321where
322 A: Address,
323 K: Clone + Eq,
324 Route: DeliveryRoute<Protocol: behavior::Protocol<Addr = A, Msg = LeaseOutcome<K>>> + Clone,
325 Route::Sends: behavior::SendsFor<TimedEvent<User<A, LeaseMessage<K, Route>>>>,
326{
327 type Protocol = Self;
328 type Event = TimedEvent<User<A, behavior::BehaviorMessage<Self>>>;
329 type Sends = LeaseSends<Route::Sends, InterpreterRequests<ScheduleAfter>>;
330 type Ph = Never;
331 type Error = Never;
332 type Birth = NoBirths;
333 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
334 Ok(match event {
335 EventLayer::Inner(event) => match event.message {
336 LeaseMessage::Acquire {
337 holder,
338 duration,
339 reply_to,
340 } => self.acquire(holder, duration, reply_to),
341 LeaseMessage::Renew {
342 holder,
343 generation,
344 duration,
345 reply_to,
346 } => {
347 let request = LeaseRequest::Renew {
348 holder: holder.clone(),
349 generation,
350 duration,
351 };
352 let successor = match self.successor(&holder, generation) {
353 Ok(successor) => successor,
354 Err(rejection) => {
355 return Ok(Self::result(
356 reply_to,
357 LeaseOutcome::Rejected {
358 request,
359 reason: rejection,
360 },
361 ));
362 }
363 };
364 let Some(fresh) = successor else {
365 return Ok(Self::result(
366 reply_to,
367 LeaseOutcome::Rejected {
368 request,
369 reason: LeaseRejection::GenerationExhausted,
370 },
371 ));
372 };
373 let next = fresh.0.checked_add(1).map(TimerGeneration);
374 self.state = LeaseState::Held {
375 holder: holder.clone(),
376 generation: fresh,
377 notify: reply_to.clone(),
378 next,
379 };
380 Actions::send(LeaseSends {
381 outcomes: reply_to.deliver(LeaseOutcome::Renewed {
382 holder,
383 generation: fresh,
384 }),
385 schedules: InterpreterRequests::one(ScheduleAfter::new(
386 self.id, fresh, duration,
387 )),
388 })
389 }
390 LeaseMessage::Release {
391 holder,
392 generation,
393 reply_to,
394 } => {
395 let request = LeaseRequest::Release {
396 holder: holder.clone(),
397 generation,
398 };
399 let next = match self.successor(&holder, generation) {
400 Ok(next) => next,
401 Err(rejection) => {
402 return Ok(Self::result(
403 reply_to,
404 LeaseOutcome::Rejected {
405 request,
406 reason: rejection,
407 },
408 ));
409 }
410 };
411 self.state =
412 next.map_or(LeaseState::Exhausted, |next| LeaseState::Vacant { next });
413 Self::result(reply_to, LeaseOutcome::Released { holder, generation })
414 }
415 },
416 EventLayer::Owned(elapsed) => {
417 if elapsed.id != self.id {
418 return Ok(Actions::cont());
419 }
420 let LeaseState::Held {
421 holder,
422 generation,
423 notify,
424 next,
425 } = &self.state
426 else {
427 return Ok(Actions::cont());
428 };
429 if elapsed.generation != *generation {
430 return Ok(Actions::cont());
431 }
432 let outcome = LeaseOutcome::Expired {
433 holder: holder.clone(),
434 generation: *generation,
435 };
436 let recipient = notify.clone();
437 let next = *next;
438 self.state = next.map_or(LeaseState::Exhausted, |next| LeaseState::Vacant { next });
439 Self::result(recipient, outcome)
440 }
441 })
442 }
443}
444
445#[cfg(test)]
446mod tests {
447 use super::*;
448 use crate::{Activate as _, TimerElapsed};
449 use behavior::MailAddr;
450 struct Reply;
451 impl behavior::Protocol for Reply {
452 type Addr = MailAddr;
453 type Msg = LeaseOutcome<u8>;
454 }
455
456 impl Behavior for Reply {
457 type Protocol = Self;
458 type Event = User<MailAddr, behavior::BehaviorMessage<Self>>;
459 type Sends = Vec<Never>;
460 type Ph = Never;
461 type Error = Never;
462 type Birth = NoBirths;
463 fn transition(&mut self, _: behavior::ActiveTurn, _: Self::Event) -> BehaviorActed<Self> {
464 Ok(Actions::cont())
465 }
466 }
467 type Subject = Lease<MailAddr, u8, Recipient<Reply>>;
468 fn reply() -> Recipient<Reply> {
469 Recipient::global(MailAddr(1))
470 }
471 fn duration() -> Duration {
472 Duration::from_secs(1)
473 }
474 #[test]
475 fn acquire_renew_release_and_stale_elapsed_are_generation_safe() {
476 let mut s = (Subject::new(TimerId(7))).initialize().unwrap().behavior;
477 let acquired = s
478 .receive(
479 MailAddr(0),
480 LeaseMessage::Acquire {
481 holder: 1,
482 duration: duration(),
483 reply_to: reply(),
484 },
485 )
486 .unwrap();
487 assert_eq!(
488 acquired.sends.schedules.as_slice()[0].generation,
489 TimerGeneration(0)
490 );
491 assert_eq!(acquired.sends.outcomes.len(), 1);
492 assert!(acquired.creates.is_empty());
493 assert_eq!(acquired.become_, behavior::Step::Continue);
494 let renewed = s
495 .receive(
496 MailAddr(0),
497 LeaseMessage::Renew {
498 holder: 1,
499 generation: TimerGeneration(0),
500 duration: duration(),
501 reply_to: reply(),
502 },
503 )
504 .unwrap();
505 assert_eq!(
506 renewed.sends.schedules.as_slice()[0].generation,
507 TimerGeneration(1)
508 );
509 assert_eq!(renewed.sends.outcomes.len(), 1);
510 assert!(renewed.creates.is_empty());
511 assert_eq!(renewed.become_, behavior::Step::Continue);
512 let stale = s
513 .on_path(TimerElapsed::new(TimerId(7), TimerGeneration(0)))
514 .unwrap();
515 assert!(stale.sends.outcomes.is_empty());
516 assert!(stale.sends.schedules.is_empty());
517 assert!(stale.creates.is_empty());
518 assert_eq!(stale.become_, behavior::Step::Continue);
519 let released = s
520 .receive(
521 MailAddr(0),
522 LeaseMessage::Release {
523 holder: 1,
524 generation: TimerGeneration(1),
525 reply_to: reply(),
526 },
527 )
528 .unwrap();
529 assert!(released.sends.schedules.is_empty());
530 assert_eq!(released.sends.outcomes.len(), 1);
531 assert!(matches!(
532 released.sends.outcomes[0].message,
533 LeaseOutcome::Released { holder: 1, .. }
534 ));
535 assert!(released.creates.is_empty());
536 assert_eq!(released.become_, behavior::Step::Continue);
537 assert!(matches!(s.state(), LeaseState::Vacant { .. }));
538 }
539 #[test]
540 fn wrong_holder_and_matching_expiry_are_distinct() {
541 let mut s = (Subject::new(TimerId(7))).initialize().unwrap().behavior;
542 let acquired = s
543 .receive(
544 MailAddr(0),
545 LeaseMessage::Acquire {
546 holder: 1,
547 duration: duration(),
548 reply_to: reply(),
549 },
550 )
551 .unwrap();
552 assert_eq!(acquired.sends.schedules.len(), 1);
553 assert_eq!(acquired.sends.outcomes.len(), 1);
554 assert!(acquired.creates.is_empty());
555 assert_eq!(acquired.become_, behavior::Step::Continue);
556 let wrong = s
557 .receive(
558 MailAddr(0),
559 LeaseMessage::Release {
560 holder: 2,
561 generation: TimerGeneration(0),
562 reply_to: reply(),
563 },
564 )
565 .unwrap();
566 assert!(matches!(
567 wrong.sends.outcomes[0].message,
568 LeaseOutcome::Rejected {
569 request: LeaseRequest::Release {
570 holder: 2,
571 generation: TimerGeneration(0),
572 },
573 reason: LeaseRejection::WrongHolder { current: 1 },
574 }
575 ));
576 let expired = s
577 .on_path(TimerElapsed::new(TimerId(7), TimerGeneration(0)))
578 .unwrap();
579 assert!(matches!(
580 expired.sends.outcomes[0].message,
581 LeaseOutcome::Expired { holder: 1, .. }
582 ));
583 assert!(matches!(s.state(), LeaseState::Vacant { .. }));
584 }
585
586 #[test]
587 fn generation_exhaustion_is_terminal_and_never_wraps() {
588 let mut definition = Subject::new(TimerId(7));
589 definition.state = LeaseState::Vacant {
590 next: TimerGeneration(u64::MAX),
591 };
592 let mut subject = (definition).initialize().unwrap().behavior;
593 let acquired = subject
594 .receive(
595 MailAddr(0),
596 LeaseMessage::Acquire {
597 holder: 1,
598 duration: duration(),
599 reply_to: reply(),
600 },
601 )
602 .unwrap();
603 assert_eq!(acquired.sends.schedules.len(), 1);
604 assert_eq!(acquired.sends.outcomes.len(), 1);
605 assert!(acquired.creates.is_empty());
606 assert_eq!(acquired.become_, behavior::Step::Continue);
607 let expired = subject
608 .on_path(TimerElapsed::new(TimerId(7), TimerGeneration(u64::MAX)))
609 .unwrap();
610 assert!(expired.sends.schedules.is_empty());
611 assert_eq!(expired.sends.outcomes.len(), 1);
612 assert!(expired.creates.is_empty());
613 assert_eq!(expired.become_, behavior::Step::Continue);
614 assert!(matches!(subject.state(), LeaseState::Exhausted));
615 let rejected = subject
616 .receive(
617 MailAddr(0),
618 LeaseMessage::Acquire {
619 holder: 2,
620 duration: duration(),
621 reply_to: reply(),
622 },
623 )
624 .unwrap();
625 assert!(matches!(
626 rejected.sends.outcomes[0].message,
627 LeaseOutcome::Rejected {
628 request: LeaseRequest::Acquire {
629 holder: 2,
630 duration: requested_duration,
631 },
632 reason: LeaseRejection::GenerationExhausted,
633 } if requested_duration == duration()
634 ));
635 }
636}