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        // Kind is deliberately not consulted. A recurring trigger pinned to a
186        // merged ticket is as unfireable as a one-shot one, and leaving it
187        // queued is demand that can never be met but is still counted.
188        Event::Completed => {
189            if trigger.state != TriggerState::Queued {
190                return Vec::new();
191            }
192            trigger.state = TriggerState::Completed;
193            vec![Effect::Complete]
194        }
195        Event::Requeued { eligible_at_ms } => {
196            if trigger.state == TriggerState::Cancelled {
197                return Vec::new();
198            }
199            trigger.state = TriggerState::Queued;
200            trigger.eligible_at_ms = Some(eligible_at_ms);
201            vec![Effect::Requeue { eligible_at_ms }]
202        }
203    }
204}
205
206/// The next instant a recurring trigger becomes due, given the instant it was
207/// due and its cadence.
208///
209/// Missed intervals collapse into one step rather than replaying: a daemon
210/// asleep for an hour on a one-minute cadence owes one run, not sixty. The
211/// result is always strictly after `now_ms`, so a rearm can never leave the
212/// trigger due again in the same tick.
213pub fn rearm_every_at(eligible_at_ms: i64, interval_ms: i64, now_ms: i64) -> Option<i64> {
214    if interval_ms <= 0 || eligible_at_ms > now_ms {
215        return None;
216    }
217    let missed = now_ms.checked_sub(eligible_at_ms)?.div_euclid(interval_ms);
218    let steps = missed.checked_add(1)?;
219    eligible_at_ms.checked_add(interval_ms.checked_mul(steps)?)
220}
221
222#[cfg(test)]
223mod tests {
224    use super::*;
225
226    fn queued(kind: TriggerKind, eligible_at_ms: Option<i64>) -> Trigger {
227        Trigger {
228            state: TriggerState::Queued,
229            kind,
230            eligible_at_ms,
231            interval_ms: None,
232        }
233    }
234
235    #[test]
236    fn kinds_and_states_round_trip_through_their_stored_spelling() {
237        for kind in [
238            TriggerKind::Immediate,
239            TriggerKind::Auto,
240            TriggerKind::At,
241            TriggerKind::Every,
242            TriggerKind::Overnight,
243        ] {
244            assert_eq!(TriggerKind::parse(kind.as_str()), Some(kind));
245        }
246        for state in [
247            TriggerState::Queued,
248            TriggerState::Completed,
249            TriggerState::Cancelled,
250        ] {
251            assert_eq!(TriggerState::parse(state.as_str()), Some(state));
252        }
253        assert_eq!(TriggerKind::parse("eventually"), None);
254        assert_eq!(TriggerState::parse("pending"), None);
255    }
256
257    #[test]
258    fn scheduleless_kinds_are_due_on_sight_and_scheduled_kinds_wait() {
259        assert!(queued(TriggerKind::Immediate, None).is_due(1_000));
260        assert!(queued(TriggerKind::Auto, None).is_due(1_000));
261        // A schedule on a scheduleless kind is ignored rather than obeyed.
262        assert!(queued(TriggerKind::Immediate, Some(9_000)).is_due(1_000));
263
264        for kind in [TriggerKind::At, TriggerKind::Every, TriggerKind::Overnight] {
265            assert!(!queued(kind, Some(1_001)).is_due(1_000), "{kind:?} early");
266            assert!(queued(kind, Some(1_000)).is_due(1_000), "{kind:?} on time");
267            assert!(queued(kind, Some(999)).is_due(1_000), "{kind:?} late");
268            // No schedule at all is corrupt, and firing is the wrong guess.
269            assert!(!queued(kind, None).is_due(1_000), "{kind:?} unscheduled");
270        }
271    }
272
273    #[test]
274    fn only_queued_triggers_are_due() {
275        for state in [TriggerState::Completed, TriggerState::Cancelled] {
276            let trigger = Trigger {
277                state,
278                kind: TriggerKind::Immediate,
279                eligible_at_ms: None,
280                interval_ms: None,
281            };
282            assert!(!trigger.is_due(1_000), "{state:?}");
283        }
284    }
285
286    #[test]
287    fn firing_retires_a_one_shot_trigger() {
288        for kind in [
289            TriggerKind::Immediate,
290            TriggerKind::Auto,
291            TriggerKind::At,
292            TriggerKind::Overnight,
293        ] {
294            let mut trigger = queued(kind, Some(500));
295            assert_eq!(
296                step(&mut trigger, Event::Fired, 1_000),
297                [Effect::Complete],
298                "{kind:?}"
299            );
300            assert_eq!(trigger.state, TriggerState::Completed);
301            assert!(!trigger.is_due(1_000));
302        }
303    }
304
305    #[test]
306    fn firing_rearms_a_recurring_trigger_past_now() {
307        let mut trigger = Trigger {
308            state: TriggerState::Queued,
309            kind: TriggerKind::Every,
310            eligible_at_ms: Some(1_000),
311            interval_ms: Some(60_000),
312        };
313        assert_eq!(
314            step(&mut trigger, Event::Fired, 1_000),
315            [Effect::Rearm {
316                eligible_at_ms: 61_000
317            }]
318        );
319        assert_eq!(trigger.state, TriggerState::Queued);
320        assert_eq!(trigger.eligible_at_ms, Some(61_000));
321        // The rearm is strictly in the future, so the same tick cannot fire it
322        // a second time.
323        assert!(!trigger.is_due(1_000));
324    }
325
326    #[test]
327    fn a_recurring_trigger_without_a_cadence_faults_instead_of_writing() {
328        for interval_ms in [None, Some(0), Some(-1)] {
329            let mut trigger = Trigger {
330                state: TriggerState::Queued,
331                kind: TriggerKind::Every,
332                eligible_at_ms: Some(1_000),
333                interval_ms,
334            };
335            assert_eq!(
336                step(&mut trigger, Event::Fired, 1_000),
337                [Effect::Fault(Fault::InvalidCadence)],
338                "{interval_ms:?}"
339            );
340            assert_eq!(trigger.state, TriggerState::Queued);
341            assert_eq!(trigger.eligible_at_ms, Some(1_000));
342        }
343    }
344
345    #[test]
346    fn completion_ignores_kind_and_terminal_states_absorb_every_event() {
347        let mut recurring = Trigger {
348            state: TriggerState::Queued,
349            kind: TriggerKind::Every,
350            eligible_at_ms: Some(1_000),
351            interval_ms: Some(60_000),
352        };
353        assert_eq!(
354            step(&mut recurring, Event::Completed, 2_000),
355            [Effect::Complete]
356        );
357        assert_eq!(recurring.state, TriggerState::Completed);
358        // Replaying the sweep writes nothing, which is what lets it run on
359        // every startup for free.
360        assert_eq!(step(&mut recurring, Event::Completed, 3_000), []);
361        assert_eq!(step(&mut recurring, Event::Fired, 3_000), []);
362    }
363
364    #[test]
365    fn requeueing_revives_a_completed_trigger_at_the_retry_instant() {
366        let mut trigger = Trigger {
367            state: TriggerState::Completed,
368            kind: TriggerKind::Immediate,
369            eligible_at_ms: None,
370            interval_ms: None,
371        };
372        assert_eq!(
373            step(
374                &mut trigger,
375                Event::Requeued {
376                    eligible_at_ms: 5_000
377                },
378                2_000
379            ),
380            [Effect::Requeue {
381                eligible_at_ms: 5_000
382            }]
383        );
384        assert_eq!(trigger.state, TriggerState::Queued);
385        // `immediate` ignores its schedule, so the revived trigger is due at
386        // once; the cooldown that set `not_before_ms` is a separate gate.
387        assert!(trigger.is_due(2_000));
388
389        let mut cancelled = Trigger {
390            state: TriggerState::Cancelled,
391            ..trigger
392        };
393        assert_eq!(
394            step(
395                &mut cancelled,
396                Event::Requeued {
397                    eligible_at_ms: 5_000
398                },
399                2_000
400            ),
401            []
402        );
403        assert_eq!(cancelled.state, TriggerState::Cancelled);
404    }
405
406    #[test]
407    fn rearm_collapses_missed_intervals_into_one_step() {
408        assert_eq!(rearm_every_at(1_000, 60_000, 1_000), Some(61_000));
409        assert_eq!(rearm_every_at(1_000, 60_000, 60_999), Some(61_000));
410        assert_eq!(rearm_every_at(1_000, 60_000, 61_000), Some(121_000));
411        // Asleep for an hour on a one-minute cadence owes one run, not sixty.
412        assert_eq!(rearm_every_at(1_000, 60_000, 3_601_000), Some(3_661_000));
413    }
414
415    #[test]
416    fn rearm_refuses_impossible_arithmetic() {
417        // Not yet due: rearming would skip the instant it was waiting for.
418        assert_eq!(rearm_every_at(2_000, 60_000, 1_000), None);
419        assert_eq!(rearm_every_at(1_000, 0, 2_000), None);
420        assert_eq!(rearm_every_at(1_000, -60_000, 2_000), None);
421        assert_eq!(rearm_every_at(0, i64::MAX, i64::MAX), None);
422    }
423}