Skip to content

Commit e98fb14

Browse files
os-zhuangclaude
andauthored
fix(service-queue,platform-objects): sys_job_queue 的 completed 行按声明式 retention 到期即清(#5179) (#5192)
* fix(service-queue,platform-objects): bound sys_job_queue — completed rows expire on a declared ADR-0057 retention (#5179) DbQueueAdapter marked delivered messages `completed` and nothing ever touched the row again: `purge()` had zero production callers, `purgeFailed()` is a manual dead-letter API, and the object declared no lifecycle policy — so the queue table only ever grew (one permanent row per queued email since #5160). sys_job_queue now declares `lifecycle: { class: 'transient', retention: { maxAge: '7d', onlyWhen: { status: 'completed' } } }`, enforced by the one platform-owned LifecycleService reaper (ADR-0057 §3.3) on its existing hourly sweep — no new sweeper in the adapter's poll loop, no new configuration. `pending`/`running` (live work) and `failed`/`dlq` (the dead-letter queue) are never swept at any age. The dedup window becomes an enforced invariant rather than a coincidence: publish dedups terminal rows by `created_at` against `idempotencyWindowMs`, the reaper cuts off on the same axis, and DbQueueAdapter now reads the declared window (`completedRetentionWindowMs()`) and throws at construction if the idempotency window is configured longer than it. `class: 'transient'` and not `telemetry`: per ADR-0057 §3.6 a telemetry/event/audit class relocates the table to the dedicated `telemetry` datasource wherever one is registered, and moving a live work queue's storage would be a migration, not a cleanup. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017MCKJaEomEqg4tvz4SzdNd * test(service-queue): pin the retention fake engine's delete() to ObjectQL's own dispatch (#4550) `check:engine-double-contract` flagged the new fake engine in job-queue-retention.test.ts: its `delete()` hand-mirrored the engine's guard (`if (opts?.where?.id == null) throw`) instead of routing through `assertEngineDeleteDispatch`. A mirror is looser than the producer on exactly the shape a copy always drops — `where: { id: { $in: [...] } }` reads as an id and is a multi-row predicate the real engine rejects without `multi` — and a double looser than the engine it stands in for is how #4434 shipped a dead REST route with its suite green. Routes through the producer's predicate, same shape as the other 14 pinned doubles, and adds the `@objectstack/objectql` devDependency the import needs (the precedent set in plugin-email by b169f21, and in plugin-approvals / plugin-sharing before it). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017MCKJaEomEqg4tvz4SzdNd --------- Co-authored-by: Claude <noreply@anthropic.com>
1 parent 7f955e5 commit e98fb14

9 files changed

Lines changed: 549 additions & 3 deletions

File tree

Lines changed: 52 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,52 @@
1+
---
2+
"@objectstack/platform-objects": minor
3+
"@objectstack/service-queue": minor
4+
---
5+
6+
fix(service-queue): `sys_job_queue` no longer grows forever — `completed` rows expire on a declared 7-day retention (#5179)
7+
8+
`DbQueueAdapter` marked a delivered message `status: 'completed'` and then
9+
**nothing ever touched that row again**. `purge()` had zero production callers
10+
(tests only), `purgeFailed()` is a manual dead-letter API, and the object
11+
declared no lifecycle policy at all — so every queue delivery left a permanent
12+
row, which since #5160 means one permanent row per queued email.
13+
14+
`sys_job_queue` now declares an ADR-0057 policy and the platform
15+
`LifecycleService` enforces it on its existing hourly sweep:
16+
17+
```ts
18+
lifecycle: {
19+
class: 'transient',
20+
retention: { maxAge: '7d', onlyWhen: { status: 'completed' } },
21+
}
22+
```
23+
24+
**Only `completed` rows are swept.** `pending` / `running` are live work, and
25+
`failed` / `dlq` are the dead-letter queue — they exist to wait for a human, so
26+
they are never deleted automatically at any age. `listFailed()` / `replay()` /
27+
`purgeFailed()` remain the only way a dead letter leaves the table. This is
28+
also why the policy is `retention` (age + row filter) rather than a `ttl` on
29+
`completed_at`: TTL has no row filter, and `dlq` rows stamp `completed_at` too.
30+
31+
**No new configuration, and no new sweeper.** ADR-0057 §3.3 puts one reaper in
32+
the platform rather than one per plugin — the same call the sibling
33+
`sys_job_run` (30d) already makes. Any kernel with a data engine already runs
34+
it, its per-sweep `[lifecycle] sweep: … ~N rows reaped` line now accounts for
35+
this table too, and the window is overridable per environment through the
36+
`lifecycle` settings namespace without touching code.
37+
38+
**The dedup window is now an enforced invariant, not a coincidence.** Publish
39+
dedups against a terminal row by comparing its `created_at` to
40+
`idempotencyWindowMs` (default 24h), and the reaper cuts off on that same
41+
`created_at` axis — so retention (7d) ≥ dedup window is what keeps "duplicate
42+
publishes inside the window are suppressed" true. `DbQueueAdapter` reads the
43+
declared window (new export `completedRetentionWindowMs()`) and **throws at
44+
construction** if `idempotencyWindowMs` is configured longer than it, instead of
45+
silently degrading into duplicate deliveries days later. If you raise
46+
`idempotencyWindowMs` past 7 days, raise the object's declared retention (or the
47+
`lifecycle` settings override) to match — the error message names both numbers.
48+
49+
`class: 'transient'` is deliberate: `telemetry`/`event`/`audit` classes
50+
relocate their table to the dedicated `telemetry` datasource wherever one is
51+
registered (ADR-0057 §3.6), and moving a live work queue's storage would be a
52+
migration, not a cleanup.

‎packages/platform-objects/src/audit/sys-job-queue.object.ts‎

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,9 @@ import { ObjectSchema, Field } from '@objectstack/spec/data';
2222
* Writers: `DbQueueAdapter` (publish/lease/complete/fail).
2323
* Readers: Studio DLQ view, ops dashboards, the adapter's worker loop.
2424
*
25+
* Retention: `completed` rows are swept by the platform LifecycleService —
26+
* see the `lifecycle` block below (#5179).
27+
*
2528
* @namespace sys
2629
*/
2730
export const SysJobQueue = ObjectSchema.create({
@@ -31,6 +34,56 @@ export const SysJobQueue = ObjectSchema.create({
3134
icon: 'inbox',
3235
isSystem: true,
3336
managedBy: 'engine-owned',
37+
38+
/**
39+
* [ADR-0057 §3.1/§3.3, #5179] The queue table only ever GREW: the adapter
40+
* marks a delivered message `completed` and nothing ever touched the row
41+
* again (`purge()` had zero production callers, `purgeFailed()` is a manual
42+
* dead-letter API). Since #5160 that is one permanent row per email.
43+
*
44+
* Bounded declaratively rather than by a sweeper inside `DbQueueAdapter`:
45+
* ADR-0057 §3.3 puts ONE reaper in the platform (`LifecycleService`), not N
46+
* per-plugin ones — the same call the sibling `sys_job_run` already makes.
47+
* That the writer is the adapter itself (never user data) is what makes an
48+
* unattended delete safe here; the declaration is where an operator can see
49+
* the window, and `lifecycle` settings can override it per environment
50+
* without a code change.
51+
*
52+
* `onlyWhen: { status: 'completed' }` is the whole safety story:
53+
* - `pending` / `running` are LIVE work — reaping them would drop
54+
* undelivered messages;
55+
* - `dlq` / `failed` are the dead-letter surface and exist precisely to
56+
* wait for a human (`listFailed` / `replay` / `purgeFailed`), so they
57+
* are never swept automatically, at any age.
58+
* This is also why the policy is `retention` (age by `created_at` + row
59+
* filter) and not `ttl` on `completed_at`: TTL has no row filter, and `dlq`
60+
* rows stamp `completed_at` too — a TTL would eat the dead-letter queue.
61+
*
62+
* Window = 7d, and it MUST stay ≥ the adapter's idempotency window
63+
* (`DbQueueAdapterOptions.idempotencyWindowMs`, default 24h): publish
64+
* dedups against terminal rows by comparing `created_at` to that window
65+
* (`db-queue-adapter.ts`), and the Reaper cuts off on the very same
66+
* `created_at` axis — so a retention ≥ the dedup window means a row the
67+
* dedup check still needs can never have been reaped, with no clock skew
68+
* between the two rules. 7d gives a week of delivery history for debugging
69+
* and 7× headroom over the default dedup window. `DbQueueAdapter` reads
70+
* this declaration and refuses to start when the two are configured the
71+
* wrong way round, so the invariant cannot drift apart silently.
72+
*
73+
* `class: 'transient'` ("workflow / ephemeral state" — ADR-0057 §3.1), not
74+
* `telemetry`: this is live work state, not a log, and per §3.6 a
75+
* `telemetry`/`event`/`audit` class RELOCATES the table to the dedicated
76+
* `telemetry` datasource wherever one is registered. Moving a live queue's
77+
* store is a migration, not a cleanup — `transient` deliberately stays on
78+
* the primary.
79+
*/
80+
lifecycle: {
81+
class: 'transient',
82+
retention: {
83+
maxAge: '7d',
84+
onlyWhen: { status: 'completed' },
85+
},
86+
},
3487
description: 'Durable job/message queue including dead letters',
3588
displayNameField: 'queue',
3689
nameField: 'queue', // [ADR-0079] canonical primary-title pointer (mirrors deprecated displayNameField)

‎packages/services/service-queue/README.md‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,29 @@ new QueueServicePlugin({
7070
new QueueServicePlugin({ adapter: 'memory' });
7171
```
7272

73+
### Retention — how `sys_job_queue` stays bounded
74+
75+
Delivered messages are not kept forever. `sys_job_queue` declares an ADR-0057
76+
lifecycle policy and the platform `LifecycleService` (shipped with
77+
`@objectstack/objectql`, armed on every kernel that has data) enforces it — no
78+
configuration, no extra scheduler:
79+
80+
| Row state | What happens |
81+
|---|---|
82+
| `completed` | deleted **7 days** after `created_at` |
83+
| `pending` / `running` | never swept — live work |
84+
| `failed` / `dlq` | never swept — the dead-letter queue waits for a human (`listFailed` / `replay` / `purgeFailed`) |
85+
86+
Two consequences worth knowing:
87+
88+
- **`idempotencyWindowMs` must not exceed the retention window.** Dedup against
89+
a terminal message compares its `created_at` to that window, so a longer
90+
setting would start accepting duplicates the moment the row was swept. The
91+
`db` adapter throws at construction instead of degrading quietly.
92+
- **The window is overridable per environment** through the `lifecycle`
93+
settings namespace (`maxAge` per object), like every other ADR-0057 policy.
94+
Keep it ≥ your idempotency window.
95+
7396
## Service API
7497

7598
Implements `IQueueService` from `@objectstack/spec/contracts`:

‎packages/services/service-queue/package.json‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424
"@objectstack/spec": "workspace:*"
2525
},
2626
"devDependencies": {
27+
"@objectstack/objectql": "workspace:*",
2728
"@types/node": "^26.1.2",
2829
"typescript": "^6.0.3",
2930
"vitest": "^4.1.10"

‎packages/services/service-queue/src/common.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,39 @@ export function nowIso(clock?: JobClock): string {
3636
return (clock?.now() ?? new Date()).toISOString();
3737
}
3838

39+
/**
40+
* Milliseconds per ADR-0057 lifecycle duration unit. Mirrors
41+
* `parseLifecycleDuration` in `@objectstack/objectql` (the canonical runtime
42+
* consumer), reproduced here rather than imported because the queue adapters
43+
* deliberately do not depend on the engine package — they duck-type
44+
* {@link JobEngine} so they stay testable without booting a kernel. Both
45+
* tables are fixed by the ADR (coarse operational bounds: `y` is 365 days),
46+
* and `job-queue-retention.test.ts` pins this one against them.
47+
*/
48+
const LIFECYCLE_UNIT_MS: Record<string, number> = {
49+
h: 3_600_000,
50+
d: 86_400_000,
51+
w: 7 * 86_400_000,
52+
y: 365 * 86_400_000,
53+
};
54+
55+
/**
56+
* Parse an ADR-0057 duration literal (`'6h'`, `'7d'`, `'12w'`, `'7y'`) into
57+
* milliseconds. Throws on anything else: declarations reach this code already
58+
* validated by `LifecycleSchema`, so a failure here is a broken declaration,
59+
* not user input — and a queue that silently guessed a window would be exactly
60+
* the silent behaviour #5179 is about.
61+
*/
62+
export function lifecycleDurationMs(literal: string): number {
63+
const m = /^(\d+)(h|d|w|y)$/.exec(literal);
64+
if (!m) {
65+
throw new Error(
66+
`[service-queue] invalid lifecycle duration literal '${literal}' — expected <n><unit> with unit h|d|w|y (e.g. '7d')`,
67+
);
68+
}
69+
return Number(m[1]) * LIFECYCLE_UNIT_MS[m[2]!]!;
70+
}
71+
3972
export function parseJson<T = unknown>(raw: unknown, fallback?: T): T | undefined {
4073
if (raw == null) return fallback;
4174
if (typeof raw === 'string') {

‎packages/services/service-queue/src/db-queue-adapter.ts‎

Lines changed: 72 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,26 +7,59 @@ import type {
77
QueueMessageRecord,
88
QueueHandler,
99
} from '@objectstack/spec/contracts';
10+
import { SysJobQueue } from '@objectstack/platform-objects/audit';
1011
import {
1112
SYSTEM_CTX,
1213
uid,
1314
nowIso,
1415
parseJson,
16+
lifecycleDurationMs,
1517
type JobEngine,
1618
type JobClock,
1719
type JobLogger,
1820
} from './common.js';
1921

2022
const QUEUE_TABLE = 'sys_job_queue';
2123

24+
/**
25+
* How long a `completed` row survives before the platform Reaper deletes it.
26+
*
27+
* Read from the object's own ADR-0057 declaration
28+
* (`sys_job_queue.lifecycle.retention`, #5179) instead of being a second
29+
* number here: the declaration is what actually runs (LifecycleService sweeps
30+
* every registered object hourly), so a copy in this file could only ever be
31+
* a copy that drifts. A missing or unparseable declaration throws: the queue's
32+
* dedup contract below is defined against this window, so "no window" is not a
33+
* state the adapter can run in.
34+
*/
35+
export function completedRetentionWindowMs(): number {
36+
const maxAge = SysJobQueue.lifecycle?.retention?.maxAge;
37+
if (!maxAge) {
38+
throw new Error(
39+
'[service-queue] sys_job_queue no longer declares lifecycle.retention — DbQueueAdapter dedups against '
40+
+ 'terminal rows by `created_at` window and relies on that declared retention to keep them (ADR-0057, #5179). '
41+
+ 'Restore the declaration in @objectstack/platform-objects rather than sweeping the table from here.',
42+
);
43+
}
44+
return lifecycleDurationMs(maxAge);
45+
}
46+
2247
export interface DbQueueAdapterOptions {
2348
/** Polling interval for the worker loop (ms, default 1000) */
2449
pollIntervalMs?: number;
2550
/** Max messages claimed per poll tick (default 10) */
2651
batchSize?: number;
2752
/** Lease duration before another worker may reclaim (ms, default 30000) */
2853
leaseMs?: number;
29-
/** Idempotency window — how long the same key blocks re-publish (ms, default 24h) */
54+
/**
55+
* Idempotency window — how long the same key blocks re-publish (ms, default 24h).
56+
*
57+
* Must not exceed `sys_job_queue`'s declared retention for `completed` rows
58+
* ({@link completedRetentionWindowMs}, 7d): the window is evaluated against
59+
* rows that are still in the table, so a longer window would silently start
60+
* accepting duplicates as soon as the Reaper swept the row it dedups
61+
* against. The constructor rejects that configuration (#5179).
62+
*/
3063
idempotencyWindowMs?: number;
3164
/** Default maxAttempts when publish doesn't specify (default 3) */
3265
defaultMaxAttempts?: number;
@@ -52,6 +85,15 @@ interface RegisteredHandler {
5285
* Idempotency: publish suppresses duplicates within a configurable
5386
* window when `(queue, idempotencyKey)` is non-null.
5487
*
88+
* Retention: this adapter does NOT sweep the table. `completed` rows are
89+
* bounded by `sys_job_queue`'s declared ADR-0057 retention (7d, filtered to
90+
* `status='completed'`), enforced by the one platform-owned
91+
* `LifecycleService` reaper — see the object definition in
92+
* `@objectstack/platform-objects` and {@link completedRetentionWindowMs}.
93+
* `dlq`/`failed` rows are never swept; they are the dead-letter surface
94+
* ({@link DbQueueAdapter.listFailed} / {@link DbQueueAdapter.replay} /
95+
* {@link DbQueueAdapter.purgeFailed}).
96+
*
5597
* Designed for SQLite and Postgres alike — uses CAS via WHERE-clauses,
5698
* not row-level locking.
5799
*/
@@ -84,6 +126,25 @@ export class DbQueueAdapter implements IQueueService {
84126
autoStart: o.autoStart ?? true,
85127
workerId: o.workerId ?? uid('worker'),
86128
};
129+
130+
// [#5179] The dedup window only means anything while the row it dedups
131+
// against still exists. `completed` rows now expire on the declared
132+
// retention window, so an idempotency window LONGER than it would quietly
133+
// degrade into "dedup for as long as the Reaper happens not to have run" —
134+
// duplicate deliveries appearing days later, with nothing in any log. The
135+
// two windows are ordered here, at construction, rather than tolerated at
136+
// publish time: the fix is a config or declaration change, and both are
137+
// named in the message.
138+
const retentionMs = completedRetentionWindowMs();
139+
if (this.opts.idempotencyWindowMs > retentionMs) {
140+
throw new Error(
141+
`[service-queue] idempotencyWindowMs (${this.opts.idempotencyWindowMs}ms) exceeds the retention window `
142+
+ `sys_job_queue declares for completed rows (${retentionMs}ms, lifecycle.retention.maxAge — ADR-0057). `
143+
+ 'Terminal-row dedup is evaluated by `created_at` against that same window, so the longer setting would '
144+
+ 'silently accept duplicates once a row is reaped. Lower idempotencyWindowMs, or raise the declared '
145+
+ 'retention (both windows are measured from `created_at`).',
146+
);
147+
}
87148
}
88149

89150
// ── IQueueService ────────────────────────────────────────────────
@@ -96,7 +157,16 @@ export class DbQueueAdapter implements IQueueService {
96157
const opts = options ?? {};
97158
const now = this.now();
98159

99-
// Idempotency check
160+
// Idempotency check.
161+
//
162+
// [#5179] This is the reason `sys_job_queue`'s retention is filtered and
163+
// generous rather than aggressive: a terminal (`completed`/`dlq`) row
164+
// blocks a re-publish only while its `created_at` is inside the
165+
// idempotency window, so the row must SURVIVE that long. The declared
166+
// retention (7d on `completed`, nothing on `dlq`) is measured on the very
167+
// same `created_at` axis and is ≥ this window — enforced in the
168+
// constructor — which makes "the reaper deleted a row the dedup check
169+
// needed" unrepresentable rather than merely unlikely.
100170
if (opts.idempotencyKey) {
101171
const windowStart = new Date(now.getTime() - this.opts.idempotencyWindowMs).toISOString();
102172
const existing = await this.engine.find(QUEUE_TABLE, {

‎packages/services/service-queue/src/index.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,6 @@ export { QueueServicePlugin } from './queue-service-plugin.js';
44
export type { QueueServicePluginOptions } from './queue-service-plugin.js';
55
export { MemoryQueueAdapter } from './memory-queue-adapter.js';
66
export type { MemoryQueueAdapterOptions } from './memory-queue-adapter.js';
7-
export { DbQueueAdapter } from './db-queue-adapter.js';
7+
export { DbQueueAdapter, completedRetentionWindowMs } from './db-queue-adapter.js';
88
export type { DbQueueAdapterOptions } from './db-queue-adapter.js';
99
export type { JobEngine, JobClock, JobLogger } from './common.js';

0 commit comments

Comments
 (0)