Skip to main content

behavior_actors/time/
lease.rs

1//! Expiring exclusive ownership over explicit timer generations.
2
3use 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
15/// Complete exclusive lease phase.
16pub enum LeaseState<K, Route> {
17    /// No holder owns the lease.
18    Vacant {
19        /// Generation committed by the next successful acquisition.
20        next: TimerGeneration,
21    },
22    /// One holder owns exactly this timer generation.
23    Held {
24        /// Application-defined holder identity.
25        holder: K,
26        /// Current expiry generation.
27        generation: TimerGeneration,
28        /// Recipient notified of expiry.
29        notify: Route,
30        /// Generation available to a renewal or later acquisition. `None`
31        /// means this held generation is the final representable one.
32        next: Option<TimerGeneration>,
33    },
34    /// No fresh expiry generation is representable and no holder exists.
35    Exhausted,
36}
37
38/// Typed lease-operation rejection.
39#[derive(Debug, Error, Clone, PartialEq, Eq)]
40pub enum LeaseRejection<K> {
41    /// Another holder currently owns the lease.
42    #[error("lease is occupied")]
43    Occupied {
44        /// Holder that currently owns the lease.
45        current: K,
46    },
47    /// Operation named a holder other than the current holder.
48    #[error("lease holder does not match")]
49    WrongHolder {
50        /// Holder that currently owns the lease.
51        current: K,
52    },
53    /// Operation names an old or future generation.
54    #[error("lease generation is stale")]
55    StaleGeneration {
56        /// Rejected generation.
57        observed: TimerGeneration,
58        /// Current generation.
59        current: TimerGeneration,
60    },
61    /// No holder currently exists.
62    #[error("lease is vacant")]
63    Vacant,
64    /// The finite generation domain has been consumed.
65    #[error("lease generation domain is exhausted")]
66    GenerationExhausted,
67}
68
69/// Complete lease operation retained when admission is rejected.
70#[derive(Debug, Clone, PartialEq, Eq)]
71pub enum LeaseRequest<K> {
72    /// Attempt to acquire a vacant lease.
73    Acquire {
74        /// Requested holder.
75        holder: K,
76        /// Requested lease lifetime.
77        duration: Duration,
78    },
79    /// Attempt to renew one exact held generation.
80    Renew {
81        /// Claimed holder.
82        holder: K,
83        /// Claimed current generation.
84        generation: TimerGeneration,
85        /// Requested renewed lifetime.
86        duration: Duration,
87    },
88    /// Attempt to release one exact held generation.
89    Release {
90        /// Claimed holder.
91        holder: K,
92        /// Claimed current generation.
93        generation: TimerGeneration,
94    },
95}
96
97/// Factual lease operation result.
98#[derive(Debug, Clone, PartialEq, Eq)]
99pub enum LeaseOutcome<K> {
100    /// Fresh ownership was committed and expiry requested.
101    Acquired {
102        /// Holder.
103        holder: K,
104        /// Committed generation.
105        generation: TimerGeneration,
106    },
107    /// Existing ownership was renewed into a fresh generation.
108    Renewed {
109        /// Holder.
110        holder: K,
111        /// Fresh committed generation.
112        generation: TimerGeneration,
113    },
114    /// Ownership was explicitly relinquished.
115    Released {
116        /// Former holder.
117        holder: K,
118        /// Released generation.
119        generation: TimerGeneration,
120    },
121    /// Matching timer evidence expired ownership.
122    Expired {
123        /// Former holder.
124        holder: K,
125        /// Expired generation.
126        generation: TimerGeneration,
127    },
128    /// Operation was rejected without changing state.
129    Rejected {
130        /// Complete operation that was not admitted.
131        request: LeaseRequest<K>,
132        /// Exhaustive reason derived from current lease state.
133        reason: LeaseRejection<K>,
134    },
135}
136
137/// User commands accepted by [`Lease`].
138pub enum LeaseMessage<K, Route> {
139    /// Acquire a vacant lease for a relative duration.
140    Acquire {
141        /// Requested holder.
142        holder: K,
143        /// Relative duration interpreted by Timers.
144        duration: Duration,
145        /// Outcome and expiry recipient.
146        reply_to: Route,
147    },
148    /// Renew the exact current incarnation.
149    Renew {
150        /// Claimed holder.
151        holder: K,
152        /// Claimed current generation.
153        generation: TimerGeneration,
154        /// New relative duration.
155        duration: Duration,
156        /// Outcome and expiry recipient.
157        reply_to: Route,
158    },
159    /// Explicitly release the exact current incarnation.
160    Release {
161        /// Claimed holder.
162        holder: K,
163        /// Claimed current generation.
164        generation: TimerGeneration,
165        /// Outcome recipient.
166        reply_to: Route,
167    },
168}
169
170/// Named effect lanes emitted by [`Lease`].
171#[derive(behavior_macros::SendProduct)]
172pub struct LeaseSends<OutcomeSends, Schedules> {
173    /// Lease facts.
174    pub outcomes: OutcomeSends,
175    /// Relative expiry requests.
176    pub schedules: Schedules,
177}
178
179/// Exclusive expiring ownership behavior.
180///
181/// Acquisition is accepted only while vacant and commits a fresh generation
182/// before emitting its schedule. Renewal and release require the exact holder
183/// and generation. Matching elapsed evidence expires ownership; wrong timer IDs
184/// and stale generations are inert. Release cannot retract an already queued
185/// elapsed observation, so later evidence is explicitly stale. Once the finite
186/// generation domain is consumed, the vacant state becomes `Exhausted` and no
187/// successful ownership can be created. Initialization is empty, no actors are
188/// created, and the host never terminates by policy. Exclusivity, checked
189/// generation progression, and commit-before-schedule ordering are Bombay
190/// policy. Scheduling and sleeping belong to Timers. The current Bombay timer
191/// adapter has no cancellation effect lane; release is semantically immediate
192/// while the stale queue entry may remain until due. No transition panics.
193pub 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    /// Construct a vacant lease using one actor-local timer key.
211    #[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    /// Borrow the complete ownership phase.
222    #[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}