Skip to main content

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}