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
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
//! # later
//!
//! A distributed background job manager and runner for Rust.
//!
//! Later stores job state and delivers work through SQLite or PostgreSQL. Every server using the same database
//! and namespace joins the same worker pool. Renewable leases coordinate job
//! execution across workers and server processes.
//!
//! ## Documentation index
//!
//! - [Workers](#how-workers-run) - local worker count, shared worker pools,
//! and delivery.
//! - [Storage and SQL data model](crate::storage) - durable records, tables,
//! retention, and transaction boundaries.
//! - [Retry behavior](#retry-behavior) - per-message and server retry rules.
//! - [Dashboard](#dashboard) - feature-gated job inspection and retention.
//! - [Prometheus metrics](#prometheus-metrics) - feature-gated monitoring.
//! - [Recurring jobs](#recurring-jobs) - cron schedules, upsert-by-identifier,
//! and the sequential (non-overlapping) mode.
//! - [Sequential partition mode](#sequential-partition-mode) - topic
//! ordering, partition keys, ownership, and failure behaviour.
//! - [Retained logs](crate::retained_log) - replayable records and consumer
//! group offsets using the same topic partitions.
//!
//! ## Retained logs
//!
//! Enable `retained-log` with `sqlite` or `postgres` to append immutable bytes
//! and replay them by partition offset. Consumer groups commit offsets and use
//! 30-second, epoch-fenced leases. This is a database-backed log API, not a
//! Kafka protocol implementation: there is no Kafka wire protocol, broker
//! replication, or cross-region quorum.
//!
//! ## Set up
//! ### 1. Import `later` and required dependencies
//!
//! ```toml
//! later = { version = "0.0.28", features = ["sqlite"] }
//! serde = { version = "1.0", features = ["derive"] }
//! ```
//!
//! ### 2. Define some types to use as a payload to the background jobs
//!
//! ```
//! use serde::{Deserialize, Serialize};
//!
//! #[derive(Serialize, Deserialize)] // <- Required derives
//! pub struct SendEmail {
//! pub address: String,
//! pub body: String,
//! }
//!
//! // ... more as required
//! ```
//!
//! ### 3. Generate the stub
//!
//! ```
//! # use serde::{Deserialize, Serialize};
//! #
//! # #[derive(Serialize, Deserialize)] // <- Required derives
//! # pub struct SendEmail {
//! # pub address: String,
//! # pub body: String,
//! # }
//! later::background_job! {
//! struct Jobs {
//! send_email: SendEmail,
//! }
//! }
//! ```
//!
//! This generates two types
//! * `JobsBuilder` - used to bootstrap the background job server - which can be used to enqueue jobs,
//! * `JobContext<T>` - used to pass application context (`T`) in the handler as well as enqueue jobs,
//!
//! ### 4. Use the generated code to bootstrap the background job server
//!
//! For `struct Jobs` a type `JobsBuilder` will be generated. Use this to bootstrap the server.
//!
//! ```ignore
//! # use serde::{Deserialize, Serialize};
//! #
//! # #[derive(Serialize, Deserialize)]
//! # pub struct SendEmail { pub address: String, pub body: String }
//! # pub struct MyContext {}
//! # later::background_job! {
//! # struct Jobs {
//! # send_email: SendEmail,
//! # }
//! # }
//! use later::{backend::SqliteBackend, storage::Sqlite, BackgroundJobServer, Config};
//!
//! # #[tokio::main]
//! # async fn main() -> anyhow::Result<()> {
//! // bootstrap the server
//! let ctx = MyContext{ /*..*/ }; // Any context to pass onto the handlers
//! let storage = Sqlite::new("sqlite://later.db").await?;
//! let backend = SqliteBackend::new("fnf-example", storage)?;
//! let ctx = JobsBuilder::new(
//! later::Config::builder()
//! .context(ctx) // Pass the context here
//! .backend(Box::new(backend)) // Configure state and delivery
//! // ...
//! .build()
//! )
//! // for each payload defined in the `struct Jobs` above
//! // the generated fn name uses the pattern "with_[name]_handler"
//! .with_send_email_handler(handle_send_email) // Pass the handler function
//! // ..
//! .build()
//! .await?;
//!
//! // use ctx.enqueue(SendEmail{ ... }) to enqueue jobs,
//! // or ctx.enqueue_continue(parent_job_id, SendEmail{ ... }) to chain jobs.
//! // this will only accept types defined inside the macro above
//! # Ok(())
//! # }
//! // define handler
//! async fn handle_send_email(
//! ctx: JobsContext<MyContext>, // JobContext is generated wrapper
//! payload: SendEmail,
//! ) -> anyhow::Result<()> {
//! // handle `payload`
//!
//! // ctx.app -> Access the MyContext passed during bootstrapping
//! // ctx.enqueue(_).await to enqueue more jobs
//! // ctx.enqueue_continue(_).await to chain jobs
//!
//! Ok(()) // or Err(_) to retry this message
//! }
//! ```
//!
//! This example uses SQLite for state and delivery.
//!
//! ---
//!
//! ## How workers run
//!
//! A [`BackgroundJobServer`] starts six worker tasks by default. Set
//! [`Config::worker_count`] to change that number. The count applies to each
//! server instance, so three server processes configured with four workers
//! can run up to twelve job handlers at once. Every worker owns one delivery
//! consumer and processes one command at a time. A slow job occupies one worker
//! while the other workers continue receiving commands.
//! [`BackgroundJobServer::add_worker`] starts another local worker at runtime.
//! [`BackgroundJobServer::remove_worker`] waits for one worker's current job
//! before stopping it. At least one worker remains.
//! For process shutdown, call [`BackgroundJobServer::shutdown`] with a
//! deadline. It stops new claims, lets active handlers commit, and then
//! deregisters a sequential-processing owner. If the deadline expires,
//! remaining durable work recovers through its normal lease.
//!
//! Server startup waits until every worker consumer is ready. Workers using
//! the same backend database and namespace compete for work, including workers
//! in other processes or on other hosts. SQLite can coordinate several
//! processes only when they can all open the same database file. PostgreSQL
//! can coordinate processes across hosts.
//!
//! For each job delivery, a worker:
//!
//! 1. claims a renewable execution lease so another worker does not run the
//! same job at the same time;
//! 2. loads the durable job and changes its stage from enqueued to running;
//! 3. calls the generated handler for the stored payload type;
//! 4. stores success, failure, or the next retry; and
//! 5. acknowledges the delivery. A delivery error is reported to the backend
//! so it can be made available again.
//!
//! A delivery backend may make a command available more than once. The
//! execution lease prevents concurrent handler calls, and the stored job stage
//! makes workers ignore duplicate commands after a job has moved on from the
//! enqueued stage. A regular (non-partitioned) job whose handler is
//! interrupted after the running stage is stored - its worker crashed or was
//! killed mid-handler - is reclaimed by the periodic `PollStuckJobs`
//! maintenance command once its execution lease has been expired for a
//! while: it's retried (or marked failed, once retries are exhausted) the
//! same way a handler error would be. Partitioned/sequential jobs don't need
//! this - a stuck partition head already recovers through its own
//! partition-lease expiry.
//!
//! The same command also repairs jobs the dashboard index shows as `running`
//! that no lease tracks (left by older releases, or whose job record has
//! since expired): a missing record drops the index row, and an ordinary job
//! still `Running` is retried or failed like an abandoned one.
//!
//! Delayed jobs and retries do not occupy a sleeping worker. They remain in
//! durable storage until a periodic command finds that they are ready, then
//! they rejoin the delivery queue. Polling and queue load can add delay, so a
//! scheduled time or retry delay is a lower bound rather than an exact start
//! time.
//!
//! A separate periodic command, `PollExpiredStorage`, physically deletes
//! already-expired rows in small batches. Expiry (`Storage::expire`, and the
//! TTLs `later` sets internally on terminal jobs and dashboard bookkeeping)
//! only ever makes a row invisible to reads - something still has to delete
//! it, or storage only ever grows. This runs automatically; there is nothing
//! to configure.
//!
//! Keep the [`BackgroundJobServer`] handle alive for as long as this process
//! should accept work. Dropping it aborts its local worker and maintenance
//! tasks, including the two periodic commands described above.
//!
//! With the `dashboard` feature, each stage change updates job state and its
//! dashboard projection in the same atomic storage operation. Dashboard work
//! therefore cannot form a second delivery backlog during producer load.
//!
//! ## Fire and forget jobs
//!
//! Fire and forget jobs are made available to a worker immediately. A failed
//! handler is retried according to its resolved retry policy.
//!
//! ```no_run
//! # #[derive(serde::Serialize, serde::Deserialize)]
//! # pub struct SendEmail { pub address: String, pub body: String }
//! # later::background_job! {
//! # struct Jobs {
//! # send_email: SendEmail,
//! # }
//! # }
//! # #[tokio::main]
//! # async fn main() -> anyhow::Result<()>{
//! # let ctx : later::BackgroundJobServerPublisher = todo!();
//! ctx.enqueue(SendEmail{
//! address: "hello@rust-lang.org".to_string(),
//! body: "You rock!".to_string()
//! }).await?;
//! # Ok(())
//! # }
//! ```
//!
//! ## Continuations
//!
//! One or many jobs can be chained into a workflow. A child becomes available
//! only after its parent succeeds.
//!
//! ```no_run
//! # #[derive(serde::Serialize, serde::Deserialize)]
//! # pub struct SendEmail { pub address: String, pub body: String }
//! # later::background_job! {
//! # struct Jobs {
//! # send_email: SendEmail,
//! # create_account: CreateAccount,
//! # }
//! # }
//! # #[derive(serde::Serialize, serde::Deserialize)]
//! # pub struct CreateAccount { id: String }
//! # #[tokio::main]
//! # async fn main() -> anyhow::Result<()>{
//! # let ctx : later::BackgroundJobServerPublisher = todo!();
//! let email_welcome = ctx.enqueue(SendEmail{
//! address: "customer@example.com".to_string(),
//! body: "Creating your account!".to_string()
//! }).await?;
//!
//! let create_account = ctx.enqueue_continue(email_welcome, CreateAccount { id: "accout-1".to_string() }).await?;
//!
//! let email_confirmation = ctx.enqueue_continue(create_account, SendEmail{
//! address: "customer@example.com".to_string(),
//! body: "Your account has been created!".to_string()
//! }).await;
//! # Ok(())
//! # }
//! ```
//!
//! ## Delayed jobs
//!
//! A delayed job becomes available at or after a chosen interval or timestamp.
//!
//! ```no_run
//! # #[derive(serde::Serialize, serde::Deserialize)]
//! # pub struct SendEmail { pub address: String, pub body: String }
//! # later::background_job! {
//! # struct Jobs {
//! # send_email: SendEmail,
//! # }
//! # }
//! # #[tokio::main]
//! # async fn main() -> anyhow::Result<()>{
//! # let ctx : later::BackgroundJobServerPublisher = todo!();
//! // delay
//! ctx.enqueue_delayed(SendEmail{
//! address: "hello@rust-lang.org".to_string(),
//! body: "You rock!".to_string()
//! }, std::time::Duration::from_secs(60)).await?;
//!
//! // specific time
//! let run_job_at : chrono::DateTime<chrono::Utc> = todo!();
//! ctx.enqueue_delayed_at(SendEmail{
//! address: "hello@rust-lang.org".to_string(),
//! body: "You rock!".to_string()
//! }, run_job_at).await?;
//! # Ok(())
//! # }
//! ```
//!
//! ## Recurring jobs
//!
//! Run a recurring job on a cron schedule. [`BackgroundJobServerPublisher::enqueue_recurring`]
//! registers a payload and cron expression under an `identifier`; calling it
//! again with the same `identifier` and cron updates the existing
//! definition in place (upsert) rather than creating a second schedule, and
//! returns `None` rather than enqueuing another occurrence - safe to call
//! on every process startup to recreate the schedule, which is the
//! intended pattern; calling it unconditionally on every startup was, once,
//! not safe, and would enqueue one more permanently-queued occurrence per
//! restart. Registering with a *different* cron is treated as a real
//! change and does enqueue a fresh occurrence, returning `Some`. A
//! background poller proactively enqueues each subsequent due occurrence,
//! independent of whether a previous occurrence's handler is still running
//! or ever ran — a crashed process cannot silently stop the schedule.
//!
//! ```no_run
//! # #[derive(serde::Serialize, serde::Deserialize)]
//! # pub struct SendNewsletter { pub address: String }
//! # later::background_job! {
//! # struct Jobs {
//! # send_newsletter: SendNewsletter,
//! # }
//! # }
//! # #[tokio::main]
//! # async fn main() -> anyhow::Result<()>{
//! # let ctx : later::BackgroundJobServerPublisher = todo!();
//! ctx.enqueue_recurring("send-newsletter-1".to_string(),
//! SendNewsletter{
//! address: "hello@rust-lang.org".to_string(),
//! },
//! "0 6 1 * * *".to_string() // 6am, 1st day of every month
//! ).await?;
//! # Ok(())
//! # }
//! ```
//!
//! By default occurrences run on the plain unordered path and may overlap,
//! the same as any other job. Use
//! [`BackgroundJobServerPublisher::enqueue_recurring_sequential`] instead
//! when at most one instance of a recurring job — including any jobs
//! chained off an occurrence with
//! [`BackgroundJobServerPublisher::enqueue_recurring_continue`] — may ever
//! be running at once. It reuses the sequential partition mode described
//! below: every occurrence, and every job chained off one, is enqueued into
//! the same topic partition (keyed by `identifier`), so they run in strict
//! order and never overlap. Requires
//! [`Config::recurring_sequential_partitions`] to be set.
//!
//! ```no_run
//! # #[derive(serde::Serialize, serde::Deserialize)]
//! # pub struct RunReport;
//! # later::background_job! {
//! # struct Jobs {
//! # run_report: RunReport,
//! # }
//! # }
//! # async fn handle_run_report(ctx: JobsContext<()>, _payload: RunReport) -> anyhow::Result<()> {
//! // Chain a follow-up step off the currently-running occurrence itself,
//! // using its own job id. Because this must happen before the handler
//! // returns, the continuation is always ordered ahead of the next
//! // scheduled occurrence.
//! ctx.enqueue_recurring_continue(ctx.job_id().clone(), RunReport).await?;
//! # Ok(())
//! # }
//! # #[tokio::main]
//! # async fn main() -> anyhow::Result<()>{
//! # let ctx : later::BackgroundJobServerPublisher = todo!();
//! ctx.enqueue_recurring_sequential("nightly-report".to_string(),
//! RunReport,
//! "0 0 3 * * *".to_string() // 3am daily
//! ).await?;
//! # Ok(())
//! # }
//! ```
//!
//! ## Sequential partition mode
//!
//! Named topics with a fixed number of Kafka-like partitions give strict,
//! in-order execution within one `(topic, partition)`, while different
//! partitions and topics still run concurrently. Jobs enqueued without a
//! topic keep the default unordered, competing-worker behaviour described
//! above; nothing changes for them.
//!
//! A message type opts in by implementing [`topic::JobPartition`] and being
//! declared with `#[topic("name")]` inside [`background_job!`]. The macro
//! then requires that type to implement [`topic::JobPartition`] and
//! generates the [`topic::JobTopic`] marker the topic-aware enqueue methods
//! need. [`topic::JobPartition::partition_key`] returns an arbitrary
//! [`topic::PartitionKey`] — a customer ID, an order ID, anything that
//! identifies "things that must stay in order relative to each other" —
//! and Later hashes it to a partition with [`topic::partition_for_key`], the
//! way a Kafka producer key is hashed by the partitioner. The caller never
//! names a partition number directly. Register every topic's fixed
//! partition count on [`Config::topics`] before starting the server; an
//! unregistered topic is rejected before anything is written.
//!
//! ```
//! use later::topic::{JobPartition, PartitionKey};
//!
//! #[derive(serde::Serialize, serde::Deserialize)]
//! pub struct ProcessOrder {
//! pub customer_id: u32,
//! pub description: String,
//! }
//!
//! // Every order from the same customer must apply in order; different
//! // customers can usually be processed at the same time.
//! impl JobPartition for ProcessOrder {
//! fn partition_key(&self) -> PartitionKey {
//! PartitionKey::from(self.customer_id)
//! }
//! }
//!
//! later::background_job! {
//! struct Jobs {
//! #[topic("orders")]
//! process_order: ProcessOrder,
//! }
//! }
//! ```
//!
//! Register the topic when building the server, and use the partition-aware
//! enqueue methods for that message type:
//!
//! ```no_run
//! # use later::topic::{JobPartition, PartitionKey};
//! # #[derive(serde::Serialize, serde::Deserialize)]
//! # pub struct ProcessOrder { pub customer_id: u32, pub description: String }
//! # impl JobPartition for ProcessOrder {
//! # fn partition_key(&self) -> PartitionKey { PartitionKey::from(self.customer_id) }
//! # }
//! # later::background_job! {
//! # struct Jobs {
//! # #[topic("orders")]
//! # process_order: ProcessOrder,
//! # }
//! # }
//! use later::topic::TopicConfig;
//!
//! # #[tokio::main]
//! # async fn main() -> anyhow::Result<()>{
//! # let context = ();
//! # let backend: Box<dyn later::backend::Backend> = todo!();
//! let config = later::Config::builder()
//! .context(context)
//! .backend(backend)
//! .topics(vec![TopicConfig::new("orders", 8)?])
//! .build();
//! # let _ = config;
//! # let ctx : later::BackgroundJobServerPublisher = todo!();
//! // Only a type declared with a topic can use the partition-aware methods.
//! ctx.enqueue_to_partition(ProcessOrder {
//! customer_id: 42,
//! description: "add item".to_string(),
//! }).await?;
//! # Ok(())
//! # }
//! ```
//!
//! [`BackgroundJobServerPublisher::enqueue_to_partition`],
//! [`BackgroundJobServerPublisher::enqueue_to_partition_delayed`], and
//! [`BackgroundJobServerPublisher::enqueue_continue_to_partition`] mirror
//! [`BackgroundJobServerPublisher::enqueue`],
//! [`BackgroundJobServerPublisher::enqueue_delayed`], and
//! [`BackgroundJobServerPublisher::enqueue_continue`], but assign the job a
//! monotonic sequence
//! within its `(topic, partition)` as part of the same atomic write. The
//! SQLite and PostgreSQL backends also expose `enqueue_to_partition_in` for
//! writing application data, the job, and its sequence in one caller-owned
//! transaction, matching `enqueue_in`.
//!
//! Ordering rules:
//!
//! - The oldest unfinished job in a partition blocks every later job in that
//! partition, including a retry (which keeps its position), a delayed
//! head (which blocks until its scheduled time), and a waiting
//! continuation (which blocks until its parent succeeds, even when the
//! parent is not itself partitioned or lives in a different partition).
//! - A terminal success or failure releases the next sequence.
//! - Different partitions, and different topics, never block each other,
//! even when they share the same partition ID.
//!
//! Every server with configured topics registers a heartbeat and runs a
//! bounded number of partition poller tasks (one per worker) that lease and
//! run the oldest ready, unowned head. Which server owns a given
//! `(topic, partition)` is decided by rendezvous hashing over every
//! currently live server (heartbeats expire after 15 seconds), so a small
//! membership change only moves the partitions whose winner changes rather
//! than reshuffling every assignment. A lease keeps one worker's claim
//! exclusive across every process sharing the same database and namespace;
//! losing a race for a lease is normal and just means another worker already
//! owns that partition's head. The lease renews itself while a handler
//! keeps running, so a slow handler does not let another worker reclaim and
//! re-run the same head concurrently. As with the default delivery path, a
//! worker crash can still run the current head job again after its lease
//! expires: order is preserved, but execution is at-least-once rather than
//! exactly-once.
//!
//! Each partition claim also has a fencing epoch. If renewal shows that a
//! lease was lost, Later cancels the handler future and rejects that claim's
//! retry or completion write. Only the owner holding the current epoch can
//! advance the stored partition head. This protects Later's own state during
//! a process pause or a reclaimed lease. It cannot undo an external side
//! effect that began before cancellation, so handlers that call external
//! systems must still be idempotent, normally using the job ID as their key.
//!
//! Adding a worker lets it start claiming newly idle or unassigned
//! partitions within a couple of membership refresh cycles, without waiting
//! for any lease to expire. A job already running when its partition's
//! assignment changes finishes under its current owner; the new owner only
//! takes over once that worker stops claiming the partition's next head.
//! Dropping a [`BackgroundJobServer`] deregisters its worker immediately
//! (best-effort) so other workers stop considering it right away; an
//! ungraceful removal (a crash) is only detected once its heartbeat
//! expires.
//!
//! SQL is always the source of order and ownership. A poller also reacts
//! immediately (instead of waiting out its idle backoff) to a local enqueue
//! or continuation. This wake-up is a pure latency optimization; nothing
//! depends on it arriving, since every poller still scans SQL on its own.
//!
//! If a head job would otherwise block its partition indefinitely (a
//! "poison" job whose retries never succeed, or one you simply need
//! unblocked sooner),
//! [`BackgroundJobServerPublisher::force_fail_partition_head`] moves it
//! straight to a terminal failed state and frees the partition, without
//! waiting out its retry policy and without needing to hold its lease
//! first.
//!
//! See
//! `<https://github.com/mustakimali/later/blob/main/plans/SEQUENTIAL_PARTITION_MODE.md>`
//! for what is planned next.
//!
//! ## Backends
//!
//! | State | Default delivery | Scale | Features |
//! | --- | --- | --- | --- |
//! | SQLite | SQLite | Several processes sharing one local file | `sqlite` |
//! | PostgreSQL | PostgreSQL | Processes on many hosts | `postgres` |
//! Give every backend a namespace. Servers that use the same database and
//! namespace share work; different namespaces stay isolated. The SQL backends
//! also accept an application-owned SQLx pool, which allows application data
//! and a job to be written in the same transaction with `enqueue_in`.
//!
//! ## Retry behavior
//!
//! The default server policy retries six times with exponential backoff and
//! bounded jitter. Override it for the whole server through [`Config`]:
//!
//! ```no_run
//! # let context = ();
//! # let backend: Box<dyn later::backend::Backend> = todo!();
//! use later::{retry::RetryPolicy, Config};
//! use std::time::Duration;
//!
//! let config = Config::builder()
//! .context(context)
//! .backend(backend)
//! .default_retry_policy(RetryPolicy::fixed(3, Duration::from_secs(5)))
//! .build();
//! ```
//!
//! A message type can take priority over that setting:
//!
//! ```
//! use later::retry::{JobRetryPolicy, RetryPolicy};
//!
//! struct AppContext;
//! struct SendEmail;
//!
//! #[later::async_trait::async_trait]
//! impl JobRetryPolicy for SendEmail {
//! type Context = AppContext;
//!
//! async fn retry_policy(&self, _context: &AppContext) -> anyhow::Result<RetryPolicy> {
//! Ok(RetryPolicy::no_retries())
//! }
//! }
//! ```
//!
//! ## Dashboard
//! Enable feature `dashboard` to enable the experimental dashboard. The host
//! application is responsible for protecting its route. Dashboard metadata is
//! retained for one day and does not include job payload bytes. Mount the
//! response returned by `BackgroundJobServerPublisher::get_dashboard` in the
//! HTTP framework used by the application.
//!
//! The bundled page also shows stage-transition throughput from the last 60
//! seconds and workers with a heartbeat in the last 15 seconds, plus average
//! and maximum queue-dispatch wait time over that same window - how long a
//! job sat enqueued/delayed/requeued before a worker picked it up, tracked
//! separately for regular jobs and sequential/partitioned ones (a
//! sequential job's wait includes time blocked behind earlier jobs in its
//! partition, so the two are not directly comparable). This is the
//! overhead `later` itself adds before a handler starts, not time spent
//! inside the handler. These stats are stored only when the `dashboard`
//! feature is enabled - no `prometheus` feature needed. Each job response
//! includes total elapsed, handler-processing, and waiting durations computed
//! from its stage history. Job details link to a continuation's parent and to
//! at most 25 jobs that directly continue from it. The jobs table shows only
//! the last 6 characters of each ID (the full ID is still in a hover tooltip
//! and used for the detail link) - later's IDs are lexicographically
//! sortable ULIDs, so the trailing characters are what actually varies
//! between IDs minted close together.
//!
//! Every topic in [`Config::topics`] also gets a badge in the topic menu
//! showing its queue depth (not-yet-completed jobs across every
//! partition), refreshed on the same cycle as everything else on the page,
//! plus a per-partition breakdown in the partition picker once a topic is
//! selected - no separate opt-in needed. Set [`Config::retained_log_lag`]
//! (needs the `retained-log` feature too) to add each configured
//! retained-log topic's consumer lag to the same badges - see
//! [`RetainedLogLag`].
//!
//! The performance panel also shows the entire database's on-disk size (not
//! just Later's own tables) when the storage backend can report one -
//! SQLite and PostgreSQL both can; a backend without a single meaningful
//! database size (Redis, in-memory) shows a dash instead. This is the whole
//! database an application shares with Later, useful for keeping an eye on
//! total footprint even when most of it is application data.
//!
//! ### How the dashboard stays fast at any backlog
//!
//! **Queued jobs are not mirrored.** On SQLite the `enqueued` stage is read
//! from the delivery queue itself - the same rows workers claim from - not
//! from a job-index row per queued job (four b-trees written per enqueue, a
//! million rows for a million-job backlog). Its size is a counter that
//! triggers on the queue table keep exact in the same transaction as each
//! publish, claim, release and delete; the list is a cursor over the queue's
//! integer key, so a page costs the page however deep the queue. A job gets
//! an index row when it starts running, and loses it if it returns to the
//! queue. Jobs on a topic partition are not in the shared queue and stay in
//! the index (their backlog is the topic's own view). PostgreSQL keeps
//! indexing queued jobs: a shared counter row there would serialise
//! concurrent claimers.
//!
//! **The other live stages are counters too** (`delayed`, `waiting`,
//! `running`, `requeued`): triggers on the job index adjust them in the same
//! transaction as each row change, so reading them is a handful of single-row
//! reads. They follow the index's rows rather than each job's previous stage,
//! so a job written twice in the same stage, or a write the revision guard
//! rejects, cannot skew them. A recount at startup seeds both kinds of
//! counter for rows that predate the triggers, in one transaction. Finished
//! stages (`success`, `failed`) are bounded by their expiry and counted from
//! the index.
//!
//! The `count` command computes one snapshot at a time and serves it for a
//! second, so many open dashboards cost one computation. At startup a
//! background task also removes, in small batches, rows earlier dashboard
//! designs left behind (the all-jobs list, per-stage lists of queued and
//! running jobs, the old `stats-*` keys, and index rows for queued jobs).
//!
//! The Monitoring tab's **Queue length** graph is drawn from a per-minute
//! reading of the shared queue's length and the partitions' combined backlog,
//! kept for a week on the server, so it shows a backlog draining over hours
//! and survives a page reload. The plan for the remaining scaling work is in
//! `plans/DASHBOARD_SCALING.md`.
//!
//! A partitioned topic's dashboard history (the list backing its
//! per-partition "partition jobs" view) is the same `later_jobs_index` row
//! every other dashboard view reads, so it shares that row's own terminal-
//! job expiry rather than needing a separate retention setting - or call
//! [`BackgroundJobServerPublisher::truncate_partition_history`] to trim an
//! already-accumulated history down to its most recent N entries in one
//! shot regardless of expiry.
//!
//! ## Prometheus metrics
//!
//! Enable `prometheus` to add process-local worker metrics. Expose the text
//! returned by `BackgroundJobServerPublisher::get_metrics` on an HTTP route
//! that your Prometheus server can scrape. The feature records:
//!
//! - `later_job_transitions_total` for durable stage changes;
//! - `later_job_handler_duration_seconds` for handler latency and outcome per
//! worker;
//! - `later_job_duration_seconds` and `later_job_wait_duration_seconds` for
//! completed job workflow latency;
//! - `later_queue_jobs` for durable ready, leased, delayed, retry-waiting, and
//! dependency-waiting job counts;
//! - `later_worker_commands_total` for successful, failed, and discarded
//! delivery commands per worker;
//! - `later_workers` for each active and busy worker task; and
//! - for servers with configured topics: `later_partition_claims_total`
//! (claim attempts, whether or not a head was found),
//! `later_partition_heads_claimed_total` and `later_partitions_active` by
//! topic, `later_partition_lease_renewals_total` by topic and outcome,
//! `later_partition_workers_live` for this process's last-observed live
//! worker count, `later_partition_backlog` /
//! `later_partition_oldest_head_age_seconds` by topic for the ready
//! backlog and how long its oldest head has been waiting, and
//! `later_partition_queue_depth` by topic and partition for not-yet-completed
//! jobs currently queued there; and
//! - with `retained-log` also enabled and [`Config::retained_log_lag`] set:
//! `later_retained_log_consumer_lag` by topic and consumer group, for
//! records appended but not yet committed by that group.
//!
//! Worker metrics use a numeric `worker_id` that is local to one server.
//! Combine it with Prometheus's scrape `instance` label to distinguish workers
//! running in different processes. Metrics never use a job ID, error message,
//! or partition worker owner ID as a label — only the configured topic name,
//! which stays bounded the way `queue` and `job_type` already are. The one
//! exception is `later_retained_log_consumer_lag`'s `group` label: it is
//! whatever name the application commits offsets under, not a fixed set
//! `later` controls, so keep the number of distinct consumer group names
//! small the same way you would for any other Prometheus label. Values
//! cover only the current process. Prometheus combines several server
//! processes when a query sums their scraped series.
//! Queue depth is refreshed from shared storage every five seconds. Each
//! server reports the same shared value, so use `max` rather than `sum` across
//! Prometheus scrape instances.
use crateBgJobHandler;
use MqPublisher;
use Persist;
use ;
use ;
use JoinHandle;
use TypedBuilder;
pub use anyhow;
pub use async_trait;
pub use futures;
pub use background_job;
pub use instrument;
/// Traits used by job payloads and generated job handlers.
/// MessagePack encoding helpers used by Later's wire format.
/// Low-level job delivery interfaces and built-in delivery implementations.
/// Replayable records stored in the existing topic partitions.
/// Retry limits and delay strategies for failed jobs.
/// Dashboard response types.
/// Low-level state storage interfaces and built-in implementations.
/// Named topics and Kafka-like partitions for strictly ordered jobs.
/// A UTC timestamp used in persisted job records.
pub type UtcDateTime = DateTime;
/// The stable identifier assigned to one job execution.
///
/// Display the value to store it in logs or links. Applications receive IDs
/// from enqueue methods and should otherwise treat them as opaque.
;
/// The caller-selected identifier for a recurring job definition.
;
// ToDo: Remove H - use Box<dyn BgJobHandler<C>>
/// A running job server and its worker tasks.
///
/// The server dereferences to [`BackgroundJobServerPublisher`], so it can also
/// enqueue jobs. [`BackgroundJobServer::add_worker`] and
/// [`BackgroundJobServer::remove_worker`] change this server's concurrency at
/// runtime. Removing a worker lets its current job finish before it exits.
/// Call [`BackgroundJobServer::shutdown`] from a process shutdown handler to
/// drain active work before exit. Dropping the server aborts its local worker
/// and maintenance tasks; durable jobs remain in the configured backend.
/// Enqueues jobs and exposes dashboard and metrics data.
///
/// Applications normally get this through the server type generated by
/// [`background_job!`], rather than constructing it directly.
/// Generates a lexicographically sortable ULID string.
/// Bounds for [`Config::worker_autoscale`].
///
/// Every 15 seconds, a maintenance task compares the regular delivery
/// queue's current backlog (ready-to-claim messages, via
/// [`crate::mq::MqClient::queue_depth`]) against the live worker count and
/// adjusts by at most one worker per tick, using
/// [`BackgroundJobServer::add_worker`]/[`BackgroundJobServer::remove_worker`]
/// directly - the same methods available for manual scaling, so autoscaling
/// and a manual call never fight over how many workers should exist, only
/// over when to change it.
///
/// Scales up when the backlog exceeds roughly 20 messages per already-live
/// worker (so a burst is absorbed within a couple of ticks rather than
/// immediately maxing out on a brief spike) and there's room below `max`.
/// Scales down after the backlog has read as empty for three consecutive
/// ticks (about 45 seconds) and there's room above `min` - the delay avoids
/// tearing a worker down moments before the next burst arrives. One worker
/// changed per tick either way, so a large swing in load takes several
/// ticks to fully reflect, deliberately - it trades reaction speed for
/// never overshooting a spike into `max` workers at once.
/// Settings used by a generated job server builder.
///
/// Construct this with [`Config::builder`]. A backend selects both durable
/// state and delivery. The application context is shared immutably with every
/// handler.
/// Wires an optional [`retained_log::LogLagSource`] into a [`Config`],
/// gated by the `retained-log` feature so [`Config`] itself needs no
/// feature flag. Build with [`RetainedLogLag::new`], or use
/// [`Default::default`] to leave lag reporting disabled.
/// Runtime settings consumed by [`BackgroundJobServer::start`].