Skip to main content

runledger_runtime/
reaper.rs

1use runledger_postgres::jobs::ReapedLeaseRecord;
2use std::sync::Arc;
3use tokio::sync::watch;
4use tracing::{info, warn};
5
6mod observers;
7mod terminal_hooks;
8
9use self::observers::ReapedObserverTasks;
10use self::terminal_hooks::notify_handlers_of_terminal_lease_expirations;
11use crate::ReaperError;
12use crate::RuntimeLoopExit;
13use crate::config::JobsConfig;
14use crate::observer::JobLifecycleObservers;
15use crate::registry::JobRegistry;
16use crate::shutdown;
17
18pub async fn run_reaper_loop(
19    pool: runledger_postgres::DbPool,
20    registry: JobRegistry,
21    config: JobsConfig,
22    shutdown: watch::Receiver<bool>,
23) -> RuntimeLoopExit {
24    run_reaper_loop_with_observer(
25        pool,
26        registry,
27        config,
28        shutdown,
29        JobLifecycleObservers::empty(),
30    )
31    .await
32}
33
34pub async fn run_reaper_loop_with_observer(
35    pool: runledger_postgres::DbPool,
36    registry: JobRegistry,
37    config: JobsConfig,
38    mut shutdown: watch::Receiver<bool>,
39    observers: JobLifecycleObservers,
40) -> RuntimeLoopExit {
41    if let Err(error) = config.validate_reaper_loop() {
42        warn!(%error, "invalid jobs config; stopping reaper loop");
43        return RuntimeLoopExit::InvalidConfig(error);
44    }
45
46    let registry = Arc::new(registry);
47    let mut reaped_observer_tasks = ReapedObserverTasks::owned();
48
49    loop {
50        reaped_observer_tasks.drain_finished();
51
52        if shutdown::is_requested_or_closed(&shutdown) {
53            return reaper_shutdown_complete(&mut reaped_observer_tasks).await;
54        }
55
56        match runledger_postgres::jobs::reap_expired_leases_with_diagnostics(
57            &pool,
58            config.claim_batch_size,
59            config.reaper_retry_delay_ms,
60        )
61        .await
62        {
63            Ok(result) => {
64                if result.summary.processed > 0 {
65                    info!(
66                        reaped = result.summary.processed,
67                        "reaper reclaimed expired leases"
68                    );
69                }
70                log_deferred_row_errors(&result);
71
72                let fanout_result = notify_reaped_lease_side_effects(
73                    registry.as_ref(),
74                    &observers,
75                    &mut reaped_observer_tasks,
76                    &result.reaped_leases,
77                    &mut shutdown,
78                )
79                .await;
80                if matches!(
81                    fanout_result,
82                    ReaperNotificationFanoutResult::InterruptedByShutdown
83                ) {
84                    return reaper_shutdown_complete(&mut reaped_observer_tasks).await;
85                }
86            }
87            Err(error) => {
88                let error = ReaperError::ReapExpiredLeases {
89                    batch_size: config.claim_batch_size,
90                    retry_delay_ms: config.reaper_retry_delay_ms,
91                    source: error,
92                };
93                warn!(%error, "reaper iteration failed");
94            }
95        }
96
97        if shutdown::wait_for_request_or_timeout(&mut shutdown, config.reaper_interval).await {
98            return reaper_shutdown_complete(&mut reaped_observer_tasks).await;
99        }
100    }
101}
102
103async fn notify_reaped_lease_side_effects(
104    registry: &JobRegistry,
105    observers: &JobLifecycleObservers,
106    reaped_observer_tasks: &mut ReapedObserverTasks,
107    reaped_leases: &[ReapedLeaseRecord],
108    shutdown: &mut watch::Receiver<bool>,
109) -> ReaperNotificationFanoutResult {
110    let terminal_hook_result =
111        notify_handlers_of_terminal_lease_expirations(registry, reaped_leases, shutdown).await;
112
113    if terminal_hook_result.interrupted_by_shutdown() || shutdown::is_requested_or_closed(shutdown)
114    {
115        return ReaperNotificationFanoutResult::InterruptedByShutdown;
116    }
117
118    reaped_observer_tasks.spawn_batch(observers, reaped_leases);
119    ReaperNotificationFanoutResult::Completed
120}
121
122fn log_deferred_row_errors(result: &runledger_postgres::jobs::ReapExpiredLeasesDetailedResult) {
123    for error in &result.deferred_row_errors {
124        warn!(
125            job_id = %error.job_id,
126            run_number = error.run_number,
127            attempt = error.attempt,
128            error_code = %error.error_code,
129            error_message = %error.error_message,
130            error_sqlstate = error.sqlstate.as_deref().unwrap_or(""),
131            deferred_row_error_count = result.deferred_row_error_count,
132            logged_deferred_row_errors = result.deferred_row_errors.len(),
133            "reaper deferred expired leased job after row-level processing error"
134        );
135    }
136
137    if result.deferred_row_error_count > result.deferred_row_errors.len() {
138        warn!(
139            deferred_row_error_count = result.deferred_row_error_count,
140            logged_deferred_row_errors = result.deferred_row_errors.len(),
141            "reaper deferred additional expired leased jobs after row-level processing errors"
142        );
143    }
144}
145
146async fn reaper_shutdown_complete(
147    reaped_observer_tasks: &mut ReapedObserverTasks,
148) -> RuntimeLoopExit {
149    reaped_observer_tasks.abort_for_shutdown().await;
150    info!("reaper shutdown complete");
151    RuntimeLoopExit::Shutdown
152}
153
154#[derive(Debug, PartialEq, Eq)]
155enum ReaperNotificationFanoutResult {
156    Completed,
157    InterruptedByShutdown,
158}
159
160#[cfg(test)]
161mod tests {
162    use std::future::pending;
163    use std::sync::atomic::{AtomicUsize, Ordering};
164    use std::sync::{Arc, Mutex};
165    use std::time::Duration;
166
167    use chrono::Utc;
168    use runledger_core::jobs::{
169        JobCompletion, JobContext, JobDeadLetterInfo, JobDeadLetterReason, JobFailure, JobType,
170    };
171    use runledger_postgres::jobs::ReapedLeaseDisposition;
172    use runledger_postgres::jobs::test_support::{
173        reaped_lease_record as postgres_reaped_lease_record,
174        reaped_lease_record_with_checkpoint as postgres_reaped_lease_record_with_checkpoint,
175    };
176    use serde_json::{Value, json};
177    use sqlx::types::Uuid;
178    use tokio::sync::{Notify, watch};
179    use tokio::time::{sleep, timeout};
180
181    use crate::observer::{JobLeaseReapedEvent, JobLifecycleObserver, JobLifecycleObservers};
182    use crate::registry::{JobHandler, JobRegistry};
183
184    use super::observers::ReapedObserverTasks;
185    use super::terminal_hooks::{
186        TerminalHookFanoutResult, notify_handlers_of_terminal_lease_expirations,
187        notify_handlers_of_terminal_lease_expirations_with_before_first_hook_admission,
188    };
189    use super::{ReaperNotificationFanoutResult, notify_reaped_lease_side_effects};
190
191    struct HangingTerminalHookHandler;
192    struct CountingHangingTerminalHookHandler {
193        started: Arc<AtomicUsize>,
194    }
195
196    struct ControlledTerminalHookHandler {
197        started: Arc<Notify>,
198        release: Arc<Notify>,
199        completions: Arc<AtomicUsize>,
200    }
201
202    struct RecordingDeadLetterHandler {
203        dead_letters: Arc<Mutex<Vec<JobDeadLetterInfo>>>,
204        contexts: Arc<Mutex<Vec<JobContext>>>,
205    }
206
207    #[derive(Clone, Default)]
208    struct RecordingObserver {
209        reaped: Arc<Mutex<Vec<JobLeaseReapedEvent>>>,
210    }
211
212    #[derive(Clone)]
213    struct SlowReapedObserver {
214        started: Arc<Notify>,
215        started_count: Arc<AtomicUsize>,
216    }
217
218    #[derive(Clone)]
219    struct CancellableReapedObserver {
220        started: Arc<Notify>,
221        started_count: Arc<AtomicUsize>,
222        active_callbacks: Arc<AtomicUsize>,
223    }
224
225    struct ActiveCallbackGuard {
226        active_callbacks: Arc<AtomicUsize>,
227    }
228
229    impl Drop for ActiveCallbackGuard {
230        fn drop(&mut self) {
231            self.active_callbacks.fetch_sub(1, Ordering::SeqCst);
232        }
233    }
234
235    #[async_trait::async_trait]
236    impl JobHandler for HangingTerminalHookHandler {
237        fn job_type(&self) -> JobType<'static> {
238            JobType::new("jobs.test.reaper.hook.hang")
239        }
240
241        async fn execute(
242            &self,
243            _context: JobContext,
244            _payload: Value,
245        ) -> Result<JobCompletion, JobFailure> {
246            Ok(JobCompletion::success())
247        }
248
249        async fn on_dead_letter(
250            &self,
251            _context: JobContext,
252            _payload: Value,
253            _dead_letter: JobDeadLetterInfo,
254        ) {
255            pending::<()>().await;
256        }
257    }
258
259    #[async_trait::async_trait]
260    impl JobHandler for CountingHangingTerminalHookHandler {
261        fn job_type(&self) -> JobType<'static> {
262            JobType::new("jobs.test.reaper.hook.counting_hang")
263        }
264
265        async fn execute(
266            &self,
267            _context: JobContext,
268            _payload: Value,
269        ) -> Result<JobCompletion, JobFailure> {
270            Ok(JobCompletion::success())
271        }
272
273        async fn on_dead_letter(
274            &self,
275            _context: JobContext,
276            _payload: Value,
277            _dead_letter: JobDeadLetterInfo,
278        ) {
279            self.started.fetch_add(1, Ordering::SeqCst);
280            pending::<()>().await;
281        }
282    }
283
284    #[async_trait::async_trait]
285    impl JobHandler for ControlledTerminalHookHandler {
286        fn job_type(&self) -> JobType<'static> {
287            JobType::new("jobs.test.reaper.hook.controlled")
288        }
289
290        async fn execute(
291            &self,
292            _context: JobContext,
293            _payload: Value,
294        ) -> Result<JobCompletion, JobFailure> {
295            Ok(JobCompletion::success())
296        }
297
298        async fn on_dead_letter(
299            &self,
300            _context: JobContext,
301            _payload: Value,
302            _dead_letter: JobDeadLetterInfo,
303        ) {
304            self.started.notify_one();
305            self.release.notified().await;
306            self.completions.fetch_add(1, Ordering::SeqCst);
307        }
308    }
309
310    #[async_trait::async_trait]
311    impl JobHandler for RecordingDeadLetterHandler {
312        fn job_type(&self) -> JobType<'static> {
313            JobType::new("jobs.test.reaper.dead_letter.record")
314        }
315
316        async fn execute(
317            &self,
318            _context: JobContext,
319            _payload: Value,
320        ) -> Result<JobCompletion, JobFailure> {
321            Ok(JobCompletion::success())
322        }
323
324        async fn on_dead_letter(
325            &self,
326            context: JobContext,
327            _payload: Value,
328            dead_letter: JobDeadLetterInfo,
329        ) {
330            self.contexts
331                .lock()
332                .expect("dead-letter context list lock should not be poisoned")
333                .push(context);
334            self.dead_letters
335                .lock()
336                .expect("dead-letter list lock should not be poisoned")
337                .push(dead_letter);
338        }
339    }
340
341    #[async_trait::async_trait]
342    impl JobLifecycleObserver for RecordingObserver {
343        async fn on_job_lease_reaped(&self, event: JobLeaseReapedEvent) {
344            self.reaped
345                .lock()
346                .expect("reaped events lock should not be poisoned")
347                .push(event);
348        }
349    }
350
351    #[async_trait::async_trait]
352    impl JobLifecycleObserver for SlowReapedObserver {
353        async fn on_job_lease_reaped(&self, _event: JobLeaseReapedEvent) {
354            self.started_count.fetch_add(1, Ordering::SeqCst);
355            self.started.notify_waiters();
356            pending::<()>().await;
357        }
358    }
359
360    #[async_trait::async_trait]
361    impl JobLifecycleObserver for CancellableReapedObserver {
362        async fn on_job_lease_reaped(&self, _event: JobLeaseReapedEvent) {
363            self.started_count.fetch_add(1, Ordering::SeqCst);
364            self.active_callbacks.fetch_add(1, Ordering::SeqCst);
365            let _guard = ActiveCallbackGuard {
366                active_callbacks: self.active_callbacks.clone(),
367            };
368            self.started.notify_waiters();
369            pending::<()>().await;
370        }
371    }
372
373    async fn wait_for_count(counter: Arc<AtomicUsize>, notify: Arc<Notify>, expected: usize) {
374        while counter.load(Ordering::SeqCst) < expected {
375            notify.notified().await;
376        }
377    }
378
379    async fn wait_for_reaped_count(
380        observer: &RecordingObserver,
381        expected: usize,
382        timeout_after: Duration,
383    ) {
384        timeout(timeout_after, async {
385            loop {
386                let count = observer
387                    .reaped
388                    .lock()
389                    .expect("reaped events lock should not be poisoned")
390                    .len();
391                if count >= expected {
392                    break;
393                }
394                sleep(Duration::from_millis(10)).await;
395            }
396        })
397        .await
398        .expect("timed out waiting for reaped observer event count");
399    }
400
401    fn reaped_lease_record(job_type: &'static str) -> runledger_postgres::jobs::ReapedLeaseRecord {
402        postgres_reaped_lease_record(
403            Uuid::now_v7(),
404            runledger_core::jobs::JobTypeName::from_static(job_type),
405            None,
406            1,
407            1,
408            3,
409            Some("worker-reaped-observer".to_owned()),
410            false,
411            ReapedLeaseDisposition::ReleasedToPending,
412        )
413    }
414
415    fn terminal_reaped_lease_record(
416        job_id: Uuid,
417        job_type: runledger_core::jobs::JobTypeName,
418        payload: Value,
419    ) -> runledger_postgres::jobs::ReapedLeaseRecord {
420        postgres_reaped_lease_record(
421            job_id,
422            job_type.clone(),
423            None,
424            1,
425            1,
426            1,
427            Some("worker-reaped-terminal".to_owned()),
428            false,
429            ReapedLeaseDisposition::DeadLetteredTerminal { payload },
430        )
431    }
432
433    fn terminal_reaped_lease_record_with_checkpoint(
434        job_id: Uuid,
435        job_type: runledger_core::jobs::JobTypeName,
436        payload: Value,
437        checkpoint: Value,
438    ) -> runledger_postgres::jobs::ReapedLeaseRecord {
439        postgres_reaped_lease_record_with_checkpoint(
440            job_id,
441            job_type,
442            None,
443            1,
444            1,
445            1,
446            Some(checkpoint),
447            Some("worker-reaped-terminal".to_owned()),
448            false,
449            ReapedLeaseDisposition::DeadLetteredTerminal { payload },
450        )
451    }
452
453    #[tokio::test]
454    async fn notify_observers_reports_reaped_retryable_lease() {
455        let observer = RecordingObserver::default();
456        let observers = JobLifecycleObservers::from_observer(observer.clone());
457        let next_run_at = Utc::now();
458        let mut tasks = ReapedObserverTasks::owned();
459
460        let job = postgres_reaped_lease_record(
461            Uuid::now_v7(),
462            runledger_core::jobs::JobTypeName::from_static("jobs.test.reaper.observer"),
463            None,
464            1,
465            2,
466            3,
467            Some("worker-lost-lease".to_owned()),
468            true,
469            ReapedLeaseDisposition::RetryScheduled {
470                retry_delay_ms: 1_000,
471                next_run_at,
472            },
473        );
474        tasks.spawn_batch(&observers, &[job]);
475        wait_for_reaped_count(&observer, 1, Duration::from_millis(500)).await;
476        tasks.abort_for_shutdown().await;
477
478        let reaped = observer
479            .reaped
480            .lock()
481            .expect("reaped events lock should not be poisoned")
482            .clone();
483        assert_eq!(reaped.len(), 1);
484        assert_eq!(reaped[0].job.worker_id, "worker-lost-lease");
485        assert_eq!(reaped[0].job.attempt, 2);
486        assert_eq!(reaped[0].job.max_attempts, 3);
487        assert_eq!(reaped[0].failure.code, "job.lease_expired");
488        assert!(reaped[0].started_without_renewal_heartbeat);
489        assert!(matches!(
490            reaped[0].disposition,
491            crate::observer::JobLeaseReapedDisposition::RetryScheduled { .. }
492        ));
493    }
494
495    #[tokio::test]
496    async fn reaped_observer_task_owner_drops_newest_callback_at_cap() {
497        let started = Arc::new(Notify::new());
498        let started_count = Arc::new(AtomicUsize::new(0));
499        let active_callbacks = Arc::new(AtomicUsize::new(0));
500        let observers = JobLifecycleObservers::from_observer(CancellableReapedObserver {
501            started: started.clone(),
502            started_count: started_count.clone(),
503            active_callbacks: active_callbacks.clone(),
504        });
505        let jobs: Vec<_> = (0..2)
506            .map(|_| reaped_lease_record("jobs.test.reaper.observer.cap"))
507            .collect();
508        let mut tasks = ReapedObserverTasks::owned_with_max_concurrency(1);
509
510        tasks.spawn_batch(&observers, &jobs);
511        timeout(
512            Duration::from_millis(500),
513            wait_for_count(started_count.clone(), started.clone(), 1),
514        )
515        .await
516        .expect("first reaped observer callback should start");
517        sleep(Duration::from_millis(25)).await;
518
519        assert_eq!(active_callbacks.load(Ordering::SeqCst), 1);
520        assert_eq!(tasks.in_flight_count(), 1);
521        tasks.abort_for_shutdown().await;
522        assert_eq!(active_callbacks.load(Ordering::SeqCst), 0);
523    }
524
525    #[tokio::test]
526    async fn notify_observers_starts_multiple_reaped_callbacks_without_serial_timeout() {
527        let started = Arc::new(Notify::new());
528        let started_count = Arc::new(AtomicUsize::new(0));
529        let active_callbacks = Arc::new(AtomicUsize::new(0));
530        let observers = JobLifecycleObservers::from_observer(CancellableReapedObserver {
531            started: started.clone(),
532            started_count: started_count.clone(),
533            active_callbacks: active_callbacks.clone(),
534        });
535        let jobs: Vec<_> = (0..3)
536            .map(|_| reaped_lease_record("jobs.test.reaper.concurrent.observer"))
537            .collect();
538        let mut tasks = ReapedObserverTasks::owned();
539
540        tasks.spawn_batch(&observers, &jobs);
541        timeout(
542            Duration::from_millis(75),
543            wait_for_count(started_count.clone(), started.clone(), 3),
544        )
545        .await
546        .expect("reaped observer callbacks should not wait for serial observer timeouts to start");
547        assert_eq!(active_callbacks.load(Ordering::SeqCst), 3);
548
549        tasks.abort_for_shutdown().await;
550        assert_eq!(active_callbacks.load(Ordering::SeqCst), 0);
551    }
552
553    #[tokio::test]
554    async fn notify_observers_shutdown_cancels_inflight_reaped_observer() {
555        let started = Arc::new(Notify::new());
556        let started_count = Arc::new(AtomicUsize::new(0));
557        let active_callbacks = Arc::new(AtomicUsize::new(0));
558        let observers = JobLifecycleObservers::from_observer(CancellableReapedObserver {
559            started: started.clone(),
560            started_count: started_count.clone(),
561            active_callbacks: active_callbacks.clone(),
562        });
563        let jobs: Vec<_> = (0..2)
564            .map(|_| reaped_lease_record("jobs.test.reaper.cancel.observer"))
565            .collect();
566        let mut tasks = ReapedObserverTasks::owned();
567
568        tasks.spawn_batch(&observers, &jobs);
569        timeout(
570            Duration::from_millis(500),
571            wait_for_count(started_count, started, 2),
572        )
573        .await
574        .expect("reaped observer callbacks should start");
575        assert_eq!(active_callbacks.load(Ordering::SeqCst), 2);
576
577        tasks.abort_for_shutdown().await;
578        assert_eq!(
579            active_callbacks.load(Ordering::SeqCst),
580            0,
581            "reaped observer callback must be cancelled before shutdown returns"
582        );
583    }
584
585    #[tokio::test]
586    async fn reaper_side_effects_notify_terminal_hooks_before_slow_reaped_observers() {
587        let hook_started = Arc::new(Notify::new());
588        let hook_release = Arc::new(Notify::new());
589        let hook_completions = Arc::new(AtomicUsize::new(0));
590        let mut registry = JobRegistry::new();
591        registry.register(ControlledTerminalHookHandler {
592            started: hook_started.clone(),
593            release: hook_release.clone(),
594            completions: hook_completions.clone(),
595        });
596
597        let observer_started = Arc::new(Notify::new());
598        let observer_started_count = Arc::new(AtomicUsize::new(0));
599        let observers = JobLifecycleObservers::from_observer(SlowReapedObserver {
600            started: observer_started.clone(),
601            started_count: observer_started_count.clone(),
602        });
603        let job_id = Uuid::now_v7();
604        let job_type =
605            runledger_core::jobs::JobTypeName::from_static("jobs.test.reaper.hook.controlled");
606        let reaped_leases = vec![terminal_reaped_lease_record(
607            job_id,
608            job_type,
609            json!({ "kind": "reaper-side-effect-order" }),
610        )];
611        let (_shutdown_tx, mut shutdown_rx) = watch::channel(false);
612
613        let mut notification_task = tokio::spawn(async move {
614            let mut reaped_observer_tasks = ReapedObserverTasks::owned();
615            let result = notify_reaped_lease_side_effects(
616                &registry,
617                &observers,
618                &mut reaped_observer_tasks,
619                &reaped_leases,
620                &mut shutdown_rx,
621            )
622            .await;
623            (result, reaped_observer_tasks)
624        });
625
626        timeout(Duration::from_millis(500), hook_started.notified())
627            .await
628            .expect("terminal dead-letter hook should start before reaped observer fanout");
629        for _ in 0..10 {
630            tokio::task::yield_now().await;
631        }
632        assert_eq!(
633            observer_started_count.load(Ordering::SeqCst),
634            0,
635            "reaped observer callback should not start while terminal hook is still running"
636        );
637
638        hook_release.notify_waiters();
639        timeout(
640            Duration::from_millis(500),
641            wait_for_count(observer_started_count, observer_started, 1),
642        )
643        .await
644        .expect("reaped observer fanout should start after terminal hook delivery");
645        assert_eq!(hook_completions.load(Ordering::SeqCst), 1);
646
647        let (result, mut reaped_observer_tasks) =
648            timeout(Duration::from_secs(2), &mut notification_task)
649                .await
650                .expect("reaper notification fanout should finish after scheduling observers")
651                .expect("notification task should not panic");
652        assert_eq!(result, ReaperNotificationFanoutResult::Completed);
653        reaped_observer_tasks.abort_for_shutdown().await;
654    }
655
656    #[tokio::test]
657    async fn notify_handlers_survives_terminal_hook_timeout() {
658        let mut registry = JobRegistry::new();
659        registry.register(HangingTerminalHookHandler);
660
661        let jobs = vec![terminal_reaped_lease_record(
662            Uuid::now_v7(),
663            runledger_core::jobs::JobTypeName::from_static("jobs.test.reaper.hook.hang"),
664            json!({ "kind": "hook-timeout" }),
665        )];
666
667        let (_shutdown_tx, mut shutdown_rx) = watch::channel(false);
668        let result = timeout(
669            Duration::from_secs(2),
670            notify_handlers_of_terminal_lease_expirations(&registry, &jobs, &mut shutdown_rx),
671        )
672        .await
673        .expect("notification pass should return even when terminal hook hangs");
674
675        assert_eq!(result, TerminalHookFanoutResult::Completed { started: 1 });
676    }
677
678    #[tokio::test]
679    async fn notify_handlers_shutdown_bounds_hanging_terminal_hook_batch() {
680        const HOOK_COUNT: usize = 80;
681
682        let started = Arc::new(AtomicUsize::new(0));
683        let mut registry = JobRegistry::new();
684        registry.register(CountingHangingTerminalHookHandler {
685            started: started.clone(),
686        });
687
688        let jobs: Vec<_> = (0..HOOK_COUNT)
689            .map(|_| {
690                terminal_reaped_lease_record(
691                    Uuid::now_v7(),
692                    runledger_core::jobs::JobTypeName::from_static(
693                        "jobs.test.reaper.hook.counting_hang",
694                    ),
695                    json!({ "kind": "bounded-shutdown-hook-drain" }),
696                )
697            })
698            .collect();
699        let (shutdown_tx, mut shutdown_rx) = watch::channel(false);
700        shutdown_tx
701            .send(true)
702            .expect("shutdown receiver should still be active");
703
704        let result = timeout(
705            Duration::from_millis(350),
706            notify_handlers_of_terminal_lease_expirations(&registry, &jobs, &mut shutdown_rx),
707        )
708        .await
709        .expect("shutdown hook fanout should not wait for every hanging hook timeout");
710
711        assert_eq!(
712            result,
713            TerminalHookFanoutResult::InterruptedByShutdown {
714                started: 8,
715                skipped: HOOK_COUNT - 8
716            }
717        );
718        assert_eq!(
719            started.load(Ordering::SeqCst),
720            8,
721            "shutdown should admit only one bounded hook concurrency window"
722        );
723    }
724
725    #[tokio::test]
726    async fn notify_handlers_shutdown_after_entry_before_first_admission_uses_bounded_window() {
727        const HOOK_COUNT: usize = 80;
728
729        let started = Arc::new(AtomicUsize::new(0));
730        let mut registry = JobRegistry::new();
731        registry.register(CountingHangingTerminalHookHandler {
732            started: started.clone(),
733        });
734
735        let jobs: Vec<_> = (0..HOOK_COUNT)
736            .map(|_| {
737                terminal_reaped_lease_record(
738                    Uuid::now_v7(),
739                    runledger_core::jobs::JobTypeName::from_static(
740                        "jobs.test.reaper.hook.counting_hang",
741                    ),
742                    json!({ "kind": "mid-fanout-shutdown-hook-drain" }),
743                )
744            })
745            .collect();
746        let (shutdown_tx, mut shutdown_rx) = watch::channel(false);
747
748        let result = timeout(
749            Duration::from_millis(350),
750            notify_handlers_of_terminal_lease_expirations_with_before_first_hook_admission(
751                &registry,
752                &jobs,
753                &mut shutdown_rx,
754                move || {
755                    shutdown_tx
756                        .send(true)
757                        .expect("shutdown receiver should still be active");
758                },
759            ),
760        )
761        .await
762        .expect("shutdown hook fanout should not skip every committed hook");
763
764        assert_eq!(
765            result,
766            TerminalHookFanoutResult::InterruptedByShutdown {
767                started: 8,
768                skipped: HOOK_COUNT - 8
769            }
770        );
771        assert_eq!(
772            started.load(Ordering::SeqCst),
773            8,
774            "shutdown first observed after entry should admit one bounded hook concurrency window"
775        );
776    }
777
778    #[tokio::test]
779    async fn reaper_side_effects_shutdown_after_committed_terminal_batch_delivers_hook() {
780        let started = Arc::new(Notify::new());
781        let release = Arc::new(Notify::new());
782        let completions = Arc::new(AtomicUsize::new(0));
783        let mut registry = JobRegistry::new();
784        registry.register(ControlledTerminalHookHandler {
785            started: started.clone(),
786            release: release.clone(),
787            completions: completions.clone(),
788        });
789        let observer_started = Arc::new(Notify::new());
790        let observer_started_count = Arc::new(AtomicUsize::new(0));
791        let observers = JobLifecycleObservers::from_observer(SlowReapedObserver {
792            started: observer_started.clone(),
793            started_count: observer_started_count.clone(),
794        });
795        let mut reaped_observer_tasks = ReapedObserverTasks::owned();
796        let (shutdown_tx, mut shutdown_rx) = watch::channel(false);
797        shutdown_tx
798            .send(true)
799            .expect("shutdown receiver should still be active");
800
801        let reaped_leases = vec![terminal_reaped_lease_record(
802            Uuid::now_v7(),
803            runledger_core::jobs::JobTypeName::from_static("jobs.test.reaper.hook.controlled"),
804            json!({ "kind": "committed-batch-shutdown" }),
805        )];
806
807        let mut notification_task = tokio::spawn(async move {
808            notify_reaped_lease_side_effects(
809                &registry,
810                &observers,
811                &mut reaped_observer_tasks,
812                &reaped_leases,
813                &mut shutdown_rx,
814            )
815            .await
816        });
817
818        timeout(Duration::from_millis(500), started.notified())
819            .await
820            .expect("terminal hook should start even though shutdown is already requested");
821        assert_eq!(completions.load(Ordering::SeqCst), 0);
822        assert!(
823            timeout(Duration::from_millis(25), observer_started.notified())
824                .await
825                .is_err(),
826            "reaped observer should not start while shutdown is requested after terminal hook delivery"
827        );
828
829        release.notify_waiters();
830        let result = timeout(Duration::from_secs(2), &mut notification_task)
831            .await
832            .expect("notification fanout should return after terminal hook drains")
833            .expect("notification task should not panic");
834
835        assert_eq!(
836            result,
837            ReaperNotificationFanoutResult::InterruptedByShutdown
838        );
839        assert_eq!(completions.load(Ordering::SeqCst), 1);
840        assert_eq!(observer_started_count.load(Ordering::SeqCst), 0);
841    }
842
843    #[tokio::test]
844    async fn notify_handlers_delivers_terminal_hook_for_committed_batch() {
845        let dead_letters = Arc::new(Mutex::new(Vec::new()));
846        let contexts = Arc::new(Mutex::new(Vec::new()));
847        let mut registry = JobRegistry::new();
848        registry.register(RecordingDeadLetterHandler {
849            dead_letters: dead_letters.clone(),
850            contexts: contexts.clone(),
851        });
852
853        let checkpoint = json!({ "cursor": 900 });
854        let jobs = vec![terminal_reaped_lease_record_with_checkpoint(
855            Uuid::now_v7(),
856            runledger_core::jobs::JobTypeName::from_static("jobs.test.reaper.dead_letter.record"),
857            json!({ "kind": "committed-terminal-hook" }),
858            checkpoint.clone(),
859        )];
860
861        let (_shutdown_tx, mut shutdown_rx) = watch::channel(false);
862        let result =
863            notify_handlers_of_terminal_lease_expirations(&registry, &jobs, &mut shutdown_rx).await;
864
865        assert_eq!(result, TerminalHookFanoutResult::Completed { started: 1 });
866        assert_eq!(
867            dead_letters
868                .lock()
869                .expect("dead-letter list lock should not be poisoned")
870                .len(),
871            1,
872            "terminal hook must still run for a committed terminal reaper batch"
873        );
874        assert_eq!(
875            contexts
876                .lock()
877                .expect("dead-letter context list lock should not be poisoned")[0]
878                .checkpoint,
879            Some(checkpoint),
880            "terminal hook must receive the checkpoint committed before lease expiry"
881        );
882    }
883
884    #[tokio::test]
885    async fn reaper_side_effects_closed_shutdown_channel_does_not_start_reaped_observer() {
886        let observer_started = Arc::new(Notify::new());
887        let observer_started_count = Arc::new(AtomicUsize::new(0));
888        let observers = JobLifecycleObservers::from_observer(SlowReapedObserver {
889            started: observer_started.clone(),
890            started_count: observer_started_count.clone(),
891        });
892        let registry = JobRegistry::new();
893        let mut reaped_observer_tasks = ReapedObserverTasks::owned();
894        let reaped_leases = vec![reaped_lease_record("jobs.test.reaper.closed.observer")];
895        let (shutdown_tx, mut shutdown_rx) = watch::channel(false);
896        drop(shutdown_tx);
897
898        let result = notify_reaped_lease_side_effects(
899            &registry,
900            &observers,
901            &mut reaped_observer_tasks,
902            &reaped_leases,
903            &mut shutdown_rx,
904        )
905        .await;
906
907        assert_eq!(
908            result,
909            ReaperNotificationFanoutResult::InterruptedByShutdown
910        );
911        assert_eq!(reaped_observer_tasks.in_flight_count(), 0);
912        assert_eq!(observer_started_count.load(Ordering::SeqCst), 0);
913        assert!(
914            timeout(Duration::from_millis(50), observer_started.notified())
915                .await
916                .is_err(),
917            "reaped observer callback must not start after shutdown channel closes"
918        );
919    }
920
921    #[tokio::test]
922    async fn notify_handlers_drains_controlled_terminal_hook() {
923        let started = Arc::new(Notify::new());
924        let release = Arc::new(Notify::new());
925        let completions = Arc::new(AtomicUsize::new(0));
926        let mut registry = JobRegistry::new();
927        registry.register(ControlledTerminalHookHandler {
928            started: started.clone(),
929            release: release.clone(),
930            completions: completions.clone(),
931        });
932
933        let jobs = vec![terminal_reaped_lease_record(
934            Uuid::now_v7(),
935            runledger_core::jobs::JobTypeName::from_static("jobs.test.reaper.hook.controlled"),
936            json!({ "kind": "non-shutdown-watch-update" }),
937        )];
938
939        tokio::spawn(async move {
940            started.notified().await;
941            sleep(Duration::from_millis(20)).await;
942            release.notify_waiters();
943        });
944
945        let (_shutdown_tx, mut shutdown_rx) = watch::channel(false);
946        let result =
947            notify_handlers_of_terminal_lease_expirations(&registry, &jobs, &mut shutdown_rx).await;
948
949        assert_eq!(result, TerminalHookFanoutResult::Completed { started: 1 });
950        assert_eq!(completions.load(Ordering::SeqCst), 1);
951    }
952
953    #[tokio::test]
954    async fn notify_handlers_reports_lease_expiration_dead_letter_reason() {
955        let dead_letters = Arc::new(Mutex::new(Vec::new()));
956        let contexts = Arc::new(Mutex::new(Vec::new()));
957        let mut registry = JobRegistry::new();
958        registry.register(RecordingDeadLetterHandler {
959            dead_letters: dead_letters.clone(),
960            contexts,
961        });
962
963        let jobs = vec![terminal_reaped_lease_record(
964            Uuid::now_v7(),
965            runledger_core::jobs::JobTypeName::from_static("jobs.test.reaper.dead_letter.record"),
966            json!({ "kind": "lease-expired" }),
967        )];
968
969        let (_shutdown_tx, mut shutdown_rx) = watch::channel(false);
970        let result =
971            notify_handlers_of_terminal_lease_expirations(&registry, &jobs, &mut shutdown_rx).await;
972
973        assert_eq!(result, TerminalHookFanoutResult::Completed { started: 1 });
974        let dead_letters = dead_letters
975            .lock()
976            .expect("dead-letter list lock should not be poisoned");
977        assert_eq!(dead_letters.len(), 1);
978        let dead_letter = &dead_letters[0];
979        assert_eq!(dead_letter.reason, JobDeadLetterReason::LeaseExpired);
980        assert_eq!(
981            dead_letter.failure.kind,
982            runledger_core::jobs::JobFailureKind::LeaseExpired
983        );
984        assert_eq!(dead_letter.max_attempts, Some(1));
985    }
986}