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
//! # Storage and SQL data model
//!
//! Applications normally choose a SQLite or PostgreSQL backend from
//! [`crate::backend`], rather than calling this module directly. Both apply
//! the bundled migrations. The tables below are Later's durable state for one
//! or more namespaces. They can share an application's database and pool,
//! provided the application does not alter Later rows.
//!
//! ## Core tables
//!
//! | Table | Stores | Notes |
//! | --- | --- | --- |
//! | `later_storage` | MessagePack job records and internal key-value state | `date_expire` makes terminal jobs and other expired values invisible. |
//! | `later_storage_range` | Ordered secondary indexes | Supports bounded pagination over jobs and dashboard indexes. Expired rows are ignored. |
//! | `later_delivery_queue` | SQL-delivery commands | Used when SQL delivery is selected. A renewable lease gives one consumer a command at a time. |
//! | `later_job_execution` | Ordinary job execution leases | Prevents two workers from executing the same non-topic job while its lease is current. |
//!
//! `namespace` separates all delivery, lease, and topic records. The
//! `later_storage` keys also include the namespace, so servers must use the
//! same namespace to share work and different namespaces to remain isolated.
//!
//! ## Topic partition tables
//!
//! These tables are present for every SQL backend but are used only when
//! [`crate::Config::topics`] is non-empty.
//!
//! | Table | Stores | Ordering role |
//! | --- | --- | --- |
//! | `later_topic` | Topic name and fixed partition count | Prevents producers with different routing layouts from sharing a topic. |
//! | `later_partition_sequence` | Next sequence for `(namespace, topic, partition)` | Enqueue assigns one monotonic sequence inside the same transaction as the job. |
//! | `later_partition_job` | The sequence-to-job mapping and its earliest runnable time | Retains every queued row for a partition. |
//! | `later_partition_head` | The current first row for every active partition | Claiming and dashboard summaries read this materialized head, avoiding a scan of retained partition rows. A delayed, retrying, or waiting head blocks later rows. |
//! | `later_partition_lease` | Current partition owner, expiry, and fencing epoch | A claim must match owner and epoch before it can retry or complete a head. A stale worker cannot advance a reclaimed head. |
//! | `later_worker` | Process heartbeat records | The live set feeds rendezvous assignment. It is an affinity hint; the partition lease is the final authority. |
//!
//! A successful partition handler removes its `later_partition_job` head,
//! advances `later_partition_head`, and stores its terminal job state in one
//! transaction. A retry keeps that head,
//! updates its runnable time, and releases it for a later claim. Lease epochs
//! increase whenever an expired partition lease is reclaimed, fencing a late
//! completion from a prior owner.

//! ## PostgreSQL throughput
//!
//! PostgreSQL keeps normal delivery commands in a leased SQL queue. When at
//! least 16 commands are ready, a consumer claims and renews a 16-command
//! batch in one query; smaller queues claim one command so newly started
//! workers can share work promptly. Initial job state and its all-jobs index
//! are also written through one CTE. Partition sequence allocation validates
//! the configured topic and allocates the next offset in one query. These are
//! implementation details, not changed delivery guarantees: writes remain
//! atomic and delivery remains at-least-once.

//! ## Retained log tables
//!
//! The `retained-log` feature adds replayable records while reusing
//! `later_topic` and its fixed partition layout. `later_log_partition` assigns
//! offsets independently from job sequences, so finishing a job cannot remove
//! a record. `later_log_record` stores immutable key and payload bytes.
//! `later_log_consumer_offset` stores each group's next offset. Member
//! heartbeats and fenced claims live in `later_log_consumer_member` and
//! `later_log_consumer_lease`.
//!
//! A group commits only while its lease owner and epoch match. Leaving a group
//! releases its leases immediately; an unresponsive member is replaced after
//! 30 seconds. Compaction keeps the newest record for a key, but only below
//! every recorded group offset. Cleanup uses the same safety boundary.
//!
//! ## Transactions and recovery
//!
//! The SQLite and PostgreSQL backends' `enqueue_in` methods let an application
//! commit its data, a normal job, and delivery intent in one caller-owned SQL
//! transaction. Their `enqueue_to_partition_in` counterparts do the same for
//! a topic-partitioned job and its sequence.
//!
//! Delivery is at-least-once. A consumer must therefore make externally visible
//! work idempotent, normally with the job ID as its key.
//! The same applies after a process crash or a lease loss: Later fences stale
//! stored-state updates, but cannot undo an external side effect that began
//! before a process stopped responding.
//!
//! The table layout is an implementation detail, not a query API. Use public
//! Later methods for changes. Migrations are additive where possible so old
//! jobs and data remain readable across library upgrades.

use crate::UtcDateTime;

/// In-process storage for tests and single-process development.
pub mod memory;
#[cfg(feature = "postgres")]
/// PostgreSQL storage.
pub mod postgres;
#[cfg(feature = "redis")]
/// Redis storage.
pub mod redis;
#[cfg(feature = "sqlite")]
/// SQLite storage.
pub mod sqlite;

mod job_index;
mod range;

#[cfg(feature = "postgres")]
pub use crate::storage::postgres::Postgres;
#[cfg(feature = "redis")]
pub use crate::storage::redis::Redis;
#[cfg(feature = "sqlite")]
pub use crate::storage::sqlite::Sqlite;
pub use job_index::{
    JobIndexChange, JobIndexMetricsBatch, JobIndexPage, JobIndexRow, JobIndexWaitStats, QueueSample,
};
pub use range::{RangeItem, RangeOrder, RangePage, StorageOperation};

#[cfg(test)]
mod tests;

#[async_trait::async_trait]
/// Atomic key-value and ordered-range persistence used by job state.
///
/// Range members are unique byte strings with stable integer cursors. Expired
/// values must behave as absent. Implementations must apply a batch atomically
/// because job state can span several keys and ranges.
pub trait Storage: Sync + Send {
    /// Gets an unexpired value, or `None` when it is absent or expired.
    async fn get(&self, key: &str) -> anyhow::Result<Option<Vec<u8>>>;

    /// Applies every mutation atomically. Either all mutations become visible
    /// or none of them do.
    async fn apply(&self, operations: Vec<StorageOperation>) -> anyhow::Result<()>;

    /// Reads at most `limit` items. `cursor` is exclusive and comes from a
    /// previous page, so work does not grow with the total range size.
    async fn range_page(
        &self,
        key: &str,
        cursor: Option<i64>,
        limit: usize,
        order: RangeOrder,
    ) -> anyhow::Result<RangePage>;

    /// Counts unexpired members in an ordered range.
    async fn range_count(&self, key: &str) -> anyhow::Result<usize>;

    /// Sets a value and clears any previous expiry.
    async fn set(&self, key: &str, value: &[u8]) -> anyhow::Result<()> {
        self.apply(vec![StorageOperation::Set {
            key: key.to_string(),
            value: value.to_vec(),
        }])
        .await
    }

    /// Deletes a value. Deleting an absent key succeeds.
    async fn del(&self, key: &str) -> anyhow::Result<()> {
        self.apply(vec![StorageOperation::Delete {
            key: key.to_string(),
        }])
        .await
    }

    /// Returns whether an unexpired value exists.
    async fn exist(&self, key: &str) -> anyhow::Result<bool> {
        Ok(self.get(key).await?.is_some())
    }

    /// Sets a time-to-live on an existing value.
    ///
    /// Expiring an absent key succeeds without creating it.
    async fn expire(&self, key: &str, ttl_seconds: usize) -> anyhow::Result<()> {
        self.apply(vec![StorageOperation::Expire {
            key: key.to_string(),
            ttl_seconds,
        }])
        .await
    }

    /// Sets a value only if the key has no unexpired value yet. Returns
    /// `true` if this call created it, `false` if a value was already
    /// there (unchanged by this call) - the compare-and-set primitive that
    /// lets a caller claim "first writer wins" for a key two concurrent
    /// callers might both be creating at once (for example, two processes
    /// registering the same brand-new recurring job identifier at the same
    /// moment), rather than both reading "absent" and both proceeding as
    /// if they were first.
    ///
    /// The default implementation is a plain get-then-set and is **not**
    /// atomic - safe only when callers are known not to race this key. All
    /// four built-in backends (SQLite, PostgreSQL, Redis, in-memory)
    /// override this with a real atomic compare-and-set; a custom
    /// [`Storage`] should too if it needs the same safety under concurrent
    /// callers.
    async fn set_if_absent(&self, key: &str, value: &[u8]) -> anyhow::Result<bool> {
        if self.exist(key).await? {
            return Ok(false);
        }
        self.set(key, value).await?;
        Ok(true)
    }

    /// Total on-disk size of the whole database this backend is connected
    /// to, in bytes - not just Later's own tables. Backs the dashboard's
    /// database-size tile, so an application sharing one database between
    /// Later and its own tables (a common setup - see
    /// [`crate::backend::SqliteBackend::enqueue_in`]) can see the real
    /// total footprint, not a number that quietly excludes most of it.
    ///
    /// `None` when a backend has no meaningful single-database size to
    /// report (for example Redis, or the in-memory backend) - the default,
    /// overridden by the SQLite and PostgreSQL backends.
    async fn total_db_size_bytes(&self) -> anyhow::Result<Option<u64>> {
        Ok(None)
    }

    /// Physically deletes up to `limit` already-expired rows total (summed
    /// across both the plain key-value table and the ordered-range table,
    /// not `limit` from each), and returns how many were removed.
    ///
    /// Expiry (`Storage::expire`, and the TTLs `later` sets internally on
    /// terminal jobs and dashboard bookkeeping) only ever makes a row
    /// invisible to reads - nothing deletes it on its own. Backends that
    /// store data in a way where that matters (SQLite, PostgreSQL) override
    /// this to actually reclaim the space; call it periodically (`later`'s
    /// own servers do, via the `PollExpiredStorage` maintenance command) or
    /// the database only ever grows. Bounded by `limit` per call so a huge
    /// backlog (built up while nothing was sweeping) can't turn one call
    /// into a multi-second write lock - loop until it returns less than
    /// `limit` to fully catch up.
    ///
    /// The default no-ops - correct for backends with no separate expired
    /// rows to reclaim (Redis expires natively; the in-memory backend's
    /// expired entries are dropped on process exit anyway).
    async fn sweep_expired(&self, _limit: usize) -> anyhow::Result<usize> {
        Ok(0)
    }

    /// Keeps the database's write-ahead log from growing without bound.
    ///
    /// Called by the periodic expired-storage poll; cheap when there is
    /// nothing to do. SQLite only resets its log when no reader is mid-read
    /// at the moment a checkpoint finishes, and a busy worker pool nearly
    /// always has one - so left alone the log kept growing (8GB beside a
    /// 1.8GB database in production). The default no-ops for backends with
    /// no such log.
    async fn checkpoint_wal(&self) -> anyhow::Result<()> {
        Ok(())
    }

    /// Reclaims disk space freed by rows [`Storage::sweep_expired`] (or an
    /// application's own deletes, for a backend sharing its database with
    /// Later) has already removed - a `DELETE` only frees pages for the
    /// database's own future reuse, it never shrinks the file on disk.
    ///
    /// Not run automatically unless [`crate::Config::vacuum_interval`] is
    /// set; see that field. The default no-ops - correct for backends with
    /// no separate file to shrink (Redis, the in-memory backend) or that
    /// already reclaim space online (PostgreSQL's autovacuum).
    async fn vacuum(&self) -> anyhow::Result<()> {
        Ok(())
    }

    /// Upserts one job's current dashboard-index row (see [`JobIndexRow`]) -
    /// called once per stage transition. Replaces what used to be a
    /// remove-from-old-list/add-to-new-list pair of range operations with a
    /// single row kept up to date in place, so a stage's listing and its
    /// count can never independently drift out of sync with each other -
    /// there is only one row per job to begin with.
    async fn job_index_upsert(&self, namespace: &str, row: JobIndexRow) -> anyhow::Result<()>;

    /// [`Storage::apply`] and [`Storage::job_index_upsert`] for one job
    /// transition. SQL backends commit both in a single transaction; the
    /// default applies them one after the other, a failed index write only
    /// logged since it affects dashboard visibility, never the job itself.
    async fn apply_with_job_index(
        &self,
        operations: Vec<StorageOperation>,
        namespace: &str,
        change: JobIndexChange,
    ) -> anyhow::Result<()> {
        self.apply(operations).await?;
        let result = match change {
            JobIndexChange::Upsert(row) => self.job_index_upsert(namespace, row).await,
            JobIndexChange::Remove(job_id) => self.job_index_remove(namespace, &job_id).await,
        };
        if let Err(error) = result {
            tracing::warn!(%error, "Could not update dashboard job index");
        }
        Ok(())
    }

    /// Deletes one job's index row, if it has one. Only called when
    /// [`Storage::queue_backs_enqueued_stage`] is true.
    async fn job_index_remove(&self, _namespace: &str, _job_id: &str) -> anyhow::Result<()> {
        Ok(())
    }

    /// Whether the dashboard's `enqueued` stage is answered from the delivery
    /// queue instead of from a job-index row per queued job.
    ///
    /// A queue of a million jobs used to be mirrored by a million index rows
    /// (four b-trees written per enqueue) purely so it could be counted and
    /// listed; the queue already holds exactly those jobs. Backends that
    /// say yes keep the queue's depth as a trigger-maintained counter and
    /// list queued jobs straight from the queue. Jobs on a topic partition
    /// are not in the shared queue and stay in the index regardless.
    fn queue_backs_enqueued_stage(&self) -> bool {
        false
    }

    /// Records the queue's length for the minute containing `at` (replacing
    /// an earlier reading of the same minute) and trims readings older than a
    /// week. Backs the dashboard's queue-length graph; the default keeps
    /// nothing, and the graph then stays empty.
    async fn queue_sample_record(
        &self,
        _namespace: &str,
        _at: UtcDateTime,
        _queued: i64,
        _partitioned: i64,
    ) -> anyhow::Result<()> {
        Ok(())
    }

    /// Readings from `since` on, oldest first.
    async fn queue_samples_since(
        &self,
        _namespace: &str,
        _since: UtcDateTime,
    ) -> anyhow::Result<Vec<QueueSample>> {
        Ok(Vec::new())
    }

    /// Removes up to `limit` rows the dashboard no longer reads - leftovers
    /// of earlier designs - and returns how many it removed; zero means there
    /// is nothing left. Called repeatedly in small batches at startup so no
    /// single delete holds the write lock for long.
    async fn purge_unused_dashboard_rows(
        &self,
        _namespace: &str,
        _limit: usize,
    ) -> anyhow::Result<usize> {
        Ok(0)
    }

    /// Every stage's live (not expired) job count for `namespace`, in one
    /// call - backs the dashboard's stage-count panel. Only stages with at
    /// least one current job are present in the result.
    ///
    /// SQL backends read the live stages from `later_stage_counts` (O(1),
    /// kept exact by triggers on the index and seeded by
    /// [`Storage::job_index_reconcile_stage_counts`]) and the finished stages
    /// from the index; memory and Redis count directly.
    async fn job_index_stage_counts(
        &self,
        namespace: &str,
    ) -> anyhow::Result<std::collections::HashMap<String, usize>>;

    /// A page of `namespace`'s jobs currently in `stage`, most-recently
    /// transitioned first. `cursor` is exclusive, from a previous page.
    async fn job_index_list_by_stage(
        &self,
        namespace: &str,
        stage: &str,
        cursor: Option<i64>,
        limit: usize,
    ) -> anyhow::Result<JobIndexPage>;

    /// A page of every job ever enqueued to one `(topic, partition)`,
    /// highest sequence (newest) first - backs the dashboard's per-partition
    /// view.
    async fn job_index_list_by_partition(
        &self,
        namespace: &str,
        topic: &str,
        partition: u32,
        cursor: Option<i64>,
        limit: usize,
    ) -> anyhow::Result<JobIndexPage>;

    /// The job holding `sequence - 1` (older) or `sequence + 1` (newer) in
    /// `(topic, partition)`, if any - backs a job detail page's
    /// previous/next-in-partition links. `older` selects which direction.
    async fn job_index_partition_neighbor(
        &self,
        namespace: &str,
        topic: &str,
        partition: u32,
        sequence: i64,
        older: bool,
    ) -> anyhow::Result<Option<JobIndexRow>>;

    /// Every job currently `Stage::Waiting` on `parent_job_id`, oldest
    /// first, up to `limit` - backs a job detail page's continuation list.
    /// Bounded, not paginated (matches the existing convention of asking
    /// for one more than needed and truncating, to detect "more exist").
    async fn job_index_list_continuations(
        &self,
        namespace: &str,
        parent_job_id: &str,
        limit: usize,
    ) -> anyhow::Result<Vec<JobIndexRow>>;

    /// Deletes the oldest rows for `(topic, partition)` beyond `keep_last`.
    /// Returns the number removed.
    async fn job_index_truncate_partition(
        &self,
        namespace: &str,
        topic: &str,
        partition: u32,
        keep_last: usize,
    ) -> anyhow::Result<usize>;

    /// Records one transition into `stage` at `at`, for the throughput
    /// panel - see `later_job_transition_metrics`'s own doc comment on why
    /// this is a small aggregate counter rather than one row per
    /// transition.
    async fn job_index_record_transition(
        &self,
        namespace: &str,
        stage: &str,
        at: UtcDateTime,
    ) -> anyhow::Result<()>;

    /// Adds a pre-aggregated batch of transition and wait samples, in one
    /// round trip where the backend allows it. Equivalent to calling
    /// [`Storage::job_index_record_transition`] and
    /// [`Storage::job_index_record_wait`] once per sample.
    async fn job_index_record_metrics_batch(
        &self,
        namespace: &str,
        batch: &JobIndexMetricsBatch,
    ) -> anyhow::Result<()>;

    /// Recounts the live stages (`delayed`, `waiting`, `enqueued`, `running`,
    /// `requeued`) from the job index and replaces the counters
    /// [`Storage::job_index_stage_counts`] reads them from, correcting drift
    /// from a crash between a transition and its flush. Backends that count
    /// directly (memory, Redis) have nothing to correct.
    async fn job_index_reconcile_stage_counts(&self, _namespace: &str) -> anyhow::Result<()> {
        Ok(())
    }

    /// Sum of `stage` transitions recorded in roughly the last
    /// `RECENT_WINDOW_SECONDS` (an approximation - see the migration's own
    /// comment on bucket granularity).
    async fn job_index_recent_transition_count(
        &self,
        namespace: &str,
        stage: &str,
        now: UtcDateTime,
    ) -> anyhow::Result<usize>;

    /// Records one queue-wait sample for `mode` ("regular"/"sequential") at
    /// `at` - see `later_job_wait_metrics`.
    async fn job_index_record_wait(
        &self,
        namespace: &str,
        mode: &str,
        wait_ms: i64,
        at: UtcDateTime,
    ) -> anyhow::Result<()>;

    /// Aggregate queue-wait stats for `mode` over roughly the last
    /// `RECENT_WINDOW_SECONDS`.
    async fn job_index_recent_wait_stats(
        &self,
        namespace: &str,
        mode: &str,
        now: UtcDateTime,
    ) -> anyhow::Result<JobIndexWaitStats>;

    /// Physically deletes up to `limit` already-expired `later_jobs_index`
    /// rows for `namespace`, plus stale metric buckets older than a couple
    /// of minutes - the job-index counterpart to [`Storage::sweep_expired`],
    /// run by the same periodic poll. Returns the number of index rows
    /// removed (metric-bucket cleanup isn't counted, since it isn't bounded
    /// by `limit` the same way - the bucket tables are already tiny).
    async fn job_index_sweep_expired(
        &self,
        namespace: &str,
        now: UtcDateTime,
        limit: usize,
    ) -> anyhow::Result<usize>;
}