Skip to main content

sloop/domain/
trigger.rs

1//! The trigger: the durable record that demand for work exists.
2//!
3//! A trigger is what makes the dispatcher pick a ticket up. This module owns
4//! the pure half of the concept: its kinds, its states, whether one is due,
5//! the arithmetic that rearms a recurring one, and the transition every write
6//! is derived from. Storage lives behind the coordination boundary in
7//! `work_state::trigger`; nothing here reads a clock, a database, or a file.
8
9/// What kind of demand a trigger records. The kind is also what decides how
10/// the trigger becomes due and what firing it does next.
11#[derive(Debug, Clone, Copy, PartialEq, Eq)]
12pub enum TriggerKind {
13    /// `sloop run` — due the moment the gates allow it.
14    Immediate,
15    /// `sloop post --auto` — the ticket asked for itself; no schedule.
16    Auto,
17    /// `sloop run --at HH:MM` — due once, at or after an instant.
18    At,
19    /// `sloop run --every <interval>` — due repeatedly, rearming per fire.
20    Every,
21    /// `sloop run --overnight` — due once, at or after the next opening of
22    /// running hours.
23    Overnight,
24}
25
26impl TriggerKind {
27    pub fn as_str(self) -> &'static str {
28        match self {
29            Self::Immediate => "immediate",
30            Self::Auto => "auto",
31            Self::At => "at",
32            Self::Every => "every",
33            Self::Overnight => "overnight",
34        }
35    }
36
37    pub fn parse(raw: &str) -> Option<Self> {
38        match raw {
39            "immediate" => Some(Self::Immediate),
40            "auto" => Some(Self::Auto),
41            "at" => Some(Self::At),
42            "every" => Some(Self::Every),
43            "overnight" => Some(Self::Overnight),
44            _ => None,
45        }
46    }
47
48    /// Whether this kind carries no schedule, so a queued trigger is due on
49    /// sight. Kinds that do carry one are due only once their instant passes.
50    pub fn fires_on_sight(self) -> bool {
51        matches!(self, Self::Immediate | Self::Auto)
52    }
53
54    /// Whether firing rearms the trigger rather than retiring it.
55    pub fn recurs(self) -> bool {
56        matches!(self, Self::Every)
57    }
58}
59
60/// A trigger's lifecycle state. `queued` is demand not yet met; the other two
61/// are terminal and never dispatch.
62#[derive(Debug, Clone, Copy, PartialEq, Eq)]
63pub enum TriggerState {
64    Queued,
65    Completed,
66    Cancelled,
67}
68
69impl TriggerState {
70    pub fn as_str(self) -> &'static str {
71        match self {
72            Self::Queued => "queued",
73            Self::Completed => "completed",
74            Self::Cancelled => "cancelled",
75        }
76    }
77
78    pub fn parse(raw: &str) -> Option<Self> {
79        match raw {
80            "queued" => Some(Self::Queued),
81            "completed" => Some(Self::Completed),
82            "cancelled" => Some(Self::Cancelled),
83            _ => None,
84        }
85    }
86}
87
88/// A trigger's schedulable state: everything due-ness and [`step`] need, and
89/// nothing about which ticket or project the demand points at. Targeting is a
90/// storage-side join; whether the demand is live is a decision.
91#[derive(Debug, Clone, Copy, PartialEq, Eq)]
92pub struct Trigger {
93    pub state: TriggerState,
94    pub kind: TriggerKind,
95    pub eligible_at_ms: Option<i64>,
96    pub interval_ms: Option<i64>,
97}
98
99impl Trigger {
100    /// The definition of due-ness. The SQL scan in `work_state::trigger`
101    /// carries a predicate that mirrors this so a large queue does not have to
102    /// be read into memory to be filtered, but this function is the
103    /// definition, and a test asserts the two agree over the whole matrix of
104    /// kinds, states, and schedules.
105    ///
106    /// A schedule-carrying kind with no `eligible_at_ms` is never due: an
107    /// unscheduled `at` is a corrupt row, and answering "due" would fire it
108    /// immediately, which is the opposite of what it asked for. That matches
109    /// SQL, where `NULL <= ?` is not true.
110    pub fn is_due(&self, now_ms: i64) -> bool {
111        self.state == TriggerState::Queued
112            && (self.kind.fires_on_sight()
113                || self
114                    .eligible_at_ms
115                    .is_some_and(|eligible_at_ms| eligible_at_ms <= now_ms))
116    }
117}
118
119/// What happened to a trigger. Events are evidence the daemon observed, never
120/// an intent to write a particular row.
121#[derive(Debug, Clone, Copy, PartialEq, Eq)]
122pub enum Event {
123    /// A claim consumed the trigger.
124    Fired,
125    /// The demand can never be met — the ticket it is pinned to has merged.
126    Completed,
127    /// A run gave the ticket back and the trigger returns to the queue, re-timed
128    /// to the retry's earliest instant.
129    Requeued { eligible_at_ms: i64 },
130}
131
132/// What a transition asks storage to persist. One variant per named write in
133/// `work_state::trigger`; there is no effect for "leave it alone", so an empty
134/// result means the event was a no-op.
135#[derive(Debug, Clone, Copy, PartialEq, Eq)]
136pub enum Effect {
137    /// Stay queued, with a new eligibility instant. Only a recurring trigger
138    /// rearms.
139    Rearm { eligible_at_ms: i64 },
140    /// Retire the trigger.
141    Complete,
142    /// Return the trigger to the queue at this instant.
143    Requeue { eligible_at_ms: i64 },
144    /// The trigger cannot be honoured and no write should follow.
145    Fault(Fault),
146}
147
148/// A trigger whose own fields contradict its kind. Faults are data problems,
149/// not race outcomes, so they surface as corruption rather than denial.
150#[derive(Debug, Clone, Copy, PartialEq, Eq)]
151pub enum Fault {
152    /// A recurring trigger with no usable `(eligible_at_ms, interval_ms)`
153    /// pair: firing it could not compute a next instant, so it would either
154    /// spin or silently stop recurring.
155    InvalidCadence,
156}
157
158/// The trigger transition. `trigger` is advanced in place and the returned
159/// effects are what storage must persist, in order.
160///
161/// Every terminal state is absorbing: replaying an event against a completed
162/// or cancelled trigger yields no effects, which is what makes the recovery
163/// and sweep paths idempotent.
164pub fn step(trigger: &mut Trigger, event: Event, now_ms: i64) -> Vec<Effect> {
165    match event {
166        Event::Fired => {
167            if trigger.state != TriggerState::Queued {
168                return Vec::new();
169            }
170            if !trigger.kind.recurs() {
171                trigger.state = TriggerState::Completed;
172                return vec![Effect::Complete];
173            }
174            let rearmed = trigger.eligible_at_ms.zip(trigger.interval_ms).and_then(
175                |(eligible_at_ms, interval_ms)| rearm_every_at(eligible_at_ms, interval_ms, now_ms),
176            );
177            match rearmed {
178                Some(eligible_at_ms) => {
179                    trigger.eligible_at_ms = Some(eligible_at_ms);
180                    vec![Effect::Rearm { eligible_at_ms }]
181                }
182                None => vec![Effect::Fault(Fault::InvalidCadence)],
183            }
184        }
185        Event::Completed => {
186            if trigger.state != TriggerState::Queued {
187                return Vec::new();
188            }
189            trigger.state = TriggerState::Completed;
190            vec![Effect::Complete]
191        }
192        Event::Requeued { eligible_at_ms } => {
193            if trigger.state == TriggerState::Cancelled {
194                return Vec::new();
195            }
196            trigger.state = TriggerState::Queued;
197            trigger.eligible_at_ms = Some(eligible_at_ms);
198            vec![Effect::Requeue { eligible_at_ms }]
199        }
200    }
201}
202
203/// The next instant a recurring trigger becomes due, given the instant it was
204/// due and its cadence.
205///
206/// Missed intervals collapse into one step rather than replaying: a daemon
207/// asleep for an hour on a one-minute cadence owes one run, not sixty. The
208/// result is always strictly after `now_ms`, so a rearm can never leave the
209/// trigger due again in the same tick.
210pub fn rearm_every_at(eligible_at_ms: i64, interval_ms: i64, now_ms: i64) -> Option<i64> {
211    if interval_ms <= 0 || eligible_at_ms > now_ms {
212        return None;
213    }
214    let missed = now_ms.checked_sub(eligible_at_ms)?.div_euclid(interval_ms);
215    let steps = missed.checked_add(1)?;
216    eligible_at_ms.checked_add(interval_ms.checked_mul(steps)?)
217}
218
219#[cfg(test)]
220mod tests {
221    use super::*;
222
223    fn queued(kind: TriggerKind, eligible_at_ms: Option<i64>) -> Trigger {
224        Trigger {
225            state: TriggerState::Queued,
226            kind,
227            eligible_at_ms,
228            interval_ms: None,
229        }
230    }
231
232    #[test]
233    fn kinds_and_states_round_trip_through_their_stored_spelling() {
234        for kind in [
235            TriggerKind::Immediate,
236            TriggerKind::Auto,
237            TriggerKind::At,
238            TriggerKind::Every,
239            TriggerKind::Overnight,
240        ] {
241            assert_eq!(TriggerKind::parse(kind.as_str()), Some(kind));
242        }
243        for state in [
244            TriggerState::Queued,
245            TriggerState::Completed,
246            TriggerState::Cancelled,
247        ] {
248            assert_eq!(TriggerState::parse(state.as_str()), Some(state));
249        }
250        assert_eq!(TriggerKind::parse("eventually"), None);
251        assert_eq!(TriggerState::parse("pending"), None);
252    }
253
254    #[test]
255    fn scheduleless_kinds_are_due_on_sight_and_scheduled_kinds_wait() {
256        assert!(queued(TriggerKind::Immediate, None).is_due(1_000));
257        assert!(queued(TriggerKind::Auto, None).is_due(1_000));
258        assert!(queued(TriggerKind::Immediate, Some(9_000)).is_due(1_000));
259
260        for kind in [TriggerKind::At, TriggerKind::Every, TriggerKind::Overnight] {
261            assert!(!queued(kind, Some(1_001)).is_due(1_000), "{kind:?} early");
262            assert!(queued(kind, Some(1_000)).is_due(1_000), "{kind:?} on time");
263            assert!(queued(kind, Some(999)).is_due(1_000), "{kind:?} late");
264            assert!(!queued(kind, None).is_due(1_000), "{kind:?} unscheduled");
265        }
266    }
267
268    #[test]
269    fn only_queued_triggers_are_due() {
270        for state in [TriggerState::Completed, TriggerState::Cancelled] {
271            let trigger = Trigger {
272                state,
273                kind: TriggerKind::Immediate,
274                eligible_at_ms: None,
275                interval_ms: None,
276            };
277            assert!(!trigger.is_due(1_000), "{state:?}");
278        }
279    }
280
281    #[test]
282    fn firing_retires_a_one_shot_trigger() {
283        for kind in [
284            TriggerKind::Immediate,
285            TriggerKind::Auto,
286            TriggerKind::At,
287            TriggerKind::Overnight,
288        ] {
289            let mut trigger = queued(kind, Some(500));
290            assert_eq!(
291                step(&mut trigger, Event::Fired, 1_000),
292                [Effect::Complete],
293                "{kind:?}"
294            );
295            assert_eq!(trigger.state, TriggerState::Completed);
296            assert!(!trigger.is_due(1_000));
297        }
298    }
299
300    #[test]
301    fn firing_rearms_a_recurring_trigger_past_now() {
302        let mut trigger = Trigger {
303            state: TriggerState::Queued,
304            kind: TriggerKind::Every,
305            eligible_at_ms: Some(1_000),
306            interval_ms: Some(60_000),
307        };
308        assert_eq!(
309            step(&mut trigger, Event::Fired, 1_000),
310            [Effect::Rearm {
311                eligible_at_ms: 61_000
312            }]
313        );
314        assert_eq!(trigger.state, TriggerState::Queued);
315        assert_eq!(trigger.eligible_at_ms, Some(61_000));
316        assert!(!trigger.is_due(1_000));
317    }
318
319    #[test]
320    fn a_recurring_trigger_without_a_cadence_faults_instead_of_writing() {
321        for interval_ms in [None, Some(0), Some(-1)] {
322            let mut trigger = Trigger {
323                state: TriggerState::Queued,
324                kind: TriggerKind::Every,
325                eligible_at_ms: Some(1_000),
326                interval_ms,
327            };
328            assert_eq!(
329                step(&mut trigger, Event::Fired, 1_000),
330                [Effect::Fault(Fault::InvalidCadence)],
331                "{interval_ms:?}"
332            );
333            assert_eq!(trigger.state, TriggerState::Queued);
334            assert_eq!(trigger.eligible_at_ms, Some(1_000));
335        }
336    }
337
338    #[test]
339    fn completion_ignores_kind_and_terminal_states_absorb_every_event() {
340        let mut recurring = Trigger {
341            state: TriggerState::Queued,
342            kind: TriggerKind::Every,
343            eligible_at_ms: Some(1_000),
344            interval_ms: Some(60_000),
345        };
346        assert_eq!(
347            step(&mut recurring, Event::Completed, 2_000),
348            [Effect::Complete]
349        );
350        assert_eq!(recurring.state, TriggerState::Completed);
351        assert_eq!(step(&mut recurring, Event::Completed, 3_000), []);
352        assert_eq!(step(&mut recurring, Event::Fired, 3_000), []);
353    }
354
355    #[test]
356    fn requeueing_revives_a_completed_trigger_at_the_retry_instant() {
357        let mut trigger = Trigger {
358            state: TriggerState::Completed,
359            kind: TriggerKind::Immediate,
360            eligible_at_ms: None,
361            interval_ms: None,
362        };
363        assert_eq!(
364            step(
365                &mut trigger,
366                Event::Requeued {
367                    eligible_at_ms: 5_000
368                },
369                2_000
370            ),
371            [Effect::Requeue {
372                eligible_at_ms: 5_000
373            }]
374        );
375        assert_eq!(trigger.state, TriggerState::Queued);
376        assert!(trigger.is_due(2_000));
377
378        let mut cancelled = Trigger {
379            state: TriggerState::Cancelled,
380            ..trigger
381        };
382        assert_eq!(
383            step(
384                &mut cancelled,
385                Event::Requeued {
386                    eligible_at_ms: 5_000
387                },
388                2_000
389            ),
390            []
391        );
392        assert_eq!(cancelled.state, TriggerState::Cancelled);
393    }
394
395    #[test]
396    fn rearm_collapses_missed_intervals_into_one_step() {
397        assert_eq!(rearm_every_at(1_000, 60_000, 1_000), Some(61_000));
398        assert_eq!(rearm_every_at(1_000, 60_000, 60_999), Some(61_000));
399        assert_eq!(rearm_every_at(1_000, 60_000, 61_000), Some(121_000));
400        assert_eq!(rearm_every_at(1_000, 60_000, 3_601_000), Some(3_661_000));
401    }
402
403    #[test]
404    fn rearm_refuses_impossible_arithmetic() {
405        assert_eq!(rearm_every_at(2_000, 60_000, 1_000), None);
406        assert_eq!(rearm_every_at(1_000, 0, 2_000), None);
407        assert_eq!(rearm_every_at(1_000, -60_000, 2_000), None);
408        assert_eq!(rearm_every_at(0, i64::MAX, i64::MAX), None);
409    }
410}