later 0.0.48

Distributed Background jobs manager and runner for Rust
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
#[cfg(feature = "prometheus")]
use crate::metrics::{HandlerOutcome, HandlerTimer};
use crate::{
    core::BgJobHandler,
    encoder,
    models::{
        AmqpCommand, ChannelCommand, DelayedStage, FailureTransition, Job, RecurringMode,
        RequeuedStage, Stage, StageName,
    },
    mq::MessageHeaders,
    JobId, RecurringJobId,
};
use async_std::channel::Sender;
use std::{sync::Arc, time::Duration};

#[tracing::instrument(skip_all, fields(command = ?command))]
pub(crate) async fn handle_amqp_command<C, H>(
    command: AmqpCommand,
    worker_id: i32,
    #[cfg(feature = "prometheus")] worker_id_label: &str,
    handler: &Arc<H>,
    // No longer read: every poll command that used to be re-triggered
    // through here now runs directly from the in-process ops loop (see
    // `bg_job_server::start_bg_worker_system_ops_inproc_cmd_to_amqp_cmd`)
    // instead of round-tripping through an `AmqpCommand` and the shared
    // delivery queue. Kept as a parameter (not removed) since every layer
    // between a worker's own channel handle and here would otherwise need
    // unwinding along with it, for no behavior change.
    _inproc_cmd_tx: &Sender<ChannelCommand>,
    headers: Option<MessageHeaders>,
) -> Result<(), anyhow::Error>
where
    C: Sync + Send,
    H: BgJobHandler<C> + Sync + Send + 'static,
{
    if let Some(headers) = headers {
        use opentelemetry::{
            propagation::TextMapPropagator, sdk::propagation::TraceContextPropagator,
        };
        use tracing_opentelemetry::OpenTelemetrySpanExt;

        let propagator = TraceContextPropagator::new();
        let context = propagator.extract(&headers);
        tracing::Span::current().set_parent(context);
    }

    tracing::debug!("Amqp Command: {:?}", command);
    match command {
        AmqpCommand::PartitionReady { topic } => {
            tracing::debug!(
                "[Worker#{}] amqp_command: PartitionReady [Topic: {}]",
                worker_id,
                topic
            );
            handler.get_publisher().partition_wake().notify_waiters();
        }
        AmqpCommand::ExecuteJob(job) => {
            tracing::debug!("[Worker#{}] amqp_command: Job [Id: {}]", worker_id, job.id);

            let Some((claimed_job, lease)) =
                handler.get_publisher().claim_and_fetch_job(&job.id).await?
            else {
                tracing::debug!(job_id = %job.id, "Job is already leased by another worker");
                return Ok(());
            };
            let mut lease_lost = lease.lease_lost();
            // Finished (the lease row deleted) once this delivery is done,
            // *unless* it turns out to be a stray redelivery landing on a
            // job some earlier, now-dead attempt left at `Running` without
            // ever finishing. Delayed/Waiting/Requeued/Success/Failed are
            // all tracked some other way, so skipping and finishing is
            // correct there - nothing else needs their lease. `Running` is
            // the one exception: an abandoned Running job's execution
            // lease is the *only* signal `PollStuckJobs` has to ever
            // recover it (see `stale_leased_job_ids`), so finishing
            // (deleting) the lease we just (re)claimed here, only to
            // discover we must not touch the job, would permanently erase
            // that signal - no future redelivery or poll could ever find
            // it again. Leaving it unfinished lets its renewal simply stop
            // (see below) and its `lease_until` go stale on its own, which
            // is exactly what makes it visible to `PollStuckJobs` again.
            let mut finish_lease = true;
            let result = async {
                if let Some(job) = claimed_job {
                    let should_run = match &job.stage {
                        Stage::Enqueued(_) => true,
                        Stage::Delayed(_)
                        | Stage::Waiting(_)
                        | Stage::Requeued(_)
                        | Stage::Success(_)
                        | Stage::Failed(_) => false,
                        Stage::Running(_) => {
                            finish_lease = false;
                            false
                        }
                    };
                    if should_run {
                        // Recurring occurrences are scheduled proactively by
                        // `AmqpCommand::PollRecurringJobs`, independent of
                        // any occurrence's own execution — see
                        // `handle_poll_recurring_job_command`.
                        handle_job(
                            job,
                            #[cfg(feature = "prometheus")]
                            worker_id_label,
                            handler.clone(),
                            &mut lease_lost,
                        )
                        .await?;
                    }
                }
                Ok::<(), anyhow::Error>(())
            }
            .await;
            let finish_result = if finish_lease {
                lease.finish().await
            } else {
                tracing::debug!(
                    job_id = %job.id,
                    "Found this job already Running from an earlier attempt - not \
                     finishing its execution lease, so PollStuckJobs can find and \
                     reclaim it once the lease actually goes stale"
                );
                drop(lease);
                Ok(())
            };
            result?;
            finish_result?;
        }
    }

    Ok(())
}

#[tracing::instrument(level = "trace", skip(handler))]
pub(crate) async fn handle_poll_delayed_job_command<C, H: BgJobHandler<C>>(
    handler: Arc<H>,
) -> anyhow::Result<()> {
    tracing::debug!("Polling delayed jobs");

    let publisher = handler.get_publisher();
    let mut cursor = None;
    loop {
        let page = publisher.storage.delayed_jobs_page(cursor).await?;
        for item in page.items {
            let job_id = encoder::decode::<JobId>(&item.value)?;
            let Some(job) = publisher.storage.get_job(job_id.clone()).await? else {
                publisher
                    .storage
                    .remove_job_from_poll_range(&job_id, &DelayedStage::get_name())
                    .await?;
                continue;
            };
            let Stage::Delayed(delay) = &job.stage else {
                publisher
                    .storage
                    .remove_job_from_poll_range(&job_id, &DelayedStage::get_name())
                    .await?;
                continue;
            };
            if delay.is_time() {
                tracing::debug!("Job {}: Waiting is finished", job.id);
                publisher.handle_job_enqueue_initial(job).await?;
            }
        }
        let Some(next_cursor) = page.next_cursor else {
            return Ok(());
        };
        cursor = Some(next_cursor);
    }
}

#[tracing::instrument(level = "trace", skip(handler))]
pub(crate) async fn handle_poll_requeued_job_command<C, H: BgJobHandler<C>>(
    handler: Arc<H>,
) -> anyhow::Result<()> {
    tracing::debug!("Polling reqd jobs");

    let publisher = handler.get_publisher();
    let mut cursor = None;
    loop {
        let page = publisher.storage.requeued_jobs_page(cursor).await?;
        for item in page.items {
            let job_id = encoder::decode::<JobId>(&item.value)?;
            let Some(job) = publisher.storage.get_job(job_id.clone()).await? else {
                publisher
                    .storage
                    .remove_job_from_poll_range(&job_id, &RequeuedStage::get_name())
                    .await?;
                continue;
            };
            let Stage::Requeued(requeued) = &job.stage else {
                publisher
                    .storage
                    .remove_job_from_poll_range(&job_id, &RequeuedStage::get_name())
                    .await?;
                continue;
            };
            if !requeued.is_ready() {
                continue;
            }
            tracing::debug!("Job {}: Requeue #{}", job.id, requeued.requeue_count);

            let enqueued = job.transition(); // Requeued -> Enqueued
            publisher.handle_job_enqueue_initial(enqueued).await?;
        }
        let Some(next_cursor) = page.next_cursor else {
            return Ok(());
        };
        cursor = Some(next_cursor);
    }
}

/// Proactively enqueues any recurring job occurrence that's due.
///
/// Runs independently of any occurrence's own execution — unlike the
/// previous approach of scheduling the next occurrence from inside
/// `ExecuteJob`, a crashed or never-run occurrence cannot silently stop the
/// schedule, since this scan is driven purely by wall-clock time and each
/// definition's own `next_run_at` bookmark.
///
/// `Normal`-mode occurrences are enqueued on the plain unordered path.
/// `Sequential`-mode occurrences are enqueued into the reserved recurring
/// topic partition keyed by the recurring job's identifier; the existing
/// partition machinery (see `later::topic`) then guarantees at most one
/// instance of that recurring job's chain runs at a time, in strict order.
/// A chain still running when its next occurrence comes due is not skipped
/// — the occurrence is enqueued anyway and simply queues up behind it.
#[tracing::instrument(level = "trace", skip(handler))]
pub(crate) async fn handle_poll_recurring_job_command<C, H: BgJobHandler<C>>(
    handler: Arc<H>,
) -> anyhow::Result<()> {
    tracing::debug!("Polling recurring jobs");

    let publisher = handler.get_publisher();
    let now = chrono::Utc::now();
    let mut cursor = None;
    loop {
        let page = publisher
            .storage
            .all_recurring_jobs_page(cursor, 100)
            .await?;
        for item in page.items {
            let id = encoder::decode::<RecurringJobId>(&item.value)?;
            let Some(mut rec_job) = publisher.storage.get_recurring_job(id).await? else {
                continue;
            };

            let next_run_at = match rec_job.next_run_at {
                Some(next_run_at) => next_run_at,
                None => rec_job.next_occurrence_after(now)?,
            };
            if next_run_at > now {
                continue;
            }

            tracing::debug!("Recurring job {}: occurrence due", rec_job.id);
            match rec_job.mode {
                RecurringMode::Normal => {
                    let job = rec_job.occurrence_job(next_run_at, None);
                    publisher.enqueue_internal_job(job).await?;
                }
                RecurringMode::Sequential => {
                    let (topic, partition) =
                        publisher.recurring_sequential_topic_partition(&rec_job.id)?;
                    let job =
                        rec_job.occurrence_job(next_run_at, Some((topic.clone(), partition.0)));
                    publisher
                        .commit_job_to_partition(job, topic, partition)
                        .await?;
                }
            }

            rec_job.next_run_at = Some(rec_job.next_occurrence_after(next_run_at)?);
            publisher.storage.save_recurring_job(&rec_job).await?;
        }
        let Some(next_cursor) = page.next_cursor else {
            return Ok(());
        };
        cursor = Some(next_cursor);
    }
}

/// Bounded batch: sweeps already-expired rows out of `later_storage` and
/// `later_storage_range`. See [`crate::models::AmqpCommand::PollExpiredStorage`]
/// for why this exists - reads only ever filter expired rows, nothing else
/// deletes them.
///
/// Loops until a call removes nothing (fully caught up) or a generous
/// iteration cap is hit, so a backlog built up before this poller existed
/// gets fully swept over a few ticks rather than trickling out one small
/// batch per minute forever.
///
/// Pauses briefly between iterations (once a real backlog needs more than
/// one) rather than hammering the database back-to-back - each batch is
/// its own write transaction, and a backend with a single database-wide
/// write lock (SQLite, even in WAL mode) serializes every writer behind
/// whichever one is currently running. Fifty 2,000-row deletes issued with
/// no gap between them starved every other writer - including reads that
/// also need a moment of the write lock, like the dashboard's own
/// `count`/`jobs_in_stage` commands used to (see
/// `storage::sqlite::Sqlite::range_count`) - for as long as this loop kept
/// finding more to remove. The pause costs this loop itself a little wall
/// time to fully catch up on a large backlog; it does not change how much
/// is removed per call.
pub(crate) async fn handle_poll_expired_storage_command<C, H: BgJobHandler<C>>(
    handler: Arc<H>,
) -> anyhow::Result<()> {
    const BATCH: usize = 2_000;
    const MAX_ITERATIONS: usize = 50;
    const PAUSE: std::time::Duration = std::time::Duration::from_millis(200);

    let publisher = handler.get_publisher();
    // Best-effort and independent of the sweep: a log that cannot be
    // checkpointed right now must not stop expired rows being removed.
    if let Err(error) = publisher.storage.inner.checkpoint_wal().await {
        tracing::warn!(%error, "Could not checkpoint the database write-ahead log");
    }
    for i in 0..MAX_ITERATIONS {
        if i > 0 {
            tokio::time::sleep(PAUSE).await;
        }
        let removed = publisher.storage.inner.sweep_expired(BATCH).await?;
        tracing::debug!(removed, "Swept expired storage rows");
        if removed < BATCH {
            break;
        }
    }

    #[cfg(feature = "dashboard")]
    {
        let namespace = publisher.storage.key_prefix().to_string();
        let now = chrono::Utc::now();
        for i in 0..MAX_ITERATIONS {
            if i > 0 {
                tokio::time::sleep(PAUSE).await;
            }
            let removed = publisher
                .storage
                .inner
                .job_index_sweep_expired(&namespace, now, BATCH)
                .await?;
            tracing::debug!(removed, "Swept expired dashboard job-index rows");
            if removed < BATCH {
                break;
            }
        }
    }
    Ok(())
}

/// Reclaims disk space already freed for reuse by `PollExpiredStorage`'s
/// deletes but never returned to the filesystem. See
/// [`crate::Config::vacuum_interval`] - only ever invoked when that's set,
/// and only as often as it allows.
pub(crate) async fn handle_vacuum_database_command<C, H: BgJobHandler<C>>(
    handler: Arc<H>,
) -> anyhow::Result<()> {
    let publisher = handler.get_publisher();
    publisher.storage.inner.vacuum().await?;
    Ok(())
}

/// Reclaims regular (non-partitioned) jobs abandoned mid-handler by a
/// crashed worker. See
/// [`crate::models::AmqpCommand::PollStuckJobs`] for the full mechanism.
///
/// Each stale-leased job ID is re-claimed through the normal execution-lease
/// path (so a job whose original worker is, in fact, still alive and finishes
/// between the scan and this claim is simply skipped - claiming fails the
/// same way it would for any other already-held lease). A job actually still
/// `Running` is failed through the same retry/backoff machinery a handler
/// error goes through, so it gets a normal retry (or ends up `Failed` once
/// retries are exhausted) instead of sitting invisible forever.
pub(crate) async fn handle_poll_stuck_jobs_command<C, H: BgJobHandler<C> + Sync>(
    handler: Arc<H>,
) -> anyhow::Result<()> {
    // Comfortably longer than the 10s renewal / 30s lease TTL both built-in
    // backends use, so a lease that's merely mid-renewal is never mistaken
    // for abandoned.
    const GRACE: Duration = Duration::from_secs(45);
    const BATCH: usize = 100;
    // One tick used to reclaim a single batch, so after a crash or a
    // lock-contention storm that stranded thousands of jobs, recovery
    // crawled at BATCH jobs per poll cycle (several seconds each). Keep
    // draining while full batches come back, bounded so one tick stays
    // short and the other maintenance polls are never starved.
    const MAX_BATCHES_PER_TICK: usize = 20;

    let publisher = handler.get_publisher();
    for _ in 0..MAX_BATCHES_PER_TICK {
        let stale_job_ids = publisher
            .committer
            .stale_leased_job_ids(GRACE, BATCH)
            .await?;
        let full_batch = stale_job_ids.len() >= BATCH;
        reclaim_stale_jobs(&handler, stale_job_ids).await?;
        if !full_batch {
            break;
        }
    }
    #[cfg(feature = "dashboard")]
    reconcile_running_index(&handler).await?;
    Ok(())
}

/// Repairs jobs the dashboard index shows as `running` that no lease tracks.
///
/// `stale_leased_job_ids` only finds jobs that still have an execution lease
/// row. A job left `Running` without one (its lease deleted by an older
/// release, or its job record expired while the index still said `running`)
/// was invisible to it forever. This walks the index's old `running` rows and
/// claims each job's lease, which only succeeds when no live worker holds it:
///
/// - no job record: the index row is dropped;
/// - record no longer `Running`: the index row is corrected;
/// - ordinary job still `Running`: failed/retried like any abandoned job;
/// - partitioned job still `Running`: left to its partition lease.
#[cfg(feature = "dashboard")]
async fn reconcile_running_index<C, H: BgJobHandler<C> + Sync>(
    handler: &Arc<H>,
) -> anyhow::Result<()> {
    // A live handler holds a lease and fails the claim below, so this only
    // bounds how recently entered a job must be before it is looked at.
    const MIN_AGE: chrono::Duration = chrono::Duration::minutes(10);
    const BATCH: usize = 100;
    const MAX_BATCHES_PER_TICK: usize = 20;

    let publisher = handler.get_publisher();
    let namespace = publisher.storage.key_prefix().to_owned();
    let mut cursor = Some((chrono::Utc::now() - MIN_AGE).timestamp_millis());
    for _ in 0..MAX_BATCHES_PER_TICK {
        let page = publisher
            .storage
            .inner
            .job_index_list_by_stage(&namespace, "running", cursor, BATCH)
            .await?;
        for row in &page.items {
            let job_id = JobId(row.job_id.clone());
            let Some((claimed_job, lease)) = publisher.claim_and_fetch_job(&job_id).await? else {
                continue;
            };
            match claimed_job {
                None => {
                    tracing::warn!(%job_id, "Dropping a dashboard `running` row whose job no longer exists");
                    publisher
                        .storage
                        .inner
                        .job_index_remove(&namespace, &row.job_id)
                        .await?;
                }
                Some(job) if !matches!(job.stage, Stage::Running(_)) => {
                    let expire = matches!(job.stage, Stage::Success(_) | Stage::Failed(_))
                        .then(|| chrono::Utc::now() + chrono::Duration::hours(1));
                    publisher.stats.record_transition(&job, expire).await;
                }
                Some(job) if job.topic.is_some() => {}
                Some(job) => {
                    tracing::warn!(%job_id, "Reclaiming a `Running` job that has no live lease");
                    fail_abandoned_running_job(handler, job).await?;
                }
            }
            lease.finish().await?;
        }
        match page.next_cursor {
            Some(next) => cursor = Some(next),
            None => break,
        }
    }
    Ok(())
}

async fn reclaim_stale_jobs<C, H: BgJobHandler<C> + Sync>(
    handler: &Arc<H>,
    stale_job_ids: Vec<crate::JobId>,
) -> anyhow::Result<()> {
    let publisher = handler.get_publisher();
    for job_id in stale_job_ids {
        let Some((claimed_job, lease)) = publisher.claim_and_fetch_job(&job_id).await? else {
            // Someone else claimed it since the scan - either it finished,
            // or another reclaim attempt is already handling it.
            continue;
        };
        let Some(job) = claimed_job else {
            lease.finish().await?;
            continue;
        };
        if !matches!(job.stage, Stage::Running(_)) {
            // Already moved on (or moved on and got reclaimed again) since
            // the scan - nothing to do.
            lease.finish().await?;
            continue;
        }

        fail_abandoned_running_job(handler, job).await?;
        lease.finish().await?;
    }
    Ok(())
}

/// Fails (with retry, if any are left) a `Running` job whose worker is gone,
/// the same way a handler error would. The caller holds the job's lease.
async fn fail_abandoned_running_job<C, H: BgJobHandler<C> + Sync>(
    handler: &Arc<H>,
    mut job: Job,
) -> anyhow::Result<()> {
    let publisher = handler.get_publisher();
    tracing::warn!(
        job_id = %job.id,
        "Reclaiming a job whose execution lease expired without completing \
         (worker likely crashed mid-handler)"
    );
    if job.config.needs_retry_policy() {
        match handler.retry_policy(&job.payload_type, &job.payload).await {
            Ok(retry_policy) => job.config.resolve_retry_policy(retry_policy.as_ref()),
            Err(policy_error) => tracing::warn!(
                "Could not resolve retry policy for job {}: {}",
                job.id,
                policy_error
            ),
        }
    }
    let reason =
        "job execution lease expired without completing (worker likely crashed)".to_string();
    match job.transition_failure(reason)? {
        FailureTransition::Retry(job) => publisher.save(&job).await?,
        FailureTransition::Failed(job) => {
            publisher
                .save_and_expire(&job, Duration::from_secs(3600))
                .await?;
        }
    }
    Ok(())
}

/// Enqueues every job waiting on `success_job`, keeping partition ordering
/// for children that belong to a topic partition.
///
/// A partitioned child was already assigned its sequence when it was
/// enqueued, so it is only transitioned and saved here; the partition head
/// scanner picks it up once it becomes the ready head, rather than being
/// published on the shared delivery queue.
pub(crate) async fn enqueue_continuations<C, H>(
    handler: Arc<H>,
    success_job: &Job,
    success_job_id: &JobId,
) -> anyhow::Result<()>
where
    C: Sync + Send,
    H: BgJobHandler<C> + Sync + Send + 'static,
{
    let publisher = handler.get_publisher();
    let waiting_jobs = publisher.storage.get_continuation_jobs(success_job).await?;
    if waiting_jobs.is_empty() {
        return Ok(());
    }
    for next in waiting_jobs {
        tracing::info!("Continuing {} -> {}", success_job_id, next.id);

        let next_job = next.transition(); // Waiting -> Enqueued
        if next_job.topic_partition().is_some() {
            publisher.save(&next_job).await?;
            publisher.partition_wake().notify_waiters();
        } else {
            publisher.handle_job_enqueue_initial(next_job).await?;
        }
    }
    publisher
        .storage
        .clear_continuation_jobs(success_job)
        .await?;
    Ok(())
}

#[tracing::instrument(skip_all, fields(job_id = %job.id, job_type = %job.payload_type))]
pub(crate) async fn handle_job<C, H>(
    job: Job,
    #[cfg(feature = "prometheus")] worker_id: &str,
    handler: Arc<H>,
    lease_lost: &mut tokio::sync::watch::Receiver<()>,
) -> Result<(), anyhow::Error>
where
    C: Sync + Send,
    H: BgJobHandler<C> + Sync + Send + 'static,
{
    let ptype = job.payload_type.clone();
    let payload = job.payload.clone();
    let job_id = job.id.clone();

    let publisher = handler.get_publisher();
    let running_job = job.transition();
    publisher.save(&running_job).await?;

    #[cfg(feature = "prometheus")]
    let handler_timer = HandlerTimer::start(publisher.metrics_queue(), worker_id, &ptype);
    // Races the handler against this job's own execution lease: if the
    // lease is lost mid-dispatch (renewal starved long enough for it to
    // expire, and another worker re-claimed the job), abandon this attempt
    // entirely rather than run to completion and unconditionally overwrite
    // whatever the new claimant does next - there is no per-job fencing on
    // these writes the way there is for the partitioned path, so not
    // writing at all is the only way to avoid clobbering it. Mirrors
    // `partition::handle_partitioned_job`'s identical `tokio::select!`.
    let mut lease_was_lost = false;
    let handler_result = tokio::select! {
        result = handler.dispatch(ptype.clone(), &payload, job_id) => result,
        changed = lease_lost.changed() => {
            lease_was_lost = true;
            match changed {
                Ok(()) => Err(anyhow::anyhow!("job execution lease lost while handler was running")),
                Err(_) => Err(anyhow::anyhow!("job lease renewal task stopped while handler was running")),
            }
        }
    };
    #[cfg(feature = "prometheus")]
    handler_timer.finish(match &handler_result {
        Ok(()) => HandlerOutcome::Success,
        Err(_) => HandlerOutcome::Error,
    });

    if lease_was_lost {
        tracing::warn!(
            job_id = %running_job.id,
            "Abandoning job after losing its execution lease mid-dispatch; deferring to whichever \
             worker now holds it"
        );
        return Ok(());
    }

    match handler_result {
        Ok(_) => {
            // success
            let success_job = running_job.transition_success()?;
            let success_job_id = success_job.id.clone();
            publisher
                .save_and_expire(&success_job, Duration::from_secs(3600))
                .await?;

            // enqueue waiting jobs
            enqueue_continuations(handler.clone(), &success_job, &success_job_id).await?;
        }
        Err(error) => {
            tracing::warn!("Failed job {}: {}", running_job.id, error);

            let mut running_job = running_job;
            if running_job.config.needs_retry_policy() {
                match handler.retry_policy(&ptype, &payload).await {
                    Ok(retry_policy) => running_job
                        .config
                        .resolve_retry_policy(retry_policy.as_ref()),
                    Err(policy_error) => tracing::warn!(
                        "Could not resolve retry policy for job {}: {}",
                        running_job.id,
                        policy_error
                    ),
                }
            }
            match running_job.transition_failure(error.to_string())? {
                FailureTransition::Retry(job) => publisher.save(&job).await?,
                FailureTransition::Failed(job) => {
                    publisher
                        .save_and_expire(&job, Duration::from_secs(3600))
                        .await?;
                }
            }
        }
    }

    Ok(())
}