-
Notifications
You must be signed in to change notification settings - Fork 5
Expand file tree
/
Copy pathsys-job-queue.object.ts
More file actions
157 lines (140 loc) · 4.85 KB
/
Copy pathsys-job-queue.object.ts
File metadata and controls
157 lines (140 loc) · 4.85 KB
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
// Copyright (c) 2025 ObjectStack. Licensed under the Apache-2.0 license.
import { ObjectSchema, Field } from '@objectstack/spec/data';
/**
* sys_job_queue — Durable Background Job / Message Queue + DLQ
*
* The persistent backing store for `DbQueueAdapter`. Each row is one
* enqueued message — pending until a worker leases (`status='running'`),
* then `completed`, `failed`, or finally `dlq` after exhausting retries.
*
* DLQ rows are simply `status='dlq'` rows in this same table; no separate
* table needed. Listing/replay is a filter on `status`.
*
* Idempotency: `idempotency_key` + `queue` are deduplicated by the adapter
* over a configurable window (default 24h) — the column is indexed.
*
* Concurrency: workers claim a row by CAS-updating `status` from `pending`
* to `running` plus setting `locked_by`/`locked_until`. Reads only return
* rows whose lock has not been claimed (or whose lease expired).
*
* Writers: `DbQueueAdapter` (publish/lease/complete/fail).
* Readers: Studio DLQ view, ops dashboards, the adapter's worker loop.
*
* @namespace sys
*/
export const SysJobQueue = ObjectSchema.create({
name: 'sys_job_queue',
label: 'Job Queue Message',
pluralLabel: 'Job Queue Messages',
icon: 'inbox',
isSystem: true,
managedBy: 'engine-owned',
description: 'Durable job/message queue including dead letters',
displayNameField: 'queue',
nameField: 'queue', // [ADR-0079] canonical primary-title pointer (mirrors deprecated displayNameField)
titleFormat: '{queue} #{id}',
highlightFields: ['queue', 'status', 'attempts', 'scheduled_for', 'last_error'],
fields: {
id: Field.text({ label: 'Message ID', required: true, readonly: true, group: 'System' }),
queue: Field.text({
label: 'Queue',
required: true,
maxLength: 255,
searchable: true,
description: 'Logical queue name (snake_case)',
group: 'Identity',
}),
idempotency_key: Field.text({
label: 'Idempotency Key',
required: false,
maxLength: 255,
description: 'Deduplication key within (queue, window)',
group: 'Identity',
}),
payload_json: Field.textarea({
label: 'Payload (JSON)',
required: false,
description: 'Serialized message body',
group: 'Content',
}),
metadata_json: Field.textarea({
label: 'Metadata (JSON)',
required: false,
description: 'Serialized metadata bag (tenant_id, source_record, ...)',
group: 'Content',
}),
status: Field.select(
['pending', 'running', 'completed', 'failed', 'dlq'],
{
label: 'Status',
required: true,
defaultValue: 'pending',
description: 'Lifecycle state',
group: 'State',
},
),
priority: Field.number({
label: 'Priority',
required: false,
defaultValue: 100,
description: 'Lower = higher priority',
group: 'Schedule',
}),
attempts: Field.number({ label: 'Attempts', required: false, defaultValue: 0, group: 'State' }),
max_attempts: Field.number({ label: 'Max Attempts', required: false, defaultValue: 3, group: 'State' }),
backoff_type: Field.select(
['fixed', 'exponential'],
{ label: 'Backoff', required: false, defaultValue: 'exponential', group: 'Schedule' },
),
backoff_delay_ms: Field.number({
label: 'Backoff Base (ms)',
required: false,
defaultValue: 1000,
group: 'Schedule',
}),
backoff_max_delay_ms: Field.number({
label: 'Backoff Cap (ms)',
required: false,
group: 'Schedule',
}),
scheduled_for: Field.datetime({
label: 'Scheduled For',
required: false,
description: 'Earliest time a worker may lease this message',
group: 'Schedule',
}),
locked_by: Field.text({
label: 'Locked By',
required: false,
maxLength: 255,
description: 'Worker id holding the lease',
group: 'Lease',
}),
locked_until: Field.datetime({
label: 'Locked Until',
required: false,
description: 'Lease expiry; if past, another worker may claim',
group: 'Lease',
}),
last_error: Field.textarea({ label: 'Last Error', required: false, group: 'State' }),
completed_at: Field.datetime({ label: 'Completed At', required: false, group: 'State' }),
created_at: Field.datetime({
label: 'Created At',
required: true,
defaultValue: 'NOW()',
readonly: true,
group: 'System',
}),
updated_at: Field.datetime({ label: 'Updated At', required: false, group: 'System' }),
},
indexes: [
{ fields: ['queue', 'status', 'scheduled_for'] },
{ fields: ['idempotency_key', 'queue'] },
{ fields: ['status'] },
],
enable: {
// [ADR-0103] Engine-owned: written only by the job queue runner (SYSTEM_CTX),
// never the generic data API. Reads stay open for the Setup grid.
apiMethods: ['get', 'list'],
},
});