1use chrono::{DateTime, Utc};
2use serde::{Deserialize, Serialize};
3use serde_json::{Map, Value};
4
5use crate::{HandlerVersion, LeaseToken, QueueName, RetryPolicy, SignalName, StepName, TaskId, TaskName, WorkerId};
6
7#[derive(Clone, Debug, Eq, PartialEq)]
8pub struct QueueConfig {
9 pub name: QueueName,
10 pub terminal_retention: std::time::Duration,
11 pub idempotency_retention: std::time::Duration,
12 pub max_outstanding_tasks: Option<std::num::NonZeroU64>,
13 pub starvation_timeout: std::time::Duration,
14}
15
16impl QueueConfig {
17 pub fn new(name: QueueName) -> Self {
18 Self {
19 name,
20 terminal_retention: std::time::Duration::from_hours(7 * 24),
21 idempotency_retention: std::time::Duration::from_hours(30 * 24),
22 max_outstanding_tasks: None,
23 starvation_timeout: std::time::Duration::from_mins(5),
24 }
25 }
26}
27
28#[derive(Clone, Debug, Eq, PartialEq)]
29pub struct Queue {
30 pub name: QueueName,
31 pub terminal_retention: std::time::Duration,
32 pub idempotency_retention: std::time::Duration,
33 pub max_outstanding_tasks: Option<std::num::NonZeroU64>,
34 pub starvation_timeout: std::time::Duration,
35 pub paused_at: Option<DateTime<Utc>>,
36 pub created_at: DateTime<Utc>,
37 pub updated_at: DateTime<Utc>,
38}
39
40#[derive(Clone, Copy, Debug, Eq, PartialEq)]
41pub struct LeaseRenewal {
42 pub task_id: TaskId,
43 pub attempt: u16,
44 pub lease_token: LeaseToken,
45}
46
47#[derive(Clone, Debug, Eq, PartialEq)]
48pub struct WorkerRecord {
49 pub id: WorkerId,
50 pub queue_name: QueueName,
51 pub version: String,
52 pub draining: bool,
53 pub started_at: DateTime<Utc>,
54 pub heartbeat_at: DateTime<Utc>,
55 pub expires_at: DateTime<Utc>,
56 pub capabilities: Vec<(TaskName, HandlerVersion)>,
57}
58
59#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
60pub struct Checkpoint {
61 pub task_id: TaskId,
62 pub handler_version: HandlerVersion,
63 pub step_name: StepName,
64 pub occurrence: u32,
65 pub value: Value,
66 pub created_at: DateTime<Utc>,
67}
68
69#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
70pub struct Signal {
71 pub task_id: TaskId,
72 pub name: SignalName,
73 pub occurrence: u32,
74 pub value: Value,
75 pub created_at: DateTime<Utc>,
76}
77
78#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
79#[serde(rename_all = "snake_case")]
80pub enum TaskState {
81 Pending,
82 Running,
83 Waiting,
84 Succeeded,
85 Failed,
86 Cancelled,
87}
88
89impl TaskState {
90 pub const fn as_str(self) -> &'static str {
91 match self {
92 Self::Pending => "pending",
93 Self::Running => "running",
94 Self::Waiting => "waiting",
95 Self::Succeeded => "succeeded",
96 Self::Failed => "failed",
97 Self::Cancelled => "cancelled",
98 }
99 }
100
101 pub const fn is_terminal(self) -> bool {
102 matches!(self, Self::Succeeded | Self::Failed | Self::Cancelled)
103 }
104}
105
106#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
107pub struct TaskResult {
108 pub state: TaskState,
109 pub result: Option<Value>,
110 pub error: Option<Value>,
111 pub completed_at: Option<DateTime<Utc>>,
112}
113
114#[derive(Clone, Debug, Deserialize, Serialize)]
115pub struct EnqueueRequest {
116 pub task_name: TaskName,
117 pub handler_version: HandlerVersion,
118 pub payload: Value,
119 pub queue_name: QueueName,
120 pub run_at: Option<DateTime<Utc>>,
121 pub priority: i16,
122 pub max_attempts: u16,
123 pub idempotency_key: Option<String>,
124 pub headers: Map<String, Value>,
125}
126
127impl EnqueueRequest {
128 pub fn new(task_name: TaskName, payload: Value) -> Self {
129 Self {
130 task_name,
131 handler_version: HandlerVersion::default(),
132 payload,
133 queue_name: QueueName::default(),
134 run_at: None,
135 priority: 0,
136 max_attempts: 5,
137 idempotency_key: None,
138 headers: Map::new(),
139 }
140 }
141}
142
143#[derive(Clone, Copy, Debug, Eq, PartialEq)]
144pub struct EnqueueResult {
145 pub task_id: TaskId,
146 pub created: bool,
147}
148
149#[derive(Clone, Debug, Deserialize, Serialize)]
150pub struct Task {
151 pub id: TaskId,
152 pub parent_task_id: Option<TaskId>,
153 pub queue_name: QueueName,
154 pub task_name: TaskName,
155 pub handler_version: HandlerVersion,
156 pub payload: Value,
157 pub headers: Map<String, Value>,
158 pub state: TaskState,
159 pub priority: i16,
160 pub run_at: DateTime<Utc>,
161 pub attempt: u16,
162 pub max_attempts: u16,
163 pub retry_policy: Option<RetryPolicy>,
164 pub lease_token: Option<LeaseToken>,
165 pub lease_owner: Option<WorkerId>,
166 pub lease_expires_at: Option<DateTime<Utc>>,
167 pub created_at: DateTime<Utc>,
168 pub updated_at: DateTime<Utc>,
169 pub completed_at: Option<DateTime<Utc>>,
170 pub result: Option<Value>,
171 pub error: Option<Value>,
172}