Skip to main content

oxios_kernel/task/
model.rs

1// Task model — SQLite-backed task lifecycle management (RFC-043)
2// Ported from LobeHub's builtin-tool-task system.
3
4use chrono::Utc;
5use serde::{Deserialize, Serialize};
6use std::collections::HashMap;
7
8// ── Enums ──
9
10#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
11#[serde(rename_all = "snake_case")]
12pub enum TaskStatus {
13    #[default]
14    Backlog,
15    Scheduled,
16    Running,
17    Paused,
18    Completed,
19    Failed,
20    Canceled,
21}
22
23impl std::fmt::Display for TaskStatus {
24    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
25        match self {
26            Self::Backlog => write!(f, "backlog"),
27            Self::Scheduled => write!(f, "scheduled"),
28            Self::Running => write!(f, "running"),
29            Self::Paused => write!(f, "paused"),
30            Self::Completed => write!(f, "completed"),
31            Self::Failed => write!(f, "failed"),
32            Self::Canceled => write!(f, "canceled"),
33        }
34    }
35}
36
37impl std::str::FromStr for TaskStatus {
38    type Err = String;
39    fn from_str(s: &str) -> Result<Self, Self::Err> {
40        match s {
41            "backlog" => Ok(Self::Backlog),
42            "scheduled" => Ok(Self::Scheduled),
43            "running" => Ok(Self::Running),
44            "paused" => Ok(Self::Paused),
45            "completed" => Ok(Self::Completed),
46            "failed" => Ok(Self::Failed),
47            "canceled" => Ok(Self::Canceled),
48            other => Err(format!("unknown task status: {other}")),
49        }
50    }
51}
52
53#[derive(Debug, Clone, Serialize, Deserialize)]
54#[serde(rename_all = "snake_case")]
55pub enum TaskAutomationMode {
56    Schedule,
57    Heartbeat,
58}
59
60impl std::fmt::Display for TaskAutomationMode {
61    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
62        match self {
63            Self::Schedule => write!(f, "schedule"),
64            Self::Heartbeat => write!(f, "heartbeat"),
65        }
66    }
67}
68
69impl std::str::FromStr for TaskAutomationMode {
70    type Err = String;
71    fn from_str(s: &str) -> Result<Self, Self::Err> {
72        match s {
73            "schedule" => Ok(Self::Schedule),
74            "heartbeat" => Ok(Self::Heartbeat),
75            other => Err(format!("unknown automation mode: {other}")),
76        }
77    }
78}
79
80#[derive(Debug, Clone, Serialize, Deserialize)]
81#[serde(rename_all = "snake_case")]
82pub enum TaskRunTrigger {
83    Manual,
84    Schedule,
85    Heartbeat,
86}
87
88impl std::fmt::Display for TaskRunTrigger {
89    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
90        match self {
91            Self::Manual => write!(f, "manual"),
92            Self::Schedule => write!(f, "schedule"),
93            Self::Heartbeat => write!(f, "heartbeat"),
94        }
95    }
96}
97
98impl std::str::FromStr for TaskRunTrigger {
99    type Err = String;
100    fn from_str(s: &str) -> Result<Self, Self::Err> {
101        match s {
102            "manual" => Ok(Self::Manual),
103            "schedule" => Ok(Self::Schedule),
104            "heartbeat" => Ok(Self::Heartbeat),
105            other => Err(format!("unknown run trigger: {other}")),
106        }
107    }
108}
109
110// ── Core structs ──
111
112#[derive(Debug, Clone, Serialize, Deserialize)]
113pub struct Task {
114    pub id: String,
115    pub identifier: String,
116    pub name: String,
117    #[serde(skip_serializing_if = "Option::is_none")]
118    pub description: Option<String>,
119    pub instruction: String,
120    #[serde(default)]
121    pub status: TaskStatus,
122    #[serde(default)]
123    pub priority: u8,
124    #[serde(skip_serializing_if = "Option::is_none")]
125    pub sort_order: Option<f64>,
126    #[serde(skip_serializing_if = "Option::is_none")]
127    pub parent_task_id: Option<String>,
128    #[serde(skip_serializing_if = "Option::is_none")]
129    pub assignee_agent_id: Option<String>,
130    #[serde(skip_serializing_if = "Option::is_none")]
131    pub created_by_agent_id: Option<String>,
132    #[serde(skip_serializing_if = "Option::is_none")]
133    pub created_by_session_id: Option<String>,
134    #[serde(skip_serializing_if = "Option::is_none")]
135    pub automation_mode: Option<TaskAutomationMode>,
136    #[serde(skip_serializing_if = "Option::is_none")]
137    pub schedule_pattern: Option<String>,
138    #[serde(skip_serializing_if = "Option::is_none")]
139    pub schedule_timezone: Option<String>,
140    #[serde(skip_serializing_if = "Option::is_none")]
141    pub heartbeat_interval_secs: Option<u64>,
142    #[serde(skip_serializing_if = "Option::is_none")]
143    pub max_executions: Option<u32>,
144    #[serde(default)]
145    pub execution_count: u32,
146    pub verify_enabled: bool,
147    #[serde(skip_serializing_if = "Option::is_none")]
148    pub verify_requirement: Option<String>,
149    #[serde(default = "default_verify_iterations")]
150    pub verify_max_iterations: u32,
151    #[serde(skip_serializing_if = "Option::is_none")]
152    pub verify_verifier_agent_id: Option<String>,
153    pub created_at: String,
154    pub updated_at: String,
155    #[serde(skip_serializing_if = "Option::is_none")]
156    pub started_at: Option<String>,
157    #[serde(skip_serializing_if = "Option::is_none")]
158    pub completed_at: Option<String>,
159    #[serde(skip_serializing_if = "Option::is_none")]
160    pub last_run_at: Option<String>,
161    #[serde(skip_serializing_if = "Option::is_none")]
162    pub next_run_at: Option<String>,
163    #[serde(skip_serializing_if = "Option::is_none")]
164    pub last_error: Option<String>,
165    #[serde(default)]
166    pub consecutive_failures: u32,
167    /// Dependencies: list of task identifiers this task depends on.
168    #[serde(default)]
169    pub dependencies: Vec<String>,
170    /// Arbitrary context metadata (lifecycle audit, origin info, etc.)
171    #[serde(default)]
172    pub context: HashMap<String, serde_json::Value>,
173}
174
175fn default_verify_iterations() -> u32 {
176    3
177}
178
179#[derive(Debug, Clone, Serialize, Deserialize)]
180pub struct TaskComment {
181    pub id: String,
182    pub task_id: String,
183    pub content: String,
184    #[serde(skip_serializing_if = "Option::is_none")]
185    pub author_agent_id: Option<String>,
186    pub created_at: String,
187    #[serde(skip_serializing_if = "Option::is_none")]
188    pub updated_at: Option<String>,
189}
190
191#[derive(Debug, Clone, Serialize, Deserialize)]
192pub struct TaskRun {
193    pub id: String,
194    pub task_id: String,
195    #[serde(skip_serializing_if = "Option::is_none")]
196    pub session_id: Option<String>,
197    pub trigger: TaskRunTrigger,
198    #[serde(default = "default_run_status")]
199    pub status: String,
200    #[serde(skip_serializing_if = "Option::is_none")]
201    pub summary: Option<String>,
202    #[serde(skip_serializing_if = "Option::is_none")]
203    pub result_content: Option<String>,
204    pub started_at: String,
205    #[serde(skip_serializing_if = "Option::is_none")]
206    pub completed_at: Option<String>,
207    #[serde(skip_serializing_if = "Option::is_none")]
208    pub error: Option<String>,
209    #[serde(skip_serializing_if = "Option::is_none")]
210    pub cost_usd: Option<f64>,
211    #[serde(skip_serializing_if = "Option::is_none")]
212    pub tokens_used: Option<u64>,
213}
214
215fn default_run_status() -> String {
216    "running".to_string()
217}
218
219// ── Create params ──
220
221#[derive(Debug, Clone, Deserialize)]
222pub struct CreateTaskParams {
223    pub name: String,
224    pub instruction: String,
225    #[serde(default)]
226    pub identifier: Option<String>,
227    #[serde(default)]
228    pub description: Option<String>,
229    #[serde(default)]
230    pub priority: Option<u8>,
231    #[serde(default)]
232    pub parent_task_id: Option<String>,
233    #[serde(default)]
234    pub assignee_agent_id: Option<String>,
235    #[serde(default)]
236    pub sort_order: Option<f64>,
237}
238
239#[derive(Debug, Clone, Deserialize, Default)]
240pub struct ListTasksParams {
241    #[serde(default)]
242    pub statuses: Option<Vec<String>>,
243    #[serde(default)]
244    pub assignee_agent_id: Option<String>,
245    #[serde(default)]
246    pub parent_task_id: Option<String>,
247    #[serde(default)]
248    pub limit: Option<u32>,
249    #[serde(default)]
250    pub offset: Option<u32>,
251}
252
253#[derive(Debug, Clone, Deserialize)]
254pub struct SetScheduleParams {
255    pub automation_mode: Option<TaskAutomationMode>,
256    pub schedule_pattern: Option<String>,
257    pub schedule_timezone: Option<String>,
258    pub heartbeat_interval_secs: Option<u64>,
259    pub max_executions: Option<u32>,
260}
261
262#[derive(Debug, Clone, Deserialize)]
263pub struct SetVerifyParams {
264    pub enabled: Option<bool>,
265    pub requirement: Option<String>,
266    pub max_iterations: Option<u32>,
267    pub verifier_agent_id: Option<String>,
268}
269
270// ── Helpers ──
271
272impl Task {
273    /// Generate a slug-style identifier from a name.
274    pub fn slug_from_name(name: &str) -> String {
275        let slug: String = name
276            .to_lowercase()
277            .chars()
278            .map(|c| if c.is_alphanumeric() { c } else { '-' })
279            .collect();
280        let slug = slug.trim_matches('-').to_string();
281        let shortened = if slug.len() > 48 { &slug[..48] } else { &slug };
282        // Add short random suffix for uniqueness
283        let suffix = &uuid::Uuid::new_v4().to_string()[..8];
284        format!("{shortened}-{suffix}")
285    }
286
287    /// Check if this task should be auto-run now based on schedule/heartbeat.
288    pub fn should_run_now(&self) -> bool {
289        if self.automation_mode.is_none() {
290            return false;
291        }
292        if let Some(ref next) = self.next_run_at
293            && let Ok(next_dt) = chrono::DateTime::parse_from_rfc3339(next)
294        {
295            return next_dt.with_timezone(&Utc) <= Utc::now();
296        }
297        false
298    }
299    /// Check if the task has hit max_executions.
300    pub fn is_exhausted(&self) -> bool {
301        self.max_executions
302            .map(|max| self.execution_count >= max)
303            .unwrap_or(false)
304    }
305}