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 ®istry,
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(®istry, &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(®istry, &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 ®istry,
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 ®istry,
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(®istry, &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 ®istry,
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(®istry, &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(®istry, &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}