Skip to main content

pgtask_core/
task.rs

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}