use std::collections::HashMap;
use chrono::{DateTime, Utc};
use croner::Cron;
use serde::{Deserialize, Serialize};
use crate::state::SharedState;
use crate::tmux;
const SCHEDULER_TICK_SECS: u64 = 15;
const REVIVAL_TIMEOUT_SECS: u64 = 30;
const TUI_READY_TIMEOUT_SECS: u64 = 30;
const REVIVAL_POLL_SECS: u64 = 2;
#[derive(
Clone, Debug, Default, Serialize, Deserialize, PartialEq, Eq, Hash, schemars::JsonSchema,
)]
#[serde(tag = "mode", rename_all = "snake_case")]
pub enum OnFire {
#[default]
ContinueSession,
NewSession,
PersistentWorktree {
#[serde(default)]
clear_context: bool,
},
DisposableWorktree,
}
impl OnFire {
pub fn clears_context(&self) -> bool {
match self {
Self::ContinueSession => false,
Self::NewSession => true,
Self::PersistentWorktree { clear_context } => *clear_context,
Self::DisposableWorktree => true,
}
}
pub fn uses_worktree(&self) -> bool {
matches!(
self,
Self::PersistentWorktree { .. } | Self::DisposableWorktree
)
}
pub fn is_disposable_worktree(&self) -> bool {
matches!(self, Self::DisposableWorktree)
}
pub fn kills_alive(&self) -> bool {
match self {
Self::ContinueSession | Self::NewSession => false,
Self::PersistentWorktree { clear_context } => *clear_context,
Self::DisposableWorktree => true,
}
}
}
#[derive(Clone, Debug, Serialize)]
pub struct ScheduledTask {
pub id: String,
pub name: String,
pub cron: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub target_session: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prompt: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reminder: Option<String>,
pub enabled: bool,
pub created_at: DateTime<Utc>,
pub next_run: Option<DateTime<Utc>>,
pub last_run: Option<DateTime<Utc>>,
pub last_status: Option<TaskRunStatus>,
pub run_count: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub project_dir: Option<String>,
#[serde(default)]
pub once: bool,
#[serde(
default,
skip_serializing_if = "Option::is_none",
alias = "claude_session_id"
)]
pub backend_session_id: Option<String>,
#[serde(default)]
pub on_fire: OnFire,
}
impl ScheduledTask {
pub fn session_name(&self) -> &str {
if matches!(self.on_fire, OnFire::ContinueSession) {
self.target_session.as_deref().unwrap_or(&self.name)
} else {
&self.name
}
}
}
impl<'de> serde::Deserialize<'de> for ScheduledTask {
fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
#[derive(Deserialize)]
struct Raw {
id: String,
name: String,
cron: String,
#[serde(default)]
target_session: Option<String>,
#[serde(default)]
prompt: Option<String>,
#[serde(default)]
reminder: Option<String>,
enabled: bool,
created_at: DateTime<Utc>,
#[serde(default)]
next_run: Option<DateTime<Utc>>,
#[serde(default)]
last_run: Option<DateTime<Utc>>,
#[serde(default)]
last_status: Option<TaskRunStatus>,
#[serde(default)]
run_count: u64,
#[serde(default)]
project_dir: Option<String>,
#[serde(default)]
once: bool,
#[serde(default, alias = "claude_session_id")]
backend_session_id: Option<String>,
#[serde(default)]
on_fire: Option<OnFire>,
#[serde(default)]
fresh: Option<bool>,
#[serde(default)]
worktree: Option<bool>,
#[serde(default)]
worktree_mode: Option<String>,
}
let raw = Raw::deserialize(deserializer)?;
let on_fire = raw.on_fire.unwrap_or_else(|| {
let fresh = raw.fresh.unwrap_or(false);
let worktree = raw.worktree.unwrap_or(false);
let worktree_mode = raw.worktree_mode.as_deref();
match (fresh, worktree, worktree_mode) {
(_, true, Some("per-fire")) => OnFire::DisposableWorktree,
(false, true, _) => OnFire::PersistentWorktree {
clear_context: false,
},
(true, true, _) => OnFire::PersistentWorktree {
clear_context: true,
},
(true, false, _) => OnFire::NewSession,
_ => OnFire::ContinueSession,
}
});
Ok(ScheduledTask {
id: raw.id,
name: raw.name,
cron: raw.cron,
target_session: raw.target_session,
prompt: raw.prompt,
reminder: raw.reminder,
enabled: raw.enabled,
created_at: raw.created_at,
next_run: raw.next_run,
last_run: raw.last_run,
last_status: raw.last_status,
run_count: raw.run_count,
project_dir: raw.project_dir,
once: raw.once,
backend_session_id: raw.backend_session_id,
on_fire,
})
}
}
#[derive(Clone, Debug, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum TaskRunStatus {
Ok,
Failed,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct TaskRun {
pub task_id: String,
pub task_name: String,
pub timestamp: DateTime<Utc>,
pub status: TaskRunStatus,
pub error: Option<String>,
pub session_name: String,
pub revived_pane: Option<String>,
}
impl TaskRun {
fn ok(task: &ScheduledTask, revived_pane: Option<String>) -> Self {
Self {
task_id: task.id.clone(),
task_name: task.name.clone(),
timestamp: Utc::now(),
status: TaskRunStatus::Ok,
error: None,
session_name: task.session_name().to_string(),
revived_pane,
}
}
fn failed(task: &ScheduledTask, error: String) -> Self {
Self {
task_id: task.id.clone(),
task_name: task.name.clone(),
timestamp: Utc::now(),
status: TaskRunStatus::Failed,
error: Some(error),
session_name: task.session_name().to_string(),
revived_pane: None,
}
}
}
pub fn validate_cron(expr: &str) -> Result<String, String> {
let cron = expr.parse::<Cron>().map_err(|e| format!("{e}"))?;
Ok(cron.pattern.to_string())
}
pub fn compute_next_run(expr: &str) -> Option<DateTime<Utc>> {
let cron = expr.parse::<Cron>().ok()?;
cron.find_next_occurrence(&Utc::now(), false).ok()
}
pub fn generate_task_id() -> String {
format!("{:08x}", rand::random::<u32>())
}
pub async fn run_scheduler(state: SharedState) {
recompute_all_next_runs(&state).await;
loop {
tokio::time::sleep(std::time::Duration::from_secs(SCHEDULER_TICK_SECS)).await;
tick(&state).await;
}
}
async fn recompute_all_next_runs(state: &SharedState) {
let mut tasks = state.scheduled_tasks.write().await;
let mut changed = false;
for task in tasks.values_mut() {
if task.enabled {
task.next_run = compute_next_run(&task.cron);
changed = true;
}
}
if changed {
state.persist_tasks_from(&tasks);
}
}
async fn tick(state: &SharedState) {
let now = Utc::now();
let due_ids: Vec<String> = {
let tasks = state.scheduled_tasks.read().await;
tasks
.values()
.filter(|t| t.enabled && t.next_run.is_some_and(|nr| nr <= now))
.map(|t| t.id.clone())
.collect()
};
for id in due_ids {
execute_task(state, &id).await;
}
}
pub async fn execute_task(state: &SharedState, task_id: &str) {
let task = {
let tasks = state.scheduled_tasks.read().await;
match tasks.get(task_id) {
Some(t) => t.clone(),
None => return,
}
};
let run = execute_injection(state, &task).await;
state
.update_task(task_id, |t| {
t.last_run = Some(run.timestamp);
t.last_status = Some(run.status.clone());
t.run_count += 1;
t.next_run = compute_next_run(&t.cron);
})
.await;
state.log_task_run(run).await;
if task.once {
state.remove_task(task_id).await;
}
}
async fn execute_injection(state: &SharedState, task: &ScheduledTask) -> TaskRun {
let session_name = task.session_name();
let session = {
let proto = state.protocol.read().await;
proto.sessions.get(session_name).cloned()
};
let Some(session) = session else {
if task.project_dir.is_some() || task.prompt.is_some() {
tracing::info!("session '{session_name}' not found, creating from task project_dir",);
return revive_from_task(state, task, None, None, None).await;
}
return TaskRun::failed(
task,
format!("session '{session_name}' not found and task has no project_dir"),
);
};
if !matches!(session.origin, crate::daemon_protocol::Origin::Local) {
return TaskRun::failed(task, "cannot target remote sessions".into());
}
let Some(pane) = &session.pane else {
return revive_from_task(
state,
task,
None,
session.metadata.model.clone(),
session.metadata.effort.clone(),
)
.await;
};
let pane_id = pane.clone();
let names: Vec<String> = state.backends.all_process_names();
let alive = tokio::task::spawn_blocking(move || {
let name_refs: Vec<&str> = names.iter().map(|s| s.as_str()).collect();
tmux::pane_alive(&pane_id, &name_refs)
})
.await
.unwrap_or(false);
let snapshot_model = session.metadata.model.clone();
let snapshot_effort = session.metadata.effort.clone();
if alive {
if task.on_fire.kills_alive() {
let dir = task
.project_dir
.as_deref()
.or(session.metadata.project_dir.as_deref())
.unwrap_or("/tmp");
return respawn_and_inject(state, task, pane, dir, snapshot_model, snapshot_effort)
.await;
}
if state
.protocol
.read()
.await
.sessions
.contains_key(session_name)
{
return TaskRun::ok(task, None);
}
tracing::info!("session '{session_name}' disappeared during alive check, reviving");
}
let project_dir = task
.project_dir
.as_deref()
.or(session.metadata.project_dir.as_deref());
revive_from_task(state, task, project_dir, snapshot_model, snapshot_effort).await
}
async fn respawn_and_inject(
state: &SharedState,
task: &ScheduledTask,
pane: &str,
dir: &str,
model: Option<String>,
effort: Option<String>,
) -> TaskRun {
let pane_id = pane.to_string();
let dir = dir.to_string();
let uses_worktree = task.on_fire.uses_worktree();
let is_disposable = task.on_fire.is_disposable_worktree();
let task_name = task.name.clone();
let backend = state.backend_for_session(task.session_name()).await;
let claude_cmd = backend.build_start_command(&crate::backend::StartOpts {
project_dir: dir.to_string(),
worktree: if uses_worktree {
if is_disposable {
Some(crate::backend::WorktreeMode::Disposable)
} else {
Some(crate::backend::WorktreeMode::Named(task_name.clone()))
}
} else {
None
},
model,
effort,
});
let full_cmd = if let Some(ref prompt) = task.prompt {
let full_text = match &task.reminder {
Some(r) => format!("{prompt}\n\n{r}"),
None => prompt.clone(),
};
let prompt_path = format!("/tmp/ouija-prompt-{}.txt", task_name.replace('/', "-"));
let _ = std::fs::write(&prompt_path, &full_text);
let escaped_pf = shell_escape(&prompt_path);
format!("{claude_cmd} \"$(cat {escaped_pf})\" ; rm -f {escaped_pf}")
} else {
claude_cmd
};
let session_name = task.session_name().to_string();
let respawn_result = tokio::task::spawn_blocking({
let pane_id = pane_id.clone();
move || -> anyhow::Result<()> {
let env_args = crate::tmux::pane_env_args(&session_name);
let mut args: Vec<&str> = vec!["respawn-pane", "-k"];
args.extend(env_args.iter().map(String::as_str));
args.extend_from_slice(&["-t", &pane_id, &full_cmd]);
let output = std::process::Command::new("tmux").args(&args).output()?;
if !output.status.success() {
anyhow::bail!(
"respawn-pane failed: {}",
String::from_utf8_lossy(&output.stderr)
);
}
Ok(())
}
})
.await;
match respawn_result {
Ok(Ok(())) => {}
Ok(Err(e)) => return TaskRun::failed(task, e.to_string()),
Err(e) => return TaskRun::failed(task, e.to_string()),
}
{
let mut proto = state.protocol.write().await;
if let Some(s) = proto.sessions.get_mut(task.session_name()) {
if s.metadata.prompt.is_none() {
s.metadata.prompt = task.prompt.clone();
}
if s.metadata.reminder.is_none() {
s.metadata.reminder = task.reminder.clone();
}
if s.metadata.on_fire.is_none() {
s.metadata.on_fire = Some(task.on_fire.clone());
}
s.metadata.backend_session_id = None;
}
}
let poll_pane = pane_id.clone();
let process_names: Vec<String> = backend
.process_names()
.iter()
.map(|s| s.to_string())
.collect();
let ready = tokio::task::spawn_blocking(move || {
let name_refs: Vec<&str> = process_names.iter().map(|s| s.as_str()).collect();
wait_for_process(&poll_pane, &name_refs, REVIVAL_TIMEOUT_SECS)
})
.await
.unwrap_or(false);
if !ready {
tracing::warn!("backend did not start in time after respawn in pane {pane_id}");
}
TaskRun::ok(task, None)
}
async fn revive_from_task(
state: &SharedState,
task: &ScheduledTask,
project_dir_override: Option<&str>,
model: Option<String>,
effort: Option<String>,
) -> TaskRun {
let project_dir = project_dir_override.or(task.project_dir.as_deref());
match revive_and_inject(state, task, project_dir, model, effort).await {
Ok(new_pane) => {
if task.on_fire.clears_context() {
let mut proto = state.protocol.write().await;
if let Some(s) = proto.sessions.get_mut(task.session_name()) {
s.metadata.backend_session_id = None;
}
}
TaskRun::ok(task, Some(new_pane))
}
Err(e) => TaskRun::failed(task, e.to_string()),
}
}
async fn revive_and_inject(
state: &SharedState,
task: &ScheduledTask,
project_dir: Option<&str>,
model: Option<String>,
effort: Option<String>,
) -> anyhow::Result<String> {
let dir = project_dir
.map(String::from)
.unwrap_or_else(|| std::env::var("HOME").unwrap_or_else(|_| "/tmp".into()));
let clears_context = task.on_fire.clears_context();
let uses_worktree = task.on_fire.uses_worktree();
let is_disposable = task.on_fire.is_disposable_worktree();
let worktree = if uses_worktree {
if is_disposable {
Some(crate::backend::WorktreeMode::Disposable)
} else {
Some(crate::backend::WorktreeMode::Named(task.name.clone()))
}
} else {
None
};
let backend = state.backend_for_session(task.session_name()).await;
let launch_cmd = if clears_context {
backend.build_start_command(&crate::backend::StartOpts {
project_dir: dir.clone(),
worktree,
model: model.clone(),
effort: effort.clone(),
})
} else {
let session_id = task
.backend_session_id
.clone()
.or_else(|| backend.detect_session_id(&dir));
backend
.build_resume_command(&crate::backend::ResumeOpts {
project_dir: dir.clone(),
session_id,
worktree,
model: model.clone(),
effort: effort.clone(),
})
.unwrap_or_else(|| {
backend.build_start_command(&crate::backend::StartOpts {
project_dir: dir.clone(),
worktree: None,
model: model.clone(),
effort: effort.clone(),
})
})
};
let is_tui = matches!(
backend.delivery_mode(),
crate::backend::DeliveryMode::TuiInjection
);
let full_launch_cmd = if clears_context || is_tui {
if let Some(ref prompt) = task.prompt {
let full_text = match &task.reminder {
Some(r) => format!("{prompt}\n\n{r}"),
None => prompt.clone(),
};
let prompt_path = format!("/tmp/ouija-prompt-{}.txt", task.name.replace('/', "-"));
let _ = std::fs::write(&prompt_path, &full_text);
let escaped_pf = shell_escape(&prompt_path);
format!("{launch_cmd} \"$(cat {escaped_pf})\" ; rm -f {escaped_pf}")
} else {
launch_cmd.clone()
}
} else {
launch_cmd.clone()
};
crate::backend::claude_code::pre_trust_workspace(&dir);
crate::backend::pre_trust_mise(&dir);
let new_pane = tokio::task::spawn_blocking({
let dir = dir.clone();
let window_name = task.session_name().to_string();
let tmux_session = crate::tmux::tmux_session_name(&dir);
move || -> anyhow::Result<String> {
let tmux_session_exists = std::process::Command::new("tmux")
.args(["has-session", "-t", &tmux_session])
.output()
.is_ok_and(|o| o.status.success());
let target = format!("{tmux_session}:");
let env_args = crate::tmux::pane_env_args(&window_name);
let output = if tmux_session_exists {
let mut args: Vec<&str> = vec!["new-window", "-d"];
args.extend(env_args.iter().map(String::as_str));
args.extend_from_slice(&[
"-t",
&target,
"-n",
&window_name,
"-P",
"-F",
"#{pane_id}",
]);
std::process::Command::new("tmux").args(&args).output()?
} else {
let mut args: Vec<&str> = vec!["new-session", "-d"];
args.extend(env_args.iter().map(String::as_str));
args.extend_from_slice(&[
"-s",
&tmux_session,
"-n",
&window_name,
"-P",
"-F",
"#{pane_id}",
]);
std::process::Command::new("tmux").args(&args).output()?
};
if !output.status.success() {
anyhow::bail!(
"tmux session/window creation failed: {}",
String::from_utf8_lossy(&output.stderr)
);
}
let pane_id = String::from_utf8_lossy(&output.stdout).trim().to_string();
let _ = std::process::Command::new("tmux")
.args([
"set-window-option",
"-t",
&pane_id,
"automatic-rename",
"off",
])
.status();
let hidden_cmd = format!(" {full_launch_cmd}");
std::process::Command::new("tmux")
.args(["send-keys", "-t", &pane_id, &hidden_cmd, "Enter"])
.status()?;
Ok(pane_id)
}
})
.await??;
let poll_pane = new_pane.clone();
let process_names: Vec<String> = backend
.process_names()
.iter()
.map(|s| s.to_string())
.collect();
let backend_name = backend.name().to_string();
let tui_pattern = backend.tui_ready_pattern().map(String::from);
let process_ready = tokio::task::spawn_blocking(move || {
let name_refs: Vec<&str> = process_names.iter().map(|s| s.as_str()).collect();
wait_for_process(&poll_pane, &name_refs, REVIVAL_TIMEOUT_SECS)
})
.await
.unwrap_or(false);
if !process_ready {
anyhow::bail!(
"{backend_name} did not start within {REVIVAL_TIMEOUT_SECS}s in pane {new_pane}"
);
}
if let Some(pattern) = tui_pattern {
let poll_pane = new_pane.clone();
let tui_ready = tokio::task::spawn_blocking(move || {
wait_for_tui_ready(&poll_pane, Some(&pattern), TUI_READY_TIMEOUT_SECS)
})
.await
.unwrap_or(false);
if !tui_ready {
tracing::warn!(
"{backend_name} TUI prompt not detected within {TUI_READY_TIMEOUT_SECS}s in pane {new_pane}, proceeding anyway"
);
}
}
let proto_meta = crate::daemon_protocol::SessionMeta {
project_dir: project_dir.map(String::from),
prompt: task.prompt.clone(),
reminder: task.reminder.clone(),
on_fire: Some(task.on_fire.clone()),
..Default::default()
};
state
.apply_and_execute(crate::daemon_protocol::Event::Register {
id: task.session_name().to_string(),
pane: Some(new_pane.clone()),
metadata: proto_meta,
})
.await;
if task.on_fire.is_disposable_worktree() {
if let Some(ref dir) = project_dir {
state
.perfire_worktree_panes
.write()
.await
.insert(new_pane.clone(), dir.to_string());
}
}
if !is_tui {
if let Some(ref prompt) = task.prompt {
let full_text = match &task.reminder {
Some(r) => format!("{prompt}\n\n{r}"),
None => prompt.clone(),
};
crate::nostr_transport::schedule_prompt_injection(
state,
task.session_name(),
new_pane.clone(),
full_text,
);
}
}
Ok(new_pane)
}
fn wait_for_process(pane: &str, names: &[&str], timeout_secs: u64) -> bool {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(timeout_secs);
while std::time::Instant::now() < deadline {
std::thread::sleep(std::time::Duration::from_secs(REVIVAL_POLL_SECS));
if let Ok(output) = std::process::Command::new("tmux")
.args([
"display-message",
"-t",
pane,
"-p",
"#{pane_current_command}",
])
.output()
{
let current = String::from_utf8_lossy(&output.stdout);
let current = current.trim();
if names.contains(¤t) {
return true;
}
}
}
false
}
fn wait_for_tui_ready(pane: &str, pattern: Option<&str>, timeout_secs: u64) -> bool {
let Some(pattern) = pattern else {
return true;
};
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(timeout_secs);
while std::time::Instant::now() < deadline {
std::thread::sleep(std::time::Duration::from_secs(REVIVAL_POLL_SECS));
if let Ok(output) = std::process::Command::new("tmux")
.args(["capture-pane", "-t", pane, "-p", "-S", "-20"])
.output()
{
if String::from_utf8_lossy(&output.stdout).contains(pattern) {
return true;
}
}
}
false
}
pub(crate) fn shell_escape(s: &str) -> String {
format!("'{}'", s.replace('\'', "'\\''"))
}
#[expect(
clippy::too_many_arguments,
reason = "flat parameters clearer than a builder for internal API"
)]
pub fn new_task(
name: String,
cron: String,
target_session: Option<String>,
prompt: Option<String>,
reminder: Option<String>,
once: bool,
backend_session_id: Option<String>,
on_fire: OnFire,
) -> ScheduledTask {
let next_run = compute_next_run(&cron);
ScheduledTask {
id: generate_task_id(),
name,
cron,
target_session,
prompt,
reminder,
enabled: true,
created_at: Utc::now(),
next_run,
last_run: None,
last_status: None,
run_count: 0,
project_dir: None,
once,
backend_session_id,
on_fire,
}
}
pub fn tasks_to_map(tasks: Vec<ScheduledTask>) -> HashMap<String, ScheduledTask> {
tasks.into_iter().map(|t| (t.id.clone(), t)).collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn validate_cron_valid() {
let result = validate_cron("*/5 * * * *");
assert!(result.is_ok(), "expected Ok, got {result:?}");
}
#[test]
fn validate_cron_invalid() {
let result = validate_cron("not a cron");
assert!(result.is_err());
}
#[test]
fn compute_next_run_returns_future() {
let next = compute_next_run("*/1 * * * *");
assert!(next.is_some());
assert!(next.unwrap() > Utc::now());
}
#[test]
fn compute_next_run_invalid_returns_none() {
assert!(compute_next_run("bad").is_none());
}
#[test]
fn task_id_is_8_hex_chars() {
let id = generate_task_id();
assert_eq!(id.len(), 8);
assert!(id.chars().all(|c| c.is_ascii_hexdigit()));
}
#[test]
fn task_serialization_round_trip() {
let task = ScheduledTask {
id: "a1b2c3d4".into(),
name: "test task".into(),
cron: "*/5 * * * *".into(),
target_session: Some("web".into()),
prompt: None,
reminder: None,
enabled: true,
created_at: Utc::now(),
next_run: Some(Utc::now()),
last_run: None,
last_status: None,
run_count: 0,
project_dir: Some("/tmp".into()),
once: false,
backend_session_id: None,
on_fire: OnFire::ContinueSession,
};
let json = serde_json::to_string(&task).unwrap();
let decoded: ScheduledTask = serde_json::from_str(&json).unwrap();
assert_eq!(decoded.id, task.id);
assert_eq!(decoded.name, task.name);
assert_eq!(decoded.project_dir, task.project_dir);
}
#[test]
fn shell_escape_basic() {
assert_eq!(shell_escape("/home/user"), "'/home/user'");
}
#[test]
fn shell_escape_with_quotes() {
assert_eq!(shell_escape("it's"), "'it'\\''s'");
}
#[test]
fn new_task_has_next_run() {
let task = new_task(
"t".into(),
"*/1 * * * *".into(),
Some("web".into()),
None,
None,
false,
None,
OnFire::ContinueSession,
);
assert!(task.next_run.is_some());
assert!(task.enabled);
assert_eq!(task.run_count, 0);
}
#[test]
fn task_worktree_serialization() {
let task = ScheduledTask {
id: "wt123456".into(),
name: "wt-task".into(),
cron: "0 9 * * *".into(),
target_session: Some("web".into()),
prompt: None,
reminder: None,
enabled: true,
created_at: Utc::now(),
next_run: None,
last_run: None,
last_status: None,
run_count: 0,
project_dir: Some("/tmp/project".into()),
once: false,
backend_session_id: None,
on_fire: OnFire::DisposableWorktree,
};
let json = serde_json::to_string(&task).unwrap();
assert!(json.contains("\"mode\":\"disposable_worktree\""));
let decoded: ScheduledTask = serde_json::from_str(&json).unwrap();
assert_eq!(decoded.on_fire, OnFire::DisposableWorktree);
}
#[test]
fn task_worktree_defaults_on_missing_fields() {
let json = r#"{"id":"x","name":"n","cron":"* * * * *","target_session":"s","enabled":true,"created_at":"2026-01-01T00:00:00Z","run_count":0,"once":false}"#;
let task: ScheduledTask = serde_json::from_str(json).unwrap();
assert_eq!(task.on_fire, OnFire::ContinueSession);
}
#[test]
fn on_fire_default_is_continue_session() {
assert_eq!(OnFire::default(), OnFire::ContinueSession);
}
#[test]
fn on_fire_serialization_round_trip() {
let variants = vec![
OnFire::ContinueSession,
OnFire::NewSession,
OnFire::PersistentWorktree {
clear_context: false,
},
OnFire::PersistentWorktree {
clear_context: true,
},
OnFire::DisposableWorktree,
];
for variant in variants {
let json = serde_json::to_string(&variant).unwrap();
let decoded: OnFire = serde_json::from_str(&json).unwrap();
assert_eq!(decoded, variant, "round-trip failed for {json}");
}
}
#[test]
fn on_fire_clear_context_defaults_false() {
let json = r#"{"mode":"persistent_worktree"}"#;
let on_fire: OnFire = serde_json::from_str(json).unwrap();
assert_eq!(
on_fire,
OnFire::PersistentWorktree {
clear_context: false
}
);
assert!(!on_fire.clears_context());
}
#[test]
fn legacy_task_json_migrates_to_on_fire() {
let json = r#"{"id":"x","name":"n","cron":"* * * * *","target_session":"s","enabled":true,"created_at":"2026-01-01T00:00:00Z","run_count":0,"fresh":true,"worktree":true,"worktree_mode":"per-fire"}"#;
let task: ScheduledTask = serde_json::from_str(json).unwrap();
assert_eq!(task.on_fire, OnFire::DisposableWorktree);
}
#[test]
fn legacy_task_fresh_only_migrates() {
let json = r#"{"id":"x","name":"n","cron":"* * * * *","target_session":"s","enabled":true,"created_at":"2026-01-01T00:00:00Z","run_count":0,"fresh":true}"#;
let task: ScheduledTask = serde_json::from_str(json).unwrap();
assert_eq!(task.on_fire, OnFire::NewSession);
}
#[test]
fn legacy_task_no_flags_migrates() {
let json = r#"{"id":"x","name":"n","cron":"* * * * *","target_session":"s","enabled":true,"created_at":"2026-01-01T00:00:00Z","run_count":0,"fresh":false}"#;
let task: ScheduledTask = serde_json::from_str(json).unwrap();
assert_eq!(task.on_fire, OnFire::ContinueSession);
}
#[test]
fn on_fire_kills_alive() {
assert!(!OnFire::ContinueSession.kills_alive());
assert!(!OnFire::NewSession.kills_alive());
assert!(
!OnFire::PersistentWorktree {
clear_context: false
}
.kills_alive()
);
assert!(
OnFire::PersistentWorktree {
clear_context: true
}
.kills_alive()
);
assert!(OnFire::DisposableWorktree.kills_alive());
}
#[test]
fn new_task_with_prompt_and_reminder() {
let task = new_task(
"test-task".into(),
"0 0 * * *".into(),
None,
Some("do the work".into()),
Some("call loop_next".into()),
false,
None,
OnFire::NewSession,
);
assert_eq!(task.prompt.as_deref(), Some("do the work"));
assert_eq!(task.reminder.as_deref(), Some("call loop_next"));
}
}