later 0.0.27

Distributed Background jobs manager and runner for Rust
Documentation
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
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
#[cfg(feature = "prometheus")]
use crate::metrics;
#[cfg(feature = "dashboard")]
use crate::stats::{DashboardRebuildPage, DashboardResponse, ResponseError, Stats};
use crate::{
    core::JobParameter,
    encoder,
    models::{
        AmqpCommand, DelayedStage, EnqueuedStage, FailedStage, Job, JobConfig, RecurringJob,
        RecurringMode, Stage, WaitingStage,
    },
    mq::MqClient,
    persist::Persist,
    storage::Storage,
    BackgroundJobServerPublisher, JobId, RecurringJobId, UtcDateTime,
};
use anyhow::Context;
use std::str::FromStr;
use std::sync::Arc;
use std::time::Duration;

impl BackgroundJobServerPublisher {
    /// Creates a publisher for a prepared backend.
    ///
    /// Generated server builders call this method. Application code normally
    /// receives a publisher through the generated job context.
    pub async fn new(
        id: String,
        mq_client: Arc<Box<dyn MqClient>>,
        storage: Box<dyn Storage>,
        committer: Arc<dyn crate::backend::JobCommitter>,
    ) -> anyhow::Result<Self> {
        Self::new_with_retry_policy(
            id,
            mq_client,
            storage,
            committer,
            crate::retry::RetryPolicy::default(),
        )
        .await
    }

    /// Creates a publisher with an explicit server retry policy.
    ///
    /// Generated server builders use this method. The value is stored when a
    /// job is enqueued, then a message-specific [`crate::retry::JobRetryPolicy`]
    /// can replace it after the first handler failure.
    pub async fn new_with_retry_policy(
        id: String,
        mq_client: Arc<Box<dyn MqClient>>,
        storage: Box<dyn Storage>,
        committer: Arc<dyn crate::backend::JobCommitter>,
        default_retry_policy: crate::retry::RetryPolicy,
    ) -> anyhow::Result<Self> {
        Self::new_with_retry_policy_and_topics(
            id,
            mq_client,
            storage,
            committer,
            default_retry_policy,
            Vec::new(),
            crate::RetainedLogLag::default(),
            None,
        )
        .await
    }

    /// Creates a publisher with an explicit server retry policy and the
    /// topics registered for partitioned enqueue.
    ///
    /// Generated server builders use this method.
    #[allow(clippy::too_many_arguments)]
    pub async fn new_with_retry_policy_and_topics(
        id: String,
        mq_client: Arc<Box<dyn MqClient>>,
        storage: Box<dyn Storage>,
        committer: Arc<dyn crate::backend::JobCommitter>,
        default_retry_policy: crate::retry::RetryPolicy,
        topics: Vec<crate::topic::TopicConfig>,
        retained_log_lag: crate::RetainedLogLag,
        partition_history_retention: Option<Duration>,
    ) -> anyhow::Result<Self> {
        let routing_key = format!("later-{}", id);
        let publisher = mq_client.new_publisher(&routing_key).await?;
        let persist = Arc::new(Persist::new(storage, routing_key.clone()));
        #[cfg(feature = "dashboard")]
        let stats = {
            tracing::info!("Enabling Dashboard Stats Collector");
            Stats::new(
                persist.clone(),
                topics.clone(),
                retained_log_lag.clone(),
                committer.clone(),
                partition_history_retention,
            )
        };
        // Only `Stats` (dashboard-only) reads this - nothing to do with it
        // on a build without the `dashboard` feature, but the parameter
        // still needs a use to avoid an unused-variable warning there.
        #[cfg(not(feature = "dashboard"))]
        let _ = partition_history_retention;
        Ok(Self {
            storage: persist.clone(),

            #[cfg(feature = "dashboard")]
            stats,
            publisher,
            routing_key,
            committer,
            default_retry_policy,
            topics,
            retained_log_lag,
            partition_wake: Arc::new(tokio::sync::Notify::new()),
        })
    }

    /// The configured source of retained-log consumer lag, if any - see
    /// [`crate::RetainedLogLag`]. Used by the partition-worker maintenance
    /// loop to refresh the `later_retained_log_consumer_lag` gauge.
    pub(crate) fn retained_log_lag(&self) -> &crate::RetainedLogLag {
        &self.retained_log_lag
    }

    /// Wakes this process's partition pollers so they scan sooner than their
    /// idle backoff, instead of waiting for the next tick.
    pub(crate) fn partition_wake(&self) -> &Arc<tokio::sync::Notify> {
        &self.partition_wake
    }

    #[cfg(feature = "dashboard")]
    /// Handles one dashboard HTTP request.
    ///
    /// `prefix` is the URL path where the application mounted the dashboard.
    /// Pass the raw URL query string as `query_string`, then copy the returned
    /// status, headers, and body into the application's HTTP response. The host
    /// application is responsible for access control.
    pub async fn get_dashboard(
        &self,
        prefix: String,
        query_string: String,
    ) -> Result<DashboardResponse, ResponseError> {
        self.stats.handle_http(prefix, query_string).await
    }

    /// Rebuilds at most `limit` dashboard records from durable job state.
    ///
    /// Pass the returned cursor to the next call until it is `None`. This is
    /// useful after enabling the dashboard for a database that already contains
    /// jobs. `limit` must be between 1 and 1,000 inclusive.
    #[cfg(feature = "dashboard")]
    pub async fn rebuild_dashboard_page(
        &self,
        cursor: Option<i64>,
        limit: usize,
    ) -> anyhow::Result<DashboardRebuildPage> {
        self.stats.rebuild_page(cursor, limit).await
    }

    #[cfg(feature = "dashboard")]
    pub(crate) async fn publish_worker_heartbeat(&self, worker_id: String) -> anyhow::Result<()> {
        self.stats.record_worker_heartbeat(worker_id).await
    }

    /// Deletes the oldest entries in one `(topic, partition)`'s dashboard
    /// job-history list until at most `keep_last` remain, and returns how
    /// many were removed.
    ///
    /// The dashboard's per-partition history (the list behind its
    /// `partition_jobs` view) is a full, append-only record by design - it
    /// is never pruned as jobs complete, unlike the job rows themselves
    /// (which expire an hour after reaching a terminal state). Left alone,
    /// a busy partition's history grows forever. Call this periodically
    /// (e.g. from a recurring job) if you want to cap it instead.
    ///
    /// This only touches the dashboard's own bookkeeping, never the actual
    /// job queue - jobs still waiting to run in this partition are
    /// unaffected either way.
    #[cfg(feature = "dashboard")]
    pub async fn truncate_partition_history(
        &self,
        topic: &str,
        partition: u32,
        keep_last: usize,
    ) -> anyhow::Result<usize> {
        self.stats
            .truncate_partition_history(topic, partition, keep_last)
            .await
    }

    /// Blocks until at least one worker is available.
    ///
    /// Generated builders use this during startup to ensure readiness.
    pub async fn ensure_worker_ready(&self) -> anyhow::Result<()> {
        self.publisher.ensure_consumer().await
    }

    #[cfg(feature = "prometheus")]
    /// Encodes Later's process metrics in Prometheus text format.
    ///
    /// The worker metrics cover this process. Queue depth is refreshed from
    /// shared storage every five seconds. Mount the returned text on an HTTP
    /// endpoint that Prometheus can scrape. Queue, worker, state, and job type
    /// are bounded labels; job IDs and error messages are never used as labels.
    pub fn get_metrics(&self) -> anyhow::Result<String> {
        metrics::encode()
    }

    #[cfg(feature = "prometheus")]
    pub(crate) fn metrics_queue(&self) -> &str {
        &self.routing_key
    }

    /// Enqueues `message` after `parent_job_id` succeeds.
    ///
    /// The child waits while the parent is unfinished. If the parent record has
    /// already expired, Later treats it as successfully completed.
    pub async fn enqueue_continue(
        &self,
        parent_job_id: JobId,
        message: impl JobParameter,
    ) -> anyhow::Result<JobId> {
        self.enqueue_internal(message, Some(parent_job_id), None, None)
            .await
    }

    /// Enqueues `message` to run after at least `delay`.
    ///
    /// Returns an error if the delay cannot be represented as a UTC timestamp.
    pub async fn enqueue_delayed(
        &self,
        message: impl JobParameter,
        delay: std::time::Duration,
    ) -> anyhow::Result<JobId> {
        let enqueue_time = chrono::Utc::now()
            .checked_add_signed(chrono::Duration::from_std(delay)?)
            .ok_or(anyhow::anyhow!("Error calculating enqueue time"))?;

        self.enqueue_delayed_at(message, enqueue_time).await
    }

    /// Enqueues `message` to run at or after a UTC timestamp.
    ///
    /// Returns an error when `time` is not in the future.
    pub async fn enqueue_delayed_at(
        &self,
        message: impl JobParameter,
        time: chrono::DateTime<chrono::Utc>,
    ) -> anyhow::Result<JobId> {
        if time <= chrono::Utc::now() {
            return Err(anyhow::anyhow!("Time must be in the future"));
        }
        self.enqueue_internal(message, None, Some(time), None).await
    }

    /// Enqueues `message` for immediate execution by an available worker.
    pub async fn enqueue(&self, message: impl JobParameter) -> anyhow::Result<JobId> {
        self.enqueue_internal(message, None, None, None).await
    }

    /// Enqueues `message` for strict, in-order execution within its topic
    /// partition.
    ///
    /// Only message types declared with a topic in [`crate::background_job!`]
    /// can be used here; they implement [`crate::topic::JobPartition`] and
    /// the generated [`crate::topic::JobTopic`] marker. The topic must be
    /// registered on this server's [`crate::Config`] and the returned
    /// partition must be below its configured partition count, or this
    /// returns an error without storing the job.
    pub async fn enqueue_to_partition<M>(&self, message: M) -> anyhow::Result<JobId>
    where
        M: JobParameter + crate::topic::JobPartition + crate::topic::JobTopic,
    {
        self.enqueue_to_partition_internal(message, None, None)
            .await
    }

    /// Enqueues `message` to its topic partition, to run at or after `delay`.
    ///
    /// The partition remains blocked on this job until it becomes ready and
    /// completes, even though later jobs in the partition may be enqueued
    /// sooner. See [`Self::enqueue_to_partition`] for the topic requirements.
    pub async fn enqueue_to_partition_delayed<M>(
        &self,
        message: M,
        delay: std::time::Duration,
    ) -> anyhow::Result<JobId>
    where
        M: JobParameter + crate::topic::JobPartition + crate::topic::JobTopic,
    {
        let enqueue_time = chrono::Utc::now()
            .checked_add_signed(chrono::Duration::from_std(delay)?)
            .ok_or(anyhow::anyhow!("Error calculating enqueue time"))?;
        self.enqueue_to_partition_internal(message, None, Some(enqueue_time))
            .await
    }

    /// Enqueues `message` to its topic partition after `parent_job_id`
    /// succeeds.
    ///
    /// The partition stays blocked on this job while its parent is
    /// unfinished, so a continuation waiting on a job in another partition
    /// can intentionally block this partition. See
    /// [`Self::enqueue_to_partition`] for the topic requirements.
    pub async fn enqueue_continue_to_partition<M>(
        &self,
        parent_job_id: JobId,
        message: M,
    ) -> anyhow::Result<JobId>
    where
        M: JobParameter + crate::topic::JobPartition + crate::topic::JobTopic,
    {
        self.enqueue_to_partition_internal(message, Some(parent_job_id), None)
            .await
    }

    #[tracing::instrument(skip(self, message), fields(ptype = message.get_ptype(), job_id), name = "init_create_partitioned_job")]
    async fn enqueue_to_partition_internal<M>(
        &self,
        message: M,
        parent_job_id: Option<JobId>,
        delay_until: Option<UtcDateTime>,
    ) -> anyhow::Result<JobId>
    where
        M: JobParameter + crate::topic::JobPartition + crate::topic::JobTopic,
    {
        let topic = M::topic_name().to_string();
        let key = message.partition_key();
        let partition = crate::topic::resolve_partition(&self.topics, &topic, &key)?;

        self.enqueue_to_resolved_partition_internal(
            message,
            parent_job_id,
            delay_until,
            topic,
            partition,
        )
        .await
    }

    /// Enqueues `message` into an already-resolved `(topic, partition)`,
    /// without requiring the message type to implement
    /// [`crate::topic::JobPartition`]/[`crate::topic::JobTopic`].
    ///
    /// [`Self::enqueue_to_partition_internal`] derives the partition from the
    /// message type's declared key; this lower-level path is for callers
    /// that already know which partition a job belongs to — recurring
    /// sequential jobs (partitioned by their identifier, not by message
    /// content) and continuations that must inherit their parent's
    /// partition.
    #[tracing::instrument(skip(self, message), fields(ptype = message.get_ptype(), job_id), name = "init_create_resolved_partitioned_job")]
    pub(crate) async fn enqueue_to_resolved_partition_internal(
        &self,
        message: impl JobParameter,
        parent_job_id: Option<JobId>,
        delay_until: Option<UtcDateTime>,
        topic: String,
        partition: crate::topic::PartitionId,
    ) -> anyhow::Result<JobId> {
        let job = crate::bg_job_server_publisher::create_job_in_topic(
            message,
            parent_job_id,
            delay_until,
            None,
            &self.default_retry_policy,
            Some((topic.clone(), partition.0)),
        )?;
        tracing::Span::current().record("job_id", job.id.to_string());

        self.commit_job_to_partition(job, topic, partition).await
    }

    /// Commits an already-built `job` (its `topic`/`partition` fields must
    /// already name `topic`/`partition`) into that topic partition.
    ///
    /// Shared by [`Self::enqueue_to_resolved_partition_internal`] (which
    /// builds the job from a [`crate::core::JobParameter`]) and the
    /// recurring-job paths, which build occurrences directly from an
    /// already-encoded [`RecurringJob`] payload instead.
    pub(crate) async fn commit_job_to_partition(
        &self,
        job: Job,
        topic: String,
        partition: crate::topic::PartitionId,
    ) -> anyhow::Result<JobId> {
        let id = job.id.clone();
        let payload_type = job.payload_type.clone();
        let earliest_run_at = match &job.stage {
            Stage::Delayed(delayed) => delayed.not_before,
            Stage::Waiting(_) | Stage::Enqueued(_) => chrono::Utc::now(),
            Stage::Running(_) | Stage::Requeued(_) | Stage::Success(_) | Stage::Failed(_) => {
                chrono::Utc::now()
            }
        };

        // Built once the transaction inside `commit_partitioned` has
        // assigned a sequence, and applied in that same transaction - not a
        // separate follow-up write. A worker can claim, run, and complete a
        // fast handler within milliseconds of this job becoming visible; a
        // second, unsynchronized write racing that completion (as this used
        // to be, only to stamp `sequence` for dashboard/lookup purposes)
        // could commit after the worker's own `Success` write and silently
        // revert it back to `Enqueued` forever, since nothing re-wakes a
        // job no longer at the head of its partition. See the
        // `JobCommitter::commit_partitioned` doc comment.
        #[cfg(feature = "prometheus")]
        metrics::record_job_transition(&self.routing_key, &job);

        let routing_key = self.routing_key.clone();
        #[cfg(feature = "dashboard")]
        let stats = self.stats.clone();
        let build_operations: Box<
            dyn FnOnce(
                    crate::backend::PartitionSequence,
                ) -> anyhow::Result<Vec<crate::storage::StorageOperation>>
                + Send,
        > = Box::new(move |sequence| {
            let job = Job {
                sequence: Some(sequence.0),
                ..job
            };
            let operations = Persist::job_operations(&routing_key, &job)?;
            #[cfg(feature = "dashboard")]
            let operations = {
                let mut operations = operations;
                operations.extend(stats.save_job_operations(&job)?);
                operations
            };
            Ok(operations)
        });
        self.committer
            .commit_partitioned(
                build_operations,
                &topic,
                partition,
                &id,
                &payload_type,
                earliest_run_at,
            )
            .await?;

        // Local pollers can pick this straight up instead of waiting out
        // their idle backoff; cross-process workers get the same signal
        // through the SQL delivery queue.
        self.partition_wake.notify_waiters();

        Ok(id)
    }

    #[tracing::instrument(skip(self, message), fields(ptype = message.get_ptype(), job_id), name = "init_create_job")]
    async fn enqueue_internal(
        &self,
        message: impl JobParameter,
        parent_job_id: Option<JobId>,
        delay_until: Option<UtcDateTime>,
        recurring_job_id: Option<RecurringJobId>,
    ) -> anyhow::Result<JobId> {
        let job = create_job(
            message,
            parent_job_id,
            delay_until,
            recurring_job_id,
            &self.default_retry_policy,
        )?;

        tracing::Span::current().record("job_id", job.id.to_string());

        self.enqueue_internal_job(job).await
    }

    pub(crate) async fn enqueue_internal_job(&self, job: Job) -> Result<JobId, anyhow::Error> {
        let id = job.id.clone();

        self.handle_job_enqueue_initial(job).await?;
        Ok(id)
    }

    /// Registers a recurring payload, enqueuing its first occurrence only if
    /// this identifier isn't already registered on the same schedule.
    ///
    /// `identifier` identifies the recurring definition and `cron` uses the
    /// cron crate's seconds-first syntax. Registering the same `identifier`
    /// again updates its payload and schedule (upsert) without touching its
    /// already-scheduled `next_run_at` or enqueuing a duplicate occurrence -
    /// see [`Self::save_recurring_job_definition`] for why this matters:
    /// applications typically call this once on every process startup, and
    /// re-enqueuing unconditionally there would add one more occurrence per
    /// restart, forever. The recurring-job poller enqueues every future
    /// occurrence regardless; this method only ever needs to enqueue the
    /// very first one, or a fresh one when the cron schedule actually
    /// changed. Returns `None` in every other case. Returns an error for an
    /// invalid cron expression.
    ///
    /// Occurrences run on schedule and may overlap with each other or with
    /// jobs chained off them via [`Self::enqueue_continue`]. Use
    /// [`Self::enqueue_recurring_sequential`] when at most one instance of
    /// this recurring job may ever be running.
    pub async fn enqueue_recurring(
        &self,
        identifier: String,
        message: impl JobParameter,
        cron: String,
    ) -> anyhow::Result<Option<JobId>> {
        let (recurring_job, should_enqueue_now) = self
            .save_recurring_job_definition(identifier, message, cron, RecurringMode::Normal)
            .await?;
        if !should_enqueue_now {
            return Ok(None);
        }

        let first_job = recurring_job.try_into()?;
        Ok(Some(self.enqueue_internal_job(first_job).await?))
    }

    /// Registers a recurring payload whose occurrences never overlap,
    /// enqueuing its first occurrence only if this identifier isn't already
    /// registered on the same schedule.
    ///
    /// Same upsert semantics as [`Self::enqueue_recurring`] - see there and
    /// [`Self::save_recurring_job_definition`] for why re-registering the
    /// same identifier does not by itself enqueue another occurrence.
    /// Every occurrence — and every job chained off one with
    /// [`Self::enqueue_recurring_continue`] — is enqueued into the same
    /// topic partition (keyed by `identifier`), so the existing sequential
    /// partition guarantee (see [`crate::topic`]) gives at most one instance
    /// of this recurring job's chain running at a time, in strict order.
    /// Occurrences due while the previous chain is still running queue up
    /// rather than being skipped or dropped.
    pub async fn enqueue_recurring_sequential(
        &self,
        identifier: String,
        message: impl JobParameter,
        cron: String,
    ) -> anyhow::Result<Option<JobId>> {
        let (recurring_job, should_enqueue_now) = self
            .save_recurring_job_definition(identifier, message, cron, RecurringMode::Sequential)
            .await?;
        if !should_enqueue_now {
            return Ok(None);
        }

        let not_before = recurring_job.next_occurrence_after(chrono::Utc::now())?;
        let (topic, partition) = self.recurring_sequential_topic_partition(&recurring_job.id)?;
        let job = recurring_job.occurrence_job(not_before, Some((topic.clone(), partition.0)));
        Ok(Some(
            self.commit_job_to_partition(job, topic, partition).await?,
        ))
    }

    /// Enqueues `message` after `parent_job_id` succeeds, inheriting its
    /// partition when it has one.
    ///
    /// Use this instead of [`Self::enqueue_continue`] to chain work off a
    /// job that may belong to a sequential recurring job's chain — including
    /// from inside that job's own handler, chaining off itself with
    /// `ctx.job_id()`. When `parent_job_id` isn't in a topic partition
    /// (a `Normal`-mode occurrence, or any non-recurring job), this behaves
    /// exactly like [`Self::enqueue_continue`].
    pub async fn enqueue_recurring_continue(
        &self,
        parent_job_id: JobId,
        message: impl JobParameter,
    ) -> anyhow::Result<JobId> {
        let parent_partition = self
            .storage
            .get_job(parent_job_id.clone())
            .await?
            .and_then(|job| job.topic_partition().map(|(t, p)| (t.to_string(), p)));

        match parent_partition {
            Some((topic, partition)) => {
                self.enqueue_to_resolved_partition_internal(
                    message,
                    Some(parent_job_id),
                    None,
                    topic,
                    crate::topic::PartitionId(partition),
                )
                .await
            }
            None => self.enqueue_continue(parent_job_id, message).await,
        }
    }

    /// Validates `cron`, prepares `message`, and upserts the `RecurringJob`
    /// definition. Shared by [`Self::enqueue_recurring`] and
    /// [`Self::enqueue_recurring_sequential`], which differ only in how they
    /// enqueue an occurrence when this returns `true`.
    ///
    /// A brand-new identifier (or one whose cron schedule actually changed)
    /// gets a freshly computed `next_run_at` and `true`, telling the caller
    /// to enqueue that first occurrence itself, the same as ever. Otherwise
    /// this is a no-op re-registration - a plain payload/config update, if
    /// anything - and the existing `next_run_at` (and `date_added`) are left
    /// untouched, with `false` telling the caller not to enqueue anything:
    /// the recurring-job poller already owns enqueuing every occurrence from
    /// here on, reading `next_run_at` straight from storage, so calling this
    /// again (as every process startup does, to recreate the schedule) must
    /// not also hand out a brand new occurrence each time - that would add
    /// one more permanently-queued job per restart, forever, not just once.
    async fn save_recurring_job_definition(
        &self,
        identifier: String,
        message: impl JobParameter,
        cron: String,
        mode: RecurringMode,
    ) -> anyhow::Result<(RecurringJob, bool)> {
        let _ = cron::Schedule::from_str(&cron).context("error parsing cron expression")?;

        let message = prepare_message(message, &self.default_retry_policy)?;
        let id = RecurringJobId(identifier);
        let existing = self.storage.get_recurring_job(id.clone()).await?;

        let Some(existing) = existing else {
            // Looks brand new - but another process could be registering
            // this exact identifier for the first time at the same moment
            // (the common case for a horizontally-scaled deployment, where
            // every replica calls this on startup). Claim it atomically
            // rather than deciding "I'm first" from this read alone: only
            // one of the racing callers' create actually lands, and that
            // one - not necessarily this one - is the one that should
            // enqueue the first occurrence.
            let mut recurring_job = RecurringJob {
                id,
                payload_type: message.payload_type,
                payload: message.payload,
                cron_schedule: cron,
                date_added: chrono::Utc::now(),
                config: message.config,
                mode,
                next_run_at: None,
            };
            recurring_job.next_run_at =
                Some(recurring_job.next_occurrence_after(chrono::Utc::now())?);
            if self
                .storage
                .create_recurring_job_if_absent(&recurring_job)
                .await?
            {
                return Ok((recurring_job, true));
            }
            let winner = self
                .storage
                .get_recurring_job(recurring_job.id.clone())
                .await?
                .ok_or_else(|| {
                    anyhow::anyhow!(
                        "recurring job {} disappeared immediately after a concurrent \
                         registration created it",
                        recurring_job.id
                    )
                })?;
            return Ok((winner, false));
        };

        let should_enqueue_now = existing.cron_schedule != cron;
        let mut recurring_job = RecurringJob {
            id: existing.id,
            payload_type: message.payload_type,
            payload: message.payload,
            cron_schedule: cron,
            date_added: existing.date_added,
            config: message.config,
            mode,
            next_run_at: existing.next_run_at,
        };
        if should_enqueue_now {
            recurring_job.next_run_at =
                Some(recurring_job.next_occurrence_after(chrono::Utc::now())?);
        }

        self.storage.save_recurring_job(&recurring_job).await?;
        Ok((recurring_job, should_enqueue_now))
    }

    /// Resolves the reserved sequential-recurring topic partition a
    /// recurring job's identifier hashes to.
    pub(crate) fn recurring_sequential_topic_partition(
        &self,
        id: &RecurringJobId,
    ) -> anyhow::Result<(String, crate::topic::PartitionId)> {
        let topic = crate::topic::RECURRING_SEQUENTIAL_TOPIC.to_string();
        let key = crate::topic::PartitionKey::from(id.to_string());
        let partition = crate::topic::resolve_partition(&self.topics, &topic, &key).context(
            "sequential recurring jobs require Config::recurring_sequential_partitions to be set",
        )?;
        Ok((topic, partition))
    }

    pub(crate) fn job_save_operations(
        &self,
        job: &Job,
    ) -> anyhow::Result<Vec<crate::storage::StorageOperation>> {
        let operations = Persist::job_operations(&self.routing_key, job)?;
        #[cfg(feature = "dashboard")]
        let operations = {
            let mut operations = operations;
            operations.extend(self.stats.save_job_operations(job)?);
            operations
        };
        Ok(operations)
    }

    pub(crate) fn job_save_and_expire_operations(
        &self,
        job: &Job,
        duration: Duration,
    ) -> anyhow::Result<Vec<crate::storage::StorageOperation>> {
        let mut operations = self.job_save_operations(job)?;
        operations.extend(Persist::expire_job_operations(
            &self.routing_key,
            job.id.clone(),
            duration,
        )?);
        #[cfg(feature = "dashboard")]
        {
            operations.extend(self.stats.expire_job_operations(job)?);
        }
        Ok(operations)
    }

    pub(crate) async fn save(&self, job: &Job) -> anyhow::Result<()> {
        let operations = self.job_save_operations(job)?;
        self.storage.inner.apply(operations).await?;

        #[cfg(feature = "prometheus")]
        metrics::record_job_transition(&self.routing_key, job);

        Ok(())
    }

    pub(crate) async fn save_and_expire(
        &self,
        job: &Job,
        duration: Duration,
    ) -> anyhow::Result<()> {
        let operations = self.job_save_and_expire_operations(job, duration)?;
        self.storage.inner.apply(operations).await?;

        #[cfg(feature = "prometheus")]
        metrics::record_job_transition(&self.routing_key, job);
        Ok(())
    }

    #[async_recursion::async_recursion]
    pub(crate) async fn handle_job_enqueue_initial(&self, job: Job) -> anyhow::Result<()> {
        tracing::debug!(
            "handle_job_enqueue_initial: Id: {}, Stage: {:?}",
            &job.id,
            &job.stage
        );

        match &job.stage {
            Stage::Delayed(delayed) => {
                self.save(&job).await?;
                // delayed job
                // should be polled

                if delayed.is_time() {
                    let job = job.transition(); // Delayed -> Enqueued
                    self.handle_job_enqueue_initial(job).await?;
                }
            }
            Stage::Waiting(waiting) => {
                self.save(&job).await?;
                // continuation
                // - enqueue if parent is already complete
                // - schedule self message to check an enqueue later (to prevent race)

                if let Some(parent_job) = self.storage.get_job(waiting.parent_id.clone()).await? {
                    if !parent_job.stage.is_success() {
                        return Ok(());
                    }

                    tracing::info!(
                        "Parent job {} is already completed, enqueuing this job immediately",
                        parent_job.id
                    );
                }

                // parent job is success or not found (means successful long time ago)
                let job = job.transition();
                self.handle_job_enqueue_initial(job).await?;
            }
            Stage::Enqueued(_) => {
                tracing::debug!("Enqueue job {}", job.id);
                self.save_and_publish_job(job).await?;
            }
            Stage::Running(_) | Stage::Requeued(_) | Stage::Success(_) | Stage::Failed(_) => {
                tracing::warn!("Invalid job here {}, Stage {:?}", job.id, &job.stage);
                //unreachable!("stage is handled in consumer")
            }
        }

        Ok(())
    }

    async fn save_and_publish_job(&self, job: Job) -> anyhow::Result<()> {
        let command = AmqpCommand::ExecuteJob(crate::models::JobAmqp {
            payload_type: job.payload_type.clone(),
            id: job.id.clone(),
        });
        let operations = Persist::job_operations(&self.routing_key, &job)?;
        #[cfg(feature = "dashboard")]
        let operations = {
            let mut operations = operations;
            operations.extend(self.stats.save_job_operations(&job)?);
            operations
        };
        self.committer
            .commit(operations, &self.routing_key, &encoder::encode(&command)?)
            .await?;

        #[cfg(feature = "prometheus")]
        metrics::record_job_transition(&self.routing_key, &job);

        Ok(())
    }

    /// Claims a job's execution lease and fetches its current stored state
    /// together, in as few round trips as the backend allows.
    ///
    /// Returns `None` when another worker already owns or completed the
    /// job. Returns `Some` with a `None` job when the lease was granted but
    /// the job's stored state is already gone — the caller still holds the
    /// lease and must finish it.
    pub(crate) async fn claim_and_fetch_job(
        &self,
        job_id: &JobId,
    ) -> anyhow::Result<Option<(Option<Job>, Box<dyn crate::backend::JobLease>)>> {
        let key = self.storage.job_key(job_id);
        let Some(claimed) = self
            .committer
            .claim_and_fetch(job_id, &key, self.storage.inner.as_ref())
            .await?
        else {
            return Ok(None);
        };
        let job = claimed
            .value
            .map(|bytes| encoder::decode::<Job>(&bytes))
            .transpose()?;
        Ok(Some((job, claimed.lease)))
    }

    /// Returns whether this server has any topics configured for partition
    /// ordering.
    pub(crate) fn has_topics(&self) -> bool {
        !self.topics.is_empty()
    }

    #[cfg(feature = "prometheus")]
    pub(crate) fn topic_names(&self) -> impl Iterator<Item = &str> {
        self.topics.iter().map(|topic| topic.name())
    }

    #[cfg(feature = "prometheus")]
    pub(crate) fn topics(&self) -> &[crate::topic::TopicConfig] {
        &self.topics
    }

    pub(crate) async fn claim_partition_head(
        &self,
        owner: &str,
        live_workers: &[String],
    ) -> anyhow::Result<Option<crate::backend::PartitionHeadClaim>> {
        self.committer
            .claim_partition_head(owner, live_workers)
            .await
    }

    pub(crate) async fn peek_assigned_partition_head(
        &self,
        owner: &str,
        topic: &str,
        partition: crate::topic::PartitionId,
        lease_epoch: i64,
    ) -> anyhow::Result<Option<crate::backend::PartitionHeadClaim>> {
        self.committer
            .peek_assigned_partition_head(owner, topic, partition, lease_epoch)
            .await
    }

    pub(crate) async fn heartbeat_partition_worker(&self, owner: &str) -> anyhow::Result<()> {
        self.committer.heartbeat_worker(owner).await
    }

    pub(crate) async fn list_live_partition_workers(&self) -> anyhow::Result<Vec<String>> {
        self.committer.list_live_workers().await
    }

    #[cfg(feature = "prometheus")]
    pub(crate) async fn partition_backlog_summary(
        &self,
    ) -> anyhow::Result<Vec<(String, u64, f64)>> {
        self.committer.partition_backlog_summary().await
    }

    /// Depth of every topic partition's queue - see
    /// [`crate::backend::JobCommitter::partition_queue_depth_summary`].
    /// Used by the dashboard's topic menu/partition picker and (with the
    /// `prometheus` feature) the `later_partition_queue_depth` gauge.
    pub(crate) async fn partition_queue_depth_summary(
        &self,
    ) -> anyhow::Result<Vec<(String, crate::topic::PartitionId, u64)>> {
        self.committer.partition_queue_depth_summary().await
    }

    pub(crate) async fn deregister_partition_worker(&self, owner: &str) -> anyhow::Result<()> {
        self.committer.deregister_worker(owner).await
    }

    pub(crate) async fn renew_partition_lease(
        &self,
        claim: &crate::backend::PartitionHeadClaim,
    ) -> anyhow::Result<bool> {
        self.committer.renew_partition_lease(claim).await
    }

    pub(crate) async fn release_partition_assignment(
        &self,
        owner: &str,
        topic: &str,
        partition: crate::topic::PartitionId,
        lease_epoch: i64,
    ) -> anyhow::Result<bool> {
        self.committer
            .release_partition_assignment(owner, topic, partition, lease_epoch)
            .await
    }

    pub(crate) async fn release_partition_head(
        &self,
        claim: &crate::backend::PartitionHeadClaim,
        earliest_run_at: UtcDateTime,
        operations: Vec<crate::storage::StorageOperation>,
    ) -> anyhow::Result<bool> {
        self.committer
            .release_partition_head(claim, earliest_run_at, operations)
            .await
    }

    pub(crate) async fn complete_partition_head(
        &self,
        claim: &crate::backend::PartitionHeadClaim,
        operations: Vec<crate::storage::StorageOperation>,
    ) -> anyhow::Result<bool> {
        self.committer
            .complete_partition_head(claim, operations)
            .await
    }

    /// Forcibly fails a topic partition's current head job so the partition
    /// can proceed, without waiting for its retry policy to exhaust.
    ///
    /// This is the escape hatch for a "poison" job that would otherwise
    /// block every later job in its partition until its own retries give
    /// up: it moves the head straight to a terminal failed state (through
    /// the normal job-expiry path, just like an ordinary final failure) and
    /// frees the partition immediately, regardless of whether any worker
    /// currently holds its lease. `reason` becomes the job's failure reason.
    ///
    /// Returns the failed job's ID, or `None` if the partition had no head
    /// at all — including if it finished on its own in the moment between
    /// looking it up and acting on it, in which case this call safely does
    /// nothing rather than clobbering the newer state. Each call always acts
    /// on whichever job is currently blocking the partition, so calling this
    /// again after a successful call targets the *next* head, not the one
    /// just failed.
    pub async fn force_fail_partition_head(
        &self,
        topic: &str,
        partition: crate::topic::PartitionId,
        reason: String,
    ) -> anyhow::Result<Option<JobId>> {
        let Some((job_id, sequence)) = self
            .committer
            .admin_peek_partition_head(topic, partition)
            .await?
        else {
            return Ok(None);
        };
        let Some(job) = self.storage.get_job(job_id.clone()).await? else {
            // The job record itself already expired; there's nothing
            // meaningful left to fail.
            return Ok(None);
        };

        let mut previous_stages = job.previous_stages.clone();
        previous_stages.push(job.stage.clone());
        let failed_job = Job {
            stage: Stage::Failed(FailedStage {
                date: chrono::Utc::now(),
                reason,
            }),
            previous_stages,
            ..job
        };
        let operations =
            self.job_save_and_expire_operations(&failed_job, Duration::from_secs(3600))?;

        let cleared = self
            .committer
            .admin_force_complete_partition_head(topic, partition, sequence, &job_id, operations)
            .await?;
        if !cleared {
            return Ok(None);
        }

        #[cfg(feature = "prometheus")]
        metrics::record_job_transition(&self.routing_key, &failed_job);

        // The next job in this partition, if any, can now run; wake local
        // pollers instead of waiting out their idle backoff.
        self.partition_wake.notify_waiters();

        Ok(Some(job_id))
    }
}

pub(crate) fn create_job(
    message: impl JobParameter,
    parent_job_id: Option<JobId>,
    delay_until: Option<chrono::DateTime<chrono::Utc>>,
    recurring_job_id: Option<RecurringJobId>,
    default_retry_policy: &crate::retry::RetryPolicy,
) -> anyhow::Result<Job> {
    create_job_in_topic(
        message,
        parent_job_id,
        delay_until,
        recurring_job_id,
        default_retry_policy,
        None,
    )
}

pub(crate) fn create_job_in_topic(
    message: impl JobParameter,
    parent_job_id: Option<JobId>,
    delay_until: Option<chrono::DateTime<chrono::Utc>>,
    recurring_job_id: Option<RecurringJobId>,
    default_retry_policy: &crate::retry::RetryPolicy,
    topic_partition: Option<(String, u32)>,
) -> anyhow::Result<Job> {
    let id = crate::generate_id();
    let message = prepare_message(message, default_retry_policy)?;
    let (topic, partition) = match topic_partition {
        Some((topic, partition)) => (Some(topic), Some(partition)),
        None => (None, None),
    };
    let job = Job {
        id: JobId(id.clone()),
        payload_type: message.payload_type,
        payload: message.payload,
        stage: {
            if let Some(parent_job_id) = parent_job_id {
                Stage::Waiting(WaitingStage {
                    date: chrono::Utc::now(),
                    parent_id: parent_job_id,
                })
            } else if let Some(delay_until) = delay_until {
                Stage::Delayed(DelayedStage {
                    date: chrono::Utc::now(),
                    not_before: delay_until,
                })
            } else {
                Stage::Enqueued(EnqueuedStage {
                    date: chrono::Utc::now(),
                })
            }
        },
        previous_stages: Vec::default(),
        config: message.config,
        recurring_job_id,
        topic,
        partition,
        sequence: None,
    };
    Ok(job)
}

struct PreparedMessage {
    payload_type: String,
    payload: Vec<u8>,
    config: JobConfig,
}

fn prepare_message(
    message: impl JobParameter,
    default_retry_policy: &crate::retry::RetryPolicy,
) -> anyhow::Result<PreparedMessage> {
    Ok(PreparedMessage {
        payload_type: message.get_ptype(),
        payload: message
            .to_bytes()
            .context("Unable to serialize the message to bytes")?,
        config: JobConfig::from_policy(default_retry_policy),
    })
}