feat: 任务页新增定时计划——指定执行智能体(默认通用助手)与间隔,后台到点自动在智能体会话中执行并记录结果
This commit is contained in:
@@ -3,6 +3,7 @@ pub mod app;
|
||||
pub mod error;
|
||||
pub mod mcp_servers;
|
||||
pub mod models;
|
||||
pub mod scheduled_tasks;
|
||||
pub mod sessions;
|
||||
pub mod settings;
|
||||
pub mod skills;
|
||||
@@ -12,5 +13,6 @@ pub use app::CoreApp;
|
||||
pub use error::{CoreError, Result};
|
||||
pub use models::ModelInfo;
|
||||
pub use agents::Agent;
|
||||
pub use scheduled_tasks::ScheduledTask;
|
||||
pub use sessions::{Conversation, Message};
|
||||
pub use workflows::Workflow;
|
||||
@@ -0,0 +1,290 @@
|
||||
use crate::error::Result;
|
||||
use rusqlite::{params, Connection};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
#[derive(Debug, Clone, Serialize, Deserialize)]
|
||||
pub struct ScheduledTask {
|
||||
pub id: String,
|
||||
pub name: String,
|
||||
pub agent_id: Option<String>,
|
||||
pub prompt: String,
|
||||
pub interval_minutes: i64,
|
||||
pub enabled: bool,
|
||||
pub next_run_at: Option<String>,
|
||||
pub last_run_at: Option<String>,
|
||||
/// idle(从未执行)/ running / success / error
|
||||
pub last_status: String,
|
||||
pub last_result: Option<String>,
|
||||
pub last_error: Option<String>,
|
||||
pub created_at: String,
|
||||
pub updated_at: String,
|
||||
}
|
||||
|
||||
const SELECT_COLUMNS: &str = "id, name, agent_id, prompt, interval_minutes, enabled, next_run_at, last_run_at, last_status, last_result, last_error, created_at, updated_at";
|
||||
|
||||
/// 把分钟数转成 SQLite datetime 修饰符字符串(如 `+30 minutes`)。
|
||||
fn interval_modifier(minutes: i64) -> String {
|
||||
format!("+{} minutes", minutes.max(1))
|
||||
}
|
||||
|
||||
pub fn list(db: &Connection) -> Result<Vec<ScheduledTask>> {
|
||||
let mut stmt = db.prepare(&format!(
|
||||
"SELECT {SELECT_COLUMNS} FROM scheduled_tasks ORDER BY updated_at DESC, created_at DESC"
|
||||
))?;
|
||||
let rows = stmt.query_map([], row_to_task)?;
|
||||
let mut out = Vec::new();
|
||||
for row in rows {
|
||||
out.push(row?);
|
||||
}
|
||||
Ok(out)
|
||||
}
|
||||
|
||||
pub fn get(db: &Connection, id: &str) -> Result<Option<ScheduledTask>> {
|
||||
let mut stmt = db.prepare(&format!(
|
||||
"SELECT {SELECT_COLUMNS} FROM scheduled_tasks WHERE id = ?1"
|
||||
))?;
|
||||
let mut rows = stmt.query_map(params![id], row_to_task)?;
|
||||
match rows.next() {
|
||||
Some(row) => Ok(Some(row?)),
|
||||
None => Ok(None),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn insert(db: &Connection, task: &ScheduledTask) -> Result<()> {
|
||||
db.execute(
|
||||
"INSERT INTO scheduled_tasks
|
||||
(id, name, agent_id, prompt, interval_minutes, enabled, next_run_at, last_status)
|
||||
VALUES (?1, ?2, ?3, ?4, ?5, ?6, datetime('now', ?7), 'idle')",
|
||||
params![
|
||||
task.id,
|
||||
task.name,
|
||||
task.agent_id,
|
||||
task.prompt,
|
||||
task.interval_minutes.max(1),
|
||||
task.enabled as i32,
|
||||
interval_modifier(task.interval_minutes),
|
||||
],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn update(db: &Connection, task: &ScheduledTask) -> Result<()> {
|
||||
db.execute(
|
||||
"UPDATE scheduled_tasks SET name = ?1, agent_id = ?2, prompt = ?3,
|
||||
interval_minutes = ?4, enabled = ?5,
|
||||
next_run_at = CASE WHEN ?6 = 1 THEN datetime('now', ?7) ELSE NULL END,
|
||||
updated_at = datetime('now')
|
||||
WHERE id = ?8",
|
||||
params![
|
||||
task.name,
|
||||
task.agent_id,
|
||||
task.prompt,
|
||||
task.interval_minutes.max(1),
|
||||
task.enabled as i32,
|
||||
task.enabled as i32,
|
||||
interval_modifier(task.interval_minutes),
|
||||
task.id,
|
||||
],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn delete(db: &Connection, id: &str) -> Result<()> {
|
||||
db.execute("DELETE FROM scheduled_tasks WHERE id = ?1", params![id])?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn set_enabled(db: &Connection, id: &str, enabled: bool) -> Result<()> {
|
||||
db.execute(
|
||||
"UPDATE scheduled_tasks SET enabled = ?1,
|
||||
next_run_at = CASE WHEN ?1 = 1 THEN datetime('now', ?2) ELSE NULL END,
|
||||
updated_at = datetime('now')
|
||||
WHERE id = ?3",
|
||||
params![
|
||||
enabled as i32,
|
||||
interval_modifier(
|
||||
db.query_row(
|
||||
"SELECT interval_minutes FROM scheduled_tasks WHERE id = ?1",
|
||||
params![id],
|
||||
|row| row.get::<_, i64>(0),
|
||||
)
|
||||
.unwrap_or(60),
|
||||
),
|
||||
id,
|
||||
],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// 标记任务开始执行(后台调度与「立即执行」共用)。
|
||||
pub fn mark_running(db: &Connection, id: &str) -> Result<()> {
|
||||
db.execute(
|
||||
"UPDATE scheduled_tasks SET last_status = 'running', updated_at = datetime('now')
|
||||
WHERE id = ?1",
|
||||
params![id],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// 写入一次执行结果,并推进下一次执行时间。
|
||||
pub fn finish_run(
|
||||
db: &Connection,
|
||||
id: &str,
|
||||
success: bool,
|
||||
result: Option<&str>,
|
||||
error: Option<&str>,
|
||||
) -> Result<()> {
|
||||
let interval: i64 = db.query_row(
|
||||
"SELECT interval_minutes FROM scheduled_tasks WHERE id = ?1",
|
||||
params![id],
|
||||
|row| row.get(0),
|
||||
)?;
|
||||
db.execute(
|
||||
"UPDATE scheduled_tasks SET last_status = ?1, last_run_at = datetime('now'),
|
||||
last_result = ?2, last_error = ?3, next_run_at = datetime('now', ?4),
|
||||
updated_at = datetime('now')
|
||||
WHERE id = ?5",
|
||||
params![
|
||||
if success { "success" } else { "error" },
|
||||
result,
|
||||
error,
|
||||
interval_modifier(interval),
|
||||
id,
|
||||
],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// 列出到点可执行且未在运行中的已启用任务。
|
||||
pub fn list_due(db: &Connection) -> Result<Vec<ScheduledTask>> {
|
||||
let mut stmt = db.prepare(&format!(
|
||||
"SELECT {SELECT_COLUMNS} FROM scheduled_tasks
|
||||
WHERE enabled = 1 AND last_status != 'running'
|
||||
AND (next_run_at IS NULL OR next_run_at <= datetime('now'))
|
||||
ORDER BY next_run_at ASC"
|
||||
))?;
|
||||
let rows = stmt.query_map([], row_to_task)?;
|
||||
let mut out = Vec::new();
|
||||
for row in rows {
|
||||
out.push(row?);
|
||||
}
|
||||
Ok(out)
|
||||
}
|
||||
|
||||
/// 应用重启后把上次中断的「运行中」任务标记为失败,避免永远卡住。
|
||||
pub fn reset_stale_running(db: &Connection) -> Result<usize> {
|
||||
let n = db.execute(
|
||||
"UPDATE scheduled_tasks SET last_status = 'error',
|
||||
last_error = '应用重启,上次运行被中断', updated_at = datetime('now')
|
||||
WHERE last_status = 'running'",
|
||||
[],
|
||||
)?;
|
||||
Ok(n)
|
||||
}
|
||||
|
||||
fn row_to_task(row: &rusqlite::Row<'_>) -> rusqlite::Result<ScheduledTask> {
|
||||
let enabled: i32 = row.get(5)?;
|
||||
Ok(ScheduledTask {
|
||||
id: row.get(0)?,
|
||||
name: row.get(1)?,
|
||||
agent_id: row.get(2)?,
|
||||
prompt: row.get(3)?,
|
||||
interval_minutes: row.get(4)?,
|
||||
enabled: enabled != 0,
|
||||
next_run_at: row.get(6)?,
|
||||
last_run_at: row.get(7)?,
|
||||
last_status: row.get(8)?,
|
||||
last_result: row.get(9)?,
|
||||
last_error: row.get(10)?,
|
||||
created_at: row.get(11)?,
|
||||
updated_at: row.get(12)?,
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use rusqlite::Connection;
|
||||
|
||||
fn test_conn() -> Connection {
|
||||
let conn = Connection::open_in_memory().unwrap();
|
||||
conn.execute_batch(include_str!("schema.sql")).unwrap();
|
||||
conn
|
||||
}
|
||||
|
||||
fn sample(id: &str) -> ScheduledTask {
|
||||
ScheduledTask {
|
||||
id: id.to_string(),
|
||||
name: "测试计划".to_string(),
|
||||
agent_id: None,
|
||||
prompt: "请总结今天的进展".to_string(),
|
||||
interval_minutes: 30,
|
||||
enabled: true,
|
||||
next_run_at: None,
|
||||
last_run_at: None,
|
||||
last_status: "idle".to_string(),
|
||||
last_result: None,
|
||||
last_error: None,
|
||||
created_at: String::new(),
|
||||
updated_at: String::new(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn insert_sets_future_next_run_and_not_due() {
|
||||
let conn = test_conn();
|
||||
insert(&conn, &sample("t1")).unwrap();
|
||||
let t = get(&conn, "t1").unwrap().unwrap();
|
||||
assert!(t.next_run_at.is_some());
|
||||
// 下次执行在未来,不应出现在到点列表
|
||||
assert!(list_due(&conn).unwrap().is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn due_then_running_then_finish_advances_schedule() {
|
||||
let conn = test_conn();
|
||||
insert(&conn, &sample("t1")).unwrap();
|
||||
conn.execute(
|
||||
"UPDATE scheduled_tasks SET next_run_at = datetime('now', '-1 minutes')",
|
||||
[],
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(list_due(&conn).unwrap().len(), 1);
|
||||
|
||||
mark_running(&conn, "t1").unwrap();
|
||||
// 运行中的计划不会被重复调度
|
||||
assert!(list_due(&conn).unwrap().is_empty());
|
||||
|
||||
finish_run(&conn, "t1", true, Some("已完成"), None).unwrap();
|
||||
let t = get(&conn, "t1").unwrap().unwrap();
|
||||
assert_eq!(t.last_status, "success");
|
||||
assert_eq!(t.last_result.as_deref(), Some("已完成"));
|
||||
assert!(t.last_run_at.is_some());
|
||||
assert!(t.next_run_at.is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn disable_clears_schedule_and_enable_recomputes() {
|
||||
let conn = test_conn();
|
||||
insert(&conn, &sample("t1")).unwrap();
|
||||
set_enabled(&conn, "t1", false).unwrap();
|
||||
let t = get(&conn, "t1").unwrap().unwrap();
|
||||
assert!(!t.enabled);
|
||||
assert!(t.next_run_at.is_none());
|
||||
|
||||
set_enabled(&conn, "t1", true).unwrap();
|
||||
let t = get(&conn, "t1").unwrap().unwrap();
|
||||
assert!(t.enabled);
|
||||
assert!(t.next_run_at.is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reset_stale_running_marks_interrupted_as_error() {
|
||||
let conn = test_conn();
|
||||
insert(&conn, &sample("t1")).unwrap();
|
||||
mark_running(&conn, "t1").unwrap();
|
||||
assert_eq!(reset_stale_running(&conn).unwrap(), 1);
|
||||
let t = get(&conn, "t1").unwrap().unwrap();
|
||||
assert_eq!(t.last_status, "error");
|
||||
}
|
||||
}
|
||||
@@ -106,3 +106,19 @@ CREATE TABLE IF NOT EXISTS mcp_servers (
|
||||
enabled INTEGER NOT NULL DEFAULT 1,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now'))
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS scheduled_tasks (
|
||||
id TEXT PRIMARY KEY,
|
||||
name TEXT NOT NULL,
|
||||
agent_id TEXT,
|
||||
prompt TEXT NOT NULL DEFAULT '',
|
||||
interval_minutes INTEGER NOT NULL DEFAULT 60,
|
||||
enabled INTEGER NOT NULL DEFAULT 1,
|
||||
next_run_at TEXT,
|
||||
last_run_at TEXT,
|
||||
last_status TEXT NOT NULL DEFAULT 'idle',
|
||||
last_result TEXT,
|
||||
last_error TEXT,
|
||||
created_at TEXT NOT NULL DEFAULT (datetime('now')),
|
||||
updated_at TEXT NOT NULL DEFAULT (datetime('now'))
|
||||
);
|
||||
Reference in New Issue
Block a user