crafty_proto/queue.rs
1//! Job queue wire types ([job-queue](../../../docs/decisions/job-queue.md)).
2
3use serde::{Deserialize, Serialize};
4
5/// Enqueue a job on stream `stream` (`POST /raft/v1/queue/enqueue`).
6#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
7pub struct QueueEnqueueRequest {
8 /// Logical queue stream name (e.g. `"jobs"` or sharded `"jobs~0"`).
9 pub stream: String,
10 /// Opaque job body handed to workers after lease.
11 pub payload: Vec<u8>,
12 /// Higher values are leased before lower (default `0`).
13 #[serde(default)]
14 pub priority: u8,
15 /// Earliest wall time (unix ms) the job may be leased; `0` = immediately.
16 #[serde(default)]
17 pub not_before_ms: u64,
18 /// Optional routing key for sharded streams (defaults to hashing `payload`).
19 #[serde(default)]
20 pub shard_key: Option<Vec<u8>>,
21 /// Idempotency key — retries return the same `job_id` while the job exists.
22 #[serde(default)]
23 pub dedup_key: Option<Vec<u8>>,
24 /// Maximum delivery attempts before dead letter (`0` = unlimited).
25 #[serde(default)]
26 pub max_attempts: u32,
27}
28
29/// Response to [`QueueEnqueueRequest`].
30#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
31pub struct QueueEnqueueReply {
32 /// Assigned job id when enqueue succeeded.
33 pub job_id: Option<u64>,
34 /// Human-readable error when enqueue failed.
35 pub error: Option<String>,
36}
37
38/// One job in a batch enqueue (`POST /raft/v1/queue/enqueue-batch`).
39#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
40pub struct QueueBatchEnqueueJob {
41 /// Opaque job body.
42 pub payload: Vec<u8>,
43 /// Higher values are leased before lower (default `0`).
44 #[serde(default)]
45 pub priority: u8,
46 /// Earliest wall time (unix ms) the job may be leased; `0` = immediately.
47 #[serde(default)]
48 pub not_before_ms: u64,
49 /// Optional routing key for sharded streams.
50 #[serde(default)]
51 pub shard_key: Option<Vec<u8>>,
52 /// Idempotency key — retries return the same `job_id` while the job exists.
53 #[serde(default)]
54 pub dedup_key: Option<Vec<u8>>,
55 /// Maximum delivery attempts before dead letter (`0` = unlimited).
56 #[serde(default)]
57 pub max_attempts: u32,
58}
59
60/// Enqueue many jobs in one leader transaction (`POST /raft/v1/queue/enqueue-batch`).
61#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
62pub struct QueueEnqueueBatchRequest {
63 /// Logical queue stream name.
64 pub stream: String,
65 /// Jobs to append (leader caps batch size).
66 pub jobs: Vec<QueueBatchEnqueueJob>,
67}
68
69/// Response to [`QueueEnqueueBatchRequest`].
70#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
71pub struct QueueEnqueueBatchReply {
72 /// Assigned ids in the same order as the request (dedup hits echo existing ids).
73 pub job_ids: Vec<u64>,
74 /// Set when the batch failed before any job was committed.
75 pub error: Option<String>,
76}
77
78/// Lease jobs for a worker (`POST /raft/v1/queue/lease`).
79#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
80pub struct QueueLeaseRequest {
81 /// Queue stream to pull from.
82 pub stream: String,
83 /// [`NodeId`](crate::NodeId) of the leasing worker (`.0` wire encoding).
84 pub worker_node: u64,
85 /// Worker actor instance id on that node.
86 pub worker_instance: u32,
87 /// Maximum jobs to lease in one call.
88 pub max: usize,
89}
90
91/// One job returned under lease on the wire.
92#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
93pub struct QueueLeasedJobWire {
94 /// Lease token — required for ack/nack.
95 pub lease_id: u64,
96 /// Job id within the stream.
97 pub job_id: u64,
98 /// Job body copied at enqueue time.
99 pub payload: Vec<u8>,
100}
101
102/// Response to [`QueueLeaseRequest`].
103#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
104pub struct QueueLeaseReply {
105 /// Leased jobs (may be empty when the queue is idle).
106 pub jobs: Vec<QueueLeasedJobWire>,
107 /// Set when the lease RPC failed.
108 pub error: Option<String>,
109}
110
111/// Acknowledge successful processing (`POST /raft/v1/queue/ack`).
112#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
113pub struct QueueAckRequest {
114 /// Queue stream the lease belongs to.
115 pub stream: String,
116 /// Leasing worker node id.
117 pub worker_node: u64,
118 /// Leasing worker instance id.
119 pub worker_instance: u32,
120 /// Lease token from [`QueueLeasedJobWire::lease_id`].
121 pub lease_id: u64,
122}
123
124/// Response to [`QueueAckRequest`].
125#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
126pub struct QueueAckReply {
127 /// Set when ack failed (unknown lease, wrong worker, etc.).
128 pub error: Option<String>,
129}
130
131/// Acknowledge many leased jobs in one leader transaction
132/// (`POST /raft/v1/queue/ack-batch`).
133#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
134pub struct QueueAckBatchRequest {
135 /// Queue stream the leases belong to.
136 pub stream: String,
137 /// Leasing worker node id.
138 pub worker_node: u64,
139 /// Leasing worker instance id.
140 pub worker_instance: u32,
141 /// Lease tokens from [`QueueLeasedJobWire::lease_id`].
142 pub lease_ids: Vec<u64>,
143}
144
145/// Response to [`QueueAckBatchRequest`].
146#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
147pub struct QueueAckBatchReply {
148 /// Set when the batch ack failed.
149 pub error: Option<String>,
150}
151
152/// Return a leased job to pending immediately (`POST /raft/v1/queue/nack`).
153#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
154pub struct QueueNackRequest {
155 /// Queue stream the lease belongs to.
156 pub stream: String,
157 /// Leasing worker node id.
158 pub worker_node: u64,
159 /// Leasing worker instance id.
160 pub worker_instance: u32,
161 /// Lease token to release.
162 pub lease_id: u64,
163}
164
165/// Response to [`QueueNackRequest`].
166#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
167pub struct QueueNackReply {
168 /// Set when nack failed.
169 pub error: Option<String>,
170}
171
172/// Read queue depth gauges (`POST /raft/v1/queue/metrics`).
173#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
174pub struct QueueMetricsRequest {
175 /// Stream to inspect.
176 pub stream: String,
177}
178
179/// Depth and age gauges for autoscale / observability.
180#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
181pub struct QueueMetricsReply {
182 /// Jobs waiting to be leased.
183 pub pending: u64,
184 /// Jobs currently leased to workers.
185 pub leased: u64,
186 /// Jobs in the dead-letter set (exhausted retries).
187 pub dead_letter: u64,
188 /// Age in ms of the oldest ready pending job (`0` when empty).
189 pub oldest_pending_age_ms: u64,
190 /// Set when metrics collection failed.
191 pub error: Option<String>,
192}
193
194/// Job lifecycle returned by [`QueueJobStatusReply`].
195#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
196#[repr(u8)]
197pub enum QueueJobLifecycleWire {
198 /// Waiting in pending (ready to lease).
199 Pending = 0,
200 /// Currently leased to a worker.
201 Leased = 1,
202 /// Delayed until `not_before`.
203 Delayed = 2,
204 /// Exhausted retry budget — not leased until requeued by an operator.
205 DeadLetter = 3,
206}
207
208/// Lookup job metadata by id (`POST /raft/v1/queue/job-status`).
209#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
210pub struct QueueJobStatusRequest {
211 /// Stream to inspect.
212 pub stream: String,
213 /// Job id within the stream (global id when sharded).
214 pub job_id: u64,
215}
216
217/// Metadata for a single job (`POST /raft/v1/queue/job-status`).
218#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
219pub struct QueueJobStatusReply {
220 /// `true` when the job exists (pending, leased, or delayed).
221 pub found: bool,
222 /// Echo of the requested id.
223 pub job_id: u64,
224 /// Set when [`Self::found`] is true.
225 pub lifecycle: Option<QueueJobLifecycleWire>,
226 /// Byte length of stored payload.
227 pub payload_len: u64,
228 /// Enqueue priority.
229 pub priority: u8,
230 /// Worker node when leased.
231 pub leased_worker_node: Option<u64>,
232 /// Worker instance when leased.
233 pub leased_worker_instance: Option<u32>,
234 /// Delivery attempts so far (including the attempt that dead-lettered).
235 pub attempts: u32,
236 /// Configured retry ceiling (`0` = unlimited).
237 pub max_attempts: u32,
238 /// Set when lookup failed.
239 pub error: Option<String>,
240}
241
242const fn default_true() -> bool {
243 true
244}
245
246/// Cron-driven recurring job registered on a queue stream.
247#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
248pub struct RecurringScheduleWire {
249 /// Unique schedule name within the stream.
250 pub name: String,
251 /// Cron expression (5-field `min hour dom month dow` or 6-field with seconds).
252 pub cron: String,
253 /// Payload enqueued on each tick.
254 pub payload: Vec<u8>,
255 /// Passed to [`QueueEnqueueRequest::priority`].
256 #[serde(default)]
257 pub priority: u8,
258 /// Passed to [`QueueEnqueueRequest::max_attempts`].
259 #[serde(default)]
260 pub max_attempts: u32,
261 /// When false the schedule is stored but does not fire.
262 #[serde(default = "default_true")]
263 pub enabled: bool,
264 /// Next fire time (unix ms); leader-maintained.
265 #[serde(default)]
266 pub next_run_ms: u64,
267}
268
269/// Idempotent state transition replicated from the queue leader to every voter
270/// (`POST /raft/v1/queue/replicate`).
271#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
272pub enum QueueReplicateOp {
273 /// Append a job and advance the stream's `next_job_id`.
274 Enqueue {
275 /// Assigned job id.
276 job_id: u64,
277 /// Job body.
278 payload: Vec<u8>,
279 /// Leader wall time at enqueue (unix ms).
280 enqueued_at_ms: u64,
281 /// Monotonic id generator after this enqueue.
282 next_job_id: u64,
283 #[serde(default)]
284 /// Lease priority (higher first).
285 priority: u8,
286 #[serde(default)]
287 /// Earliest lease time (unix ms).
288 not_before_ms: u64,
289 #[serde(default)]
290 /// Optional dedup key index update.
291 dedup_key: Option<Vec<u8>>,
292 #[serde(default)]
293 /// Attempts already recorded for this job.
294 attempts: u32,
295 #[serde(default)]
296 /// Retry ceiling (`0` = unlimited).
297 max_attempts: u32,
298 },
299 /// Move a job from pending to leased.
300 Lease {
301 /// New lease token.
302 lease_id: u64,
303 /// Job being leased.
304 job_id: u64,
305 /// Worker node id.
306 worker_node: u64,
307 /// Worker instance id.
308 worker_instance: u32,
309 /// Lease expiry (unix ms; followers may use local timeout).
310 expires_at_ms: u64,
311 /// Monotonic lease id generator after this lease.
312 next_lease_id: u64,
313 },
314 /// Job completed — remove job and lease rows.
315 Ack {
316 /// Released lease.
317 lease_id: u64,
318 /// Completed job.
319 job_id: u64,
320 },
321 /// Worker rejected the job — return to pending or dead letter.
322 Nack {
323 /// Released lease.
324 lease_id: u64,
325 /// Requeued or dead-lettered job.
326 job_id: u64,
327 #[serde(default)]
328 /// Attempt count after this failure.
329 attempts: u32,
330 #[serde(default)]
331 /// When true the job is in the dead-letter set, not pending.
332 dead_letter: bool,
333 #[serde(default)]
334 /// Earliest re-lease time (unix ms) when requeued.
335 not_before_ms: u64,
336 },
337 /// Visibility timeout expired — job returns to pending or dead letter.
338 Reclaim {
339 /// Expired lease.
340 lease_id: u64,
341 /// Requeued or dead-lettered job.
342 job_id: u64,
343 #[serde(default)]
344 /// Attempt count after this failure.
345 attempts: u32,
346 #[serde(default)]
347 /// When true the job is in the dead-letter set, not pending.
348 dead_letter: bool,
349 #[serde(default)]
350 /// Earliest re-lease time (unix ms) when requeued.
351 not_before_ms: u64,
352 },
353 /// Operator moved a dead-letter job back to pending.
354 RequeueDeadLetter {
355 /// Job id to retry.
356 job_id: u64,
357 #[serde(default)]
358 /// Reset attempt counter (usually `0`).
359 attempts: u32,
360 },
361 /// Upsert a cron schedule (builder / operator).
362 UpsertSchedule {
363 /// Schedule body.
364 schedule: RecurringScheduleWire,
365 },
366 /// Leader advanced a schedule after enqueueing its tick.
367 UpdateScheduleNextRun {
368 /// Schedule name within the stream.
369 name: String,
370 /// Next fire time (unix ms).
371 next_run_ms: u64,
372 },
373}
374
375/// Batch of replication ops from the queue leader (`POST /raft/v1/queue/replicate`).
376#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
377pub struct QueueReplicateRequest {
378 /// Target stream.
379 pub stream: String,
380 /// Idempotent mutations to apply in order.
381 pub ops: Vec<QueueReplicateOp>,
382 /// Declared Raft leader id (must match the receiver's leader hint).
383 pub leader_id: u64,
384}
385
386/// Response to [`QueueReplicateRequest`].
387#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
388pub struct QueueReplicateReply {
389 /// Set when replication apply failed.
390 pub error: Option<String>,
391}
392
393/// Requeue a dead-letter job for another delivery attempt
394/// (`POST /raft/v1/queue/requeue-dead-letter`).
395#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
396pub struct QueueRequeueDeadLetterRequest {
397 /// Queue stream.
398 pub stream: String,
399 /// Job id in the dead-letter set.
400 pub job_id: u64,
401}
402
403/// Response to [`QueueRequeueDeadLetterRequest`].
404#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
405pub struct QueueRequeueDeadLetterReply {
406 /// Set when requeue failed.
407 pub error: Option<String>,
408}