返回 CodeWhale
automation_manager.rs
根目录 / crates / tui / src / automation_manager.rs
1 //! Durable automation records and scheduler support.
2 //!
3 //! Automations are local-first recurring jobs that enqueue standard background
4 //! tasks. This module stores automation definitions and run history under
5 //! `~/.codewhale/automations` (or `DEEPSEEK_AUTOMATIONS_DIR` override).
6
7 use std::collections::BTreeMap;
8 use std::fs;
9 use std::future::Future;
10 use std::path::{Path, PathBuf};
11 use std::sync::Arc;
12
13 use anyhow::{Context, Result, bail};
14 use chrono::{
15 DateTime, Datelike, Duration, Local, NaiveDateTime, TimeZone, Timelike, Utc, Weekday,
16 };
17 use serde::{Deserialize, Serialize};
18 use tokio::sync::Mutex;
19 use tokio::time::sleep;
20 use tokio_util::sync::CancellationToken;
21 use uuid::Uuid;
22
23 use crate::task_manager::{NewTaskRequest, SharedTaskManager, TaskStatus};
24 use crate::utils::spawn_supervised;
25
26 /// Current automation record schema. `pub(crate)` so the Operate keepalive
27 /// can build a fixed-id record directly (no create/delete id swap).
28 // v2 pins provider identity. Older runtimes must reject a pinned definition
29 // instead of silently sending its model through their current provider.
30 pub(crate) const CURRENT_AUTOMATION_SCHEMA_VERSION: u32 = 3;
31 const CURRENT_RUN_SCHEMA_VERSION: u32 = 3;
32 const CURRENT_TRIGGER_SCHEMA_VERSION: u32 = 3;
33 const DEFAULT_AUTOMATION_MODE: &str = "agent";
34 const DEFAULT_AUTOMATION_ALLOW_SHELL: bool = false;
35 const DEFAULT_AUTOMATION_TRUST_MODE: bool = false;
36 const DEFAULT_AUTOMATION_AUTO_APPROVE: bool = false;
37 const DEFAULT_AUTOMATION_DELIVERY_MODE: AutomationDeliveryMode = AutomationDeliveryMode::Task;
38 pub const AUTOMATION_WATCHER_NO_REPORT_SENTINEL: &str = "NOTHING_TO_REPORT";
39 const MAX_HOURLY_SEARCH_STEPS: usize = 24 * 21;
40 const MAX_CRON_SEARCH_MINUTES: usize = 60 * 24 * 366 * 5;
41 const fn default_automation_schema_version() -> u32 {
42 CURRENT_AUTOMATION_SCHEMA_VERSION
43 }
44
45 const fn default_run_schema_version() -> u32 {
46 CURRENT_RUN_SCHEMA_VERSION
47 }
48
49 const fn default_trigger_schema_version() -> u32 {
50 CURRENT_TRIGGER_SCHEMA_VERSION
51 }
52
53 // ── Delayed-trigger types ──────────────────────────────────────────────────
54
55 /// Status of a one-shot delayed trigger.
56 #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
57 #[serde(rename_all = "snake_case")]
58 pub enum DelayedTriggerStatus {
59 /// Waiting to fire.
60 Pending,
61 /// Durable admission owns this trigger; task acceptance is being recovered.
62 Dispatching,
63 /// The trigger was fired and a task was enqueued.
64 Fired,
65 /// The trigger was explicitly canceled before it fired.
66 Canceled,
67 /// The trigger fired but failed to enqueue a task.
68 Failed,
69 }
70
71 /// A durable one-shot delayed continuation record.
72 ///
73 /// Stored under `~/.codewhale/automations/triggers/{trigger_id}.json`.
74 #[derive(Debug, Clone, Serialize, Deserialize)]
75 pub struct DelayedTriggerRecord {
76 #[serde(default = "default_trigger_schema_version")]
77 pub schema_version: u32,
78 pub trigger_id: String,
79 /// Absolute UTC time at which the trigger should fire.
80 pub fire_at: DateTime<Utc>,
81 /// The message that will be submitted as a new task when the trigger fires.
82 pub message: String,
83 /// Working directory for the task that fires when the trigger trips.
84 #[serde(skip_serializing_if = "Option::is_none")]
85 pub workspace: Option<PathBuf>,
86 /// Session that scheduled this trigger. Missing legacy ownership fails closed.
87 #[serde(default, skip_serializing_if = "Option::is_none")]
88 pub owner_session_id: Option<String>,
89 pub status: DelayedTriggerStatus,
90 pub created_at: DateTime<Utc>,
91 #[serde(skip_serializing_if = "Option::is_none")]
92 pub fired_at: Option<DateTime<Utc>>,
93 #[serde(skip_serializing_if = "Option::is_none")]
94 pub task_id: Option<String>,
95 #[serde(skip_serializing_if = "Option::is_none")]
96 pub thread_id: Option<String>,
97 #[serde(skip_serializing_if = "Option::is_none")]
98 pub error: Option<String>,
99 /// Optional lineage: the trigger id that scheduled this one (for re-arm chains).
100 #[serde(skip_serializing_if = "Option::is_none")]
101 pub parent_trigger_id: Option<String>,
102 #[serde(default, skip_serializing_if = "Option::is_none")]
103 pub dispatch: Option<AutomationDispatch>,
104 /// Bound by the trusted service, independently of visibility ownership.
105 #[serde(default, skip_serializing_if = "Option::is_none")]
106 pub execution_scope: Option<String>,
107 }
108
109 /// Input for creating a new delayed trigger.
110 #[derive(Debug, Clone)]
111 pub struct CreateDelayedTriggerRequest {
112 /// Absolute fire time. Callers must resolve `delay_minutes` → `fire_at`
113 /// before calling this function.
114 pub fire_at: DateTime<Utc>,
115 /// Message to submit as a new task when the trigger fires.
116 pub message: String,
117 /// Optional workspace directory for the fired task.
118 pub workspace: Option<PathBuf>,
119 /// Session that owns controls and the task created when this trigger fires.
120 pub owner_session_id: Option<String>,
121 /// Optional parent trigger id for re-arm lineage tracking.
122 pub parent_trigger_id: Option<String>,
123 }
124
125 // ── End delayed-trigger types ──────────────────────────────────────────────
126
127 #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
128 #[serde(rename_all = "snake_case")]
129 pub enum AutomationStatus {
130 Active,
131 Paused,
132 }
133
134 #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
135 #[serde(rename_all = "snake_case")]
136 pub enum AutomationRunStatus {
137 Queued,
138 Running,
139 Completed,
140 Failed,
141 Canceled,
142 }
143
144 #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
145 #[serde(rename_all = "snake_case")]
146 pub enum AutomationDeliveryMode {
147 #[default]
148 Task,
149 Watcher,
150 }
151
152 #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
153 pub struct AutomationRecord {
154 #[serde(default = "default_automation_schema_version")]
155 pub schema_version: u32,
156 pub id: String,
157 pub name: String,
158 pub prompt: String,
159 pub rrule: String,
160 #[serde(default)]
161 pub cwds: Vec<PathBuf>,
162 #[serde(default, skip_serializing_if = "Option::is_none")]
163 pub model: Option<String>,
164 /// Exact provider provenance for a pinned model; absent on legacy records.
165 #[serde(default, skip_serializing_if = "Option::is_none")]
166 pub model_provider: Option<String>,
167 #[serde(default, skip_serializing_if = "Option::is_none")]
168 pub model_provider_id: Option<String>,
169 #[serde(default, skip_serializing_if = "Option::is_none")]
170 pub mode: Option<String>,
171 #[serde(default, skip_serializing_if = "Option::is_none")]
172 pub allow_shell: Option<bool>,
173 #[serde(default, skip_serializing_if = "Option::is_none")]
174 pub trust_mode: Option<bool>,
175 #[serde(default, skip_serializing_if = "Option::is_none")]
176 pub auto_approve: Option<bool>,
177 #[serde(default, skip_serializing_if = "Option::is_none")]
178 pub delivery_mode: Option<AutomationDeliveryMode>,
179 pub status: AutomationStatus,
180 pub created_at: DateTime<Utc>,
181 pub updated_at: DateTime<Utc>,
182 #[serde(skip_serializing_if = "Option::is_none")]
183 pub next_run_at: Option<DateTime<Utc>>,
184 #[serde(skip_serializing_if = "Option::is_none")]
185 pub last_run_at: Option<DateTime<Utc>>,
186 /// Bound by the trusted service, independently of visibility ownership.
187 #[serde(default, skip_serializing_if = "Option::is_none")]
188 pub execution_scope: Option<String>,
189 }
190
191 impl AutomationRecord {
192 fn task_mode(&self) -> String {
193 self.mode
194 .as_deref()
195 .map(str::trim)
196 .filter(|mode| !mode.is_empty())
197 .unwrap_or(DEFAULT_AUTOMATION_MODE)
198 .to_string()
199 }
200
201 fn task_allow_shell(&self) -> bool {
202 self.allow_shell.unwrap_or(DEFAULT_AUTOMATION_ALLOW_SHELL)
203 }
204
205 fn task_trust_mode(&self) -> bool {
206 self.trust_mode.unwrap_or(DEFAULT_AUTOMATION_TRUST_MODE)
207 }
208
209 fn task_auto_approve(&self) -> bool {
210 self.auto_approve.unwrap_or(DEFAULT_AUTOMATION_AUTO_APPROVE)
211 }
212
213 fn delivery_mode(&self) -> AutomationDeliveryMode {
214 self.delivery_mode
215 .unwrap_or(DEFAULT_AUTOMATION_DELIVERY_MODE)
216 }
217 }
218
219 #[derive(Debug, Clone, Serialize, Deserialize)]
220 pub struct AutomationRunRecord {
221 #[serde(default = "default_run_schema_version")]
222 pub schema_version: u32,
223 pub id: String,
224 pub automation_id: String,
225 pub scheduled_for: DateTime<Utc>,
226 pub status: AutomationRunStatus,
227 pub created_at: DateTime<Utc>,
228 #[serde(skip_serializing_if = "Option::is_none")]
229 pub started_at: Option<DateTime<Utc>>,
230 #[serde(skip_serializing_if = "Option::is_none")]
231 pub ended_at: Option<DateTime<Utc>>,
232 #[serde(skip_serializing_if = "Option::is_none")]
233 pub task_id: Option<String>,
234 #[serde(skip_serializing_if = "Option::is_none")]
235 pub thread_id: Option<String>,
236 #[serde(skip_serializing_if = "Option::is_none")]
237 pub turn_id: Option<String>,
238 #[serde(skip_serializing_if = "Option::is_none")]
239 pub error: Option<String>,
240 #[serde(default, skip_serializing_if = "Option::is_none")]
241 pub dispatch: Option<AutomationDispatch>,
242 }
243
244 /// Immutable request and store bound to a durable occurrence before enqueue.
245 /// `accepted` records task promotion, not provider execution or completion.
246 #[derive(Debug, Clone, Serialize, Deserialize)]
247 pub struct AutomationDispatch {
248 #[serde(default, skip_serializing_if = "Option::is_none")]
249 execution_scope: Option<String>,
250 request: NewTaskRequest,
251 task_data_dir: PathBuf,
252 #[serde(default)]
253 accepted: bool,
254 #[serde(default)]
255 delivery_mode: AutomationDeliveryMode,
256 #[serde(default)]
257 suppress_report: bool,
258 #[serde(default, skip_serializing_if = "Option::is_none")]
259 schedule: Option<AdmittedSchedule>,
260 }
261
262 #[derive(Debug, Clone, Serialize, Deserialize)]
263 struct AdmittedSchedule {
264 updated_at: DateTime<Utc>,
265 rrule: String,
266 }
267
268 #[derive(Debug, Clone, Serialize, Deserialize)]
269 pub struct CreateAutomationRequest {
270 pub name: String,
271 pub prompt: String,
272 pub rrule: String,
273 #[serde(default)]
274 pub cwds: Vec<PathBuf>,
275 #[serde(default)]
276 pub model: Option<String>,
277 #[serde(default)]
278 pub model_provider: Option<String>,
279 #[serde(default)]
280 pub model_provider_id: Option<String>,
281 #[serde(default)]
282 pub mode: Option<String>,
283 #[serde(default)]
284 pub allow_shell: Option<bool>,
285 #[serde(default)]
286 pub trust_mode: Option<bool>,
287 #[serde(default)]
288 pub auto_approve: Option<bool>,
289 #[serde(default)]
290 pub delivery_mode: Option<AutomationDeliveryMode>,
291 #[serde(default)]
292 pub status: Option<AutomationStatus>,
293 }
294
295 #[derive(Debug, Clone, Serialize, Deserialize, Default)]
296 pub struct UpdateAutomationRequest {
297 pub name: Option<String>,
298 pub prompt: Option<String>,
299 pub rrule: Option<String>,
300 pub cwds: Option<Vec<PathBuf>>,
301 pub model: Option<String>,
302 pub model_provider: Option<String>,
303 pub model_provider_id: Option<String>,
304 pub mode: Option<String>,
305 pub allow_shell: Option<bool>,
306 pub trust_mode: Option<bool>,
307 pub auto_approve: Option<bool>,
308 pub delivery_mode: Option<AutomationDeliveryMode>,
309 pub status: Option<AutomationStatus>,
310 }
311
312 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
313 enum AutomationFrequency {
314 Hourly,
315 Weekly,
316 }
317
318 #[derive(Debug, Clone)]
319 pub enum AutomationSchedule {
320 Once {
321 at: DateTime<Utc>,
322 },
323 Hourly {
324 interval_hours: u32,
325 byday: Option<Vec<Weekday>>,
326 anchor_hour: Option<u32>,
327 anchor_minute: Option<u32>,
328 },
329 Weekly {
330 byday: Vec<Weekday>,
331 byhour: u32,
332 byminute: u32,
333 },
334 Cron {
335 expr: String,
336 },
337 }
338
339 impl AutomationSchedule {
340 pub fn parse_rrule(rrule: &str) -> Result<Self> {
341 let mut parts: BTreeMap<String, String> = BTreeMap::new();
342 for raw in rrule.split(';') {
343 let item = raw.trim();
344 if item.is_empty() {
345 continue;
346 }
347 let Some((k, v)) = item.split_once('=') else {
348 bail!("Invalid RRULE segment '{item}'");
349 };
350 parts.insert(k.trim().to_ascii_uppercase(), v.trim().to_string());
351 }
352
353 let freq = match parts
354 .get("FREQ")
355 .map(|value| value.trim().to_ascii_uppercase())
356 .as_deref()
357 {
358 Some("ONCE") => return parse_once_schedule(&parts),
359 Some("HOURLY") => AutomationFrequency::Hourly,
360 Some("WEEKLY") => AutomationFrequency::Weekly,
361 Some("CRON") => return parse_cron_schedule(&parts),
362 Some(other) => {
363 bail!("Unsupported RRULE FREQ '{other}'. Supported: ONCE, HOURLY, WEEKLY, and CRON")
364 }
365 None => bail!("RRULE must include FREQ"),
366 };
367
368 match freq {
369 AutomationFrequency::Hourly => {
370 for key in parts.keys() {
371 if key != "FREQ"
372 && key != "INTERVAL"
373 && key != "BYDAY"
374 && key != "BYHOUR"
375 && key != "BYMINUTE"
376 {
377 bail!(
378 "Unsupported RRULE field '{key}' for HOURLY. Allowed: FREQ,INTERVAL,BYDAY,BYHOUR,BYMINUTE"
379 );
380 }
381 }
382 let interval_hours = parts
383 .get("INTERVAL")
384 .map(|v| v.parse::<u32>())
385 .transpose()
386 .context("Failed to parse INTERVAL")?
387 .unwrap_or(1);
388 if interval_hours == 0 {
389 bail!("INTERVAL must be >= 1 for HOURLY schedules");
390 }
391 let byday = parts
392 .get("BYDAY")
393 .map(|value| parse_byday(&value.to_ascii_uppercase()))
394 .transpose()?;
395 let anchor_hour = parts
396 .get("BYHOUR")
397 .map(|value| value.parse::<u32>())
398 .transpose()
399 .context("Failed to parse BYHOUR")?;
400 let anchor_minute = parts
401 .get("BYMINUTE")
402 .map(|value| value.parse::<u32>())
403 .transpose()
404 .context("Failed to parse BYMINUTE")?;
405 if anchor_hour.is_some_and(|hour| hour > 23) {
406 bail!("BYHOUR must be between 0 and 23");
407 }
408 if anchor_minute.is_some_and(|minute| minute > 59) {
409 bail!("BYMINUTE must be between 0 and 59");
410 }
411 Ok(Self::Hourly {
412 interval_hours,
413 byday,
414 anchor_hour,
415 anchor_minute,
416 })
417 }
418 AutomationFrequency::Weekly => {
419 for key in parts.keys() {
420 if key != "FREQ" && key != "BYDAY" && key != "BYHOUR" && key != "BYMINUTE" {
421 bail!(
422 "Unsupported RRULE field '{key}' for WEEKLY. Allowed: FREQ,BYDAY,BYHOUR,BYMINUTE"
423 );
424 }
425 }
426 let byday_raw = parts
427 .get("BYDAY")
428 .ok_or_else(|| anyhow::anyhow!("WEEKLY schedules require BYDAY"))?;
429 let byday = parse_byday(&byday_raw.to_ascii_uppercase())?;
430 if byday.is_empty() {
431 bail!("BYDAY cannot be empty for WEEKLY schedules");
432 }
433 let byhour = parts
434 .get("BYHOUR")
435 .ok_or_else(|| anyhow::anyhow!("WEEKLY schedules require BYHOUR"))?
436 .parse::<u32>()
437 .context("Failed to parse BYHOUR")?;
438 let byminute = parts
439 .get("BYMINUTE")
440 .ok_or_else(|| anyhow::anyhow!("WEEKLY schedules require BYMINUTE"))?
441 .parse::<u32>()
442 .context("Failed to parse BYMINUTE")?;
443
444 if byhour > 23 {
445 bail!("BYHOUR must be between 0 and 23");
446 }
447 if byminute > 59 {
448 bail!("BYMINUTE must be between 0 and 59");
449 }
450
451 Ok(Self::Weekly {
452 byday,
453 byhour,
454 byminute,
455 })
456 }
457 }
458 }
459
460 pub(crate) fn next_after_with_anchor(
461 &self,
462 after: DateTime<Utc>,
463 anchor_reference: DateTime<Utc>,
464 ) -> Result<DateTime<Utc>> {
465 self.next_after_in_timezone(after, anchor_reference, &Local)
466 }
467
468 fn next_after_in_timezone<Tz: TimeZone>(
469 &self,
470 after: DateTime<Utc>,
471 anchor_reference: DateTime<Utc>,
472 timezone: &Tz,
473 ) -> Result<DateTime<Utc>> {
474 let local_after = after.with_timezone(timezone);
475 match self {
476 Self::Once { at } => {
477 if *at > after {
478 Ok(*at)
479 } else {
480 bail!(
481 "Once schedule has no future run after {}",
482 after.to_rfc3339()
483 )
484 }
485 }
486 Self::Hourly {
487 interval_hours,
488 byday,
489 anchor_hour,
490 anchor_minute,
491 } => {
492 if anchor_hour.is_some() || anchor_minute.is_some() {
493 let local_anchor_reference = anchor_reference.with_timezone(timezone);
494 let hour = anchor_hour.unwrap_or(local_anchor_reference.hour());
495 let minute = anchor_minute.unwrap_or(0);
496 let anchor_naive = local_anchor_reference
497 .date_naive()
498 .and_hms_opt(hour, minute, 0)
499 .ok_or_else(|| anyhow::anyhow!("Unable to construct HOURLY anchor"))?;
500 let interval_seconds = i64::from(*interval_hours) * 60 * 60;
501 let elapsed_seconds = local_after
502 .naive_local()
503 .signed_duration_since(anchor_naive)
504 .num_seconds();
505 let mut steps = if elapsed_seconds < 0 {
506 0
507 } else {
508 elapsed_seconds / interval_seconds + 1
509 };
510
511 for _ in 0..MAX_HOURLY_SEARCH_STEPS {
512 let hours = i64::from(*interval_hours)
513 .checked_mul(steps)
514 .ok_or_else(|| anyhow::anyhow!("HOURLY schedule exceeded its range"))?;
515 let delta = Duration::try_hours(hours)
516 .ok_or_else(|| anyhow::anyhow!("HOURLY schedule exceeded its range"))?;
517 let candidate_naive = anchor_naive
518 .checked_add_signed(delta)
519 .ok_or_else(|| anyhow::anyhow!("HOURLY schedule exceeded its range"))?;
520
521 if byday
522 .as_ref()
523 .is_none_or(|days| days.contains(&candidate_naive.weekday()))
524 && let Some(candidate) =
525 resolve_local_datetime(timezone, candidate_naive)
526 {
527 let candidate = candidate.with_timezone(&Utc);
528 if candidate > after {
529 return Ok(candidate);
530 }
531 }
532
533 steps = steps
534 .checked_add(1)
535 .ok_or_else(|| anyhow::anyhow!("HOURLY schedule exceeded its range"))?;
536 }
537 bail!("Unable to compute next anchored HOURLY run");
538 }
539
540 let after_second = local_after.second();
541 let after_nanosecond = local_after.nanosecond();
542 let mut candidate = local_after + Duration::hours(i64::from(*interval_hours))
543 - Duration::seconds(i64::from(after_second))
544 - Duration::nanoseconds(i64::from(after_nanosecond));
545
546 if let Some(days) = byday {
547 for _ in 0..(24 * 21) {
548 if days.contains(&candidate.weekday()) {
549 return Ok(candidate.with_timezone(&Utc));
550 }
551 candidate += Duration::hours(i64::from(*interval_hours));
552 }
553 bail!("Unable to compute next HOURLY run for BYDAY filter");
554 }
555
556 Ok(candidate.with_timezone(&Utc))
557 }
558 Self::Weekly {
559 byday,
560 byhour,
561 byminute,
562 } => {
563 for day_offset in 0..15 {
564 let date = local_after.date_naive() + Duration::days(i64::from(day_offset));
565 if !byday.contains(&date.weekday()) {
566 continue;
567 }
568 let Some(candidate_naive) = date.and_hms_opt(*byhour, *byminute, 0) else {
569 continue;
570 };
571 if let Some(candidate) = resolve_local_datetime(timezone, candidate_naive)
572 && candidate.with_timezone(&Utc) > after
573 {
574 return Ok(candidate.with_timezone(&Utc));
575 }
576 }
577 bail!("Unable to compute next WEEKLY run");
578 }
579 Self::Cron { expr } => {
580 let cron = ParsedCronExpr::parse(expr)?;
581 let mut candidate_naive = local_after
582 .naive_local()
583 .with_second(0)
584 .and_then(|dt| dt.with_nanosecond(0))
585 .ok_or_else(|| anyhow::anyhow!("Unable to round CRON search start"))?
586 .checked_add_signed(Duration::minutes(1))
587 .ok_or_else(|| anyhow::anyhow!("CRON schedule exceeded its range"))?;
588
589 for _ in 0..MAX_CRON_SEARCH_MINUTES {
590 if cron.matches(candidate_naive)
591 && let Some(candidate) = resolve_local_datetime(timezone, candidate_naive)
592 {
593 let candidate = candidate.with_timezone(&Utc);
594 if candidate > after {
595 return Ok(candidate);
596 }
597 }
598 candidate_naive = candidate_naive
599 .checked_add_signed(Duration::minutes(1))
600 .ok_or_else(|| anyhow::anyhow!("CRON schedule exceeded its range"))?;
601 }
602 bail!("Unable to compute next CRON run within 5 years");
603 }
604 }
605 }
606
607 fn next_after_slot(
608 &self,
609 slot: DateTime<Utc>,
610 anchor_reference: DateTime<Utc>,
611 ) -> Result<Option<DateTime<Utc>>> {
612 match self {
613 Self::Once { .. } => Ok(None),
614 _ => self
615 .next_after_with_anchor(slot, anchor_reference)
616 .map(Some),
617 }
618 }
619
620 /// First slot after `slot` that is still in the future at `now`.
621 ///
622 /// Missed slots coalesce: downtime, a paused window, or an in-flight
623 /// occurrence earn one receipt for the oldest owed slot, then the
624 /// schedule resumes on its own grid instead of replaying one stale slot
625 /// per tick. Calendar-anchored schedules (anchored HOURLY, WEEKLY, CRON)
626 /// live on a fixed wall-clock grid, so the first slot after `now` is
627 /// exactly the slot plain chaining would converge to; computing from
628 /// `now` directly skips the whole backlog in one step. Unanchored HOURLY
629 /// is a relative cadence with no calendar grid — hop along its
630 /// established `slot + k * interval` chain so a late recovery does not
631 /// re-phase the schedule to the recovery instant.
632 fn next_unskipped_slot(
633 &self,
634 slot: DateTime<Utc>,
635 now: DateTime<Utc>,
636 anchor_reference: DateTime<Utc>,
637 ) -> Result<Option<DateTime<Utc>>> {
638 if let Self::Hourly {
639 interval_hours,
640 anchor_hour: None,
641 anchor_minute: None,
642 ..
643 } = self
644 {
645 let first = self.next_after_with_anchor(slot, anchor_reference)?;
646 if first > now {
647 return Ok(Some(first));
648 }
649 // Jump whole intervals on the established UTC grid, then reuse
650 // the schedule's weekday filter for the next eligible slot.
651 // Minute normalization happens in the first advance above.
652 let interval_seconds = i64::from(*interval_hours) * 60 * 60;
653 let elapsed = (now - first).num_seconds();
654 let delta = Duration::seconds(elapsed / interval_seconds * interval_seconds);
655 let previous = first
656 .checked_add_signed(delta)
657 .context("HOURLY catch-up exceeded its range")?;
658 self.next_after_slot(previous, anchor_reference)
659 } else {
660 self.next_after_slot(slot.max(now), anchor_reference)
661 }
662 }
663 }
664
665 /// Resolve one calendar-local schedule slot.
666 ///
667 /// Nonexistent wall times in a forward clock change are skipped rather than
668 /// shifted to a different clock time. Ambiguous wall times in a backward clock
669 /// change use the first occurrence only, preventing a recurring automation from
670 /// running twice for one calendar slot.
671 fn resolve_local_datetime<Tz: TimeZone>(
672 timezone: &Tz,
673 naive: NaiveDateTime,
674 ) -> Option<DateTime<Tz>> {
675 timezone.from_local_datetime(&naive).earliest()
676 }
677
678 fn parse_byday(value: &str) -> Result<Vec<Weekday>> {
679 let mut days = Vec::new();
680 for token in value.split(',') {
681 let day = match token.trim().to_ascii_uppercase().as_str() {
682 "MO" => Weekday::Mon,
683 "TU" => Weekday::Tue,
684 "WE" => Weekday::Wed,
685 "TH" => Weekday::Thu,
686 "FR" => Weekday::Fri,
687 "SA" => Weekday::Sat,
688 "SU" => Weekday::Sun,
689 other => bail!("Invalid BYDAY value '{other}'"),
690 };
691 if !days.contains(&day) {
692 days.push(day);
693 }
694 }
695 Ok(days)
696 }
697
698 fn parse_once_schedule(parts: &BTreeMap<String, String>) -> Result<AutomationSchedule> {
699 for key in parts.keys() {
700 if key != "FREQ" && key != "AT" {
701 bail!("Unsupported RRULE field '{key}' for ONCE. Allowed: FREQ,AT");
702 }
703 }
704 let raw_at = parts
705 .get("AT")
706 .ok_or_else(|| anyhow::anyhow!("ONCE schedules require AT"))?;
707 let at = parse_once_at(raw_at)?;
708 Ok(AutomationSchedule::Once { at })
709 }
710
711 fn parse_cron_schedule(parts: &BTreeMap<String, String>) -> Result<AutomationSchedule> {
712 for key in parts.keys() {
713 if key != "FREQ" && key != "EXPR" {
714 bail!("Unsupported RRULE field '{key}' for CRON. Allowed: FREQ,EXPR");
715 }
716 }
717 let expr = parts
718 .get("EXPR")
719 .map(|value| value.trim().to_string())
720 .filter(|value| !value.is_empty())
721 .ok_or_else(|| anyhow::anyhow!("CRON schedules require EXPR"))?;
722 ParsedCronExpr::parse(&expr)?;
723 Ok(AutomationSchedule::Cron { expr })
724 }
725
726 fn parse_once_at(raw: &str) -> Result<DateTime<Utc>> {
727 let trimmed = raw.trim();
728 if let Ok(at) = DateTime::parse_from_rfc3339(trimmed) {
729 return Ok(at.with_timezone(&Utc));
730 }
731 for format in ["%Y-%m-%dT%H:%M:%S", "%Y-%m-%dT%H:%M"] {
732 if let Ok(naive) = NaiveDateTime::parse_from_str(trimmed, format) {
733 return resolve_local_datetime(&Local, naive)
734 .map(|value| value.with_timezone(&Utc))
735 .ok_or_else(|| anyhow::anyhow!("ONCE local time does not exist: {trimmed}"));
736 }
737 }
738 bail!("Failed to parse ONCE AT '{trimmed}'. Use local YYYY-MM-DDTHH:MM[:SS] or RFC3339")
739 }
740
741 #[derive(Debug, Clone)]
742 struct ParsedCronExpr {
743 minute: CronField,
744 hour: CronField,
745 day_of_month: CronField,
746 month: CronField,
747 day_of_week: CronField,
748 }
749
750 impl ParsedCronExpr {
751 fn parse(expr: &str) -> Result<Self> {
752 let fields: Vec<&str> = expr.split_whitespace().collect();
753 if fields.len() != 5 {
754 bail!(
755 "CRON EXPR must have exactly 5 fields: minute hour day-of-month month day-of-week"
756 );
757 }
758 let parsed = Self {
759 minute: CronField::parse(fields[0], 0, 59, CronNameMap::none(), "minute")?,
760 hour: CronField::parse(fields[1], 0, 23, CronNameMap::none(), "hour")?,
761 day_of_month: CronField::parse(fields[2], 1, 31, CronNameMap::none(), "day-of-month")?,
762 month: CronField::parse(fields[3], 1, 12, CronNameMap::month(), "month")?,
763 day_of_week: CronField::parse(fields[4], 0, 7, CronNameMap::weekday(), "day-of-week")?
764 .normalized_day_of_week(),
765 };
766 parsed.validate_date_space()?;
767 Ok(parsed)
768 }
769
770 fn matches(&self, candidate: NaiveDateTime) -> bool {
771 if !self.minute.contains(candidate.minute())
772 || !self.hour.contains(candidate.hour())
773 || !self.month.contains(candidate.month())
774 {
775 return false;
776 }
777
778 let day_of_month = self.day_of_month.contains(candidate.day());
779 let weekday = self
780 .day_of_week
781 .contains(weekday_to_cron(candidate.weekday()));
782 if self.day_of_month.is_wildcard && self.day_of_week.is_wildcard {
783 true
784 } else if self.day_of_month.is_wildcard {
785 weekday
786 } else if self.day_of_week.is_wildcard {
787 day_of_month
788 } else {
789 day_of_month || weekday
790 }
791 }
792
793 fn validate_date_space(&self) -> Result<()> {
794 if self.day_of_month.is_wildcard {
795 return Ok(());
796 }
797 let months = self.month.values();
798 let days = self.day_of_month.values();
799 let valid = months.iter().copied().any(|month| {
800 let common = days_in_month(2025, month);
801 let leap = days_in_month(2024, month);
802 days.iter().copied().any(|day| day <= common || day <= leap)
803 });
804 if valid {
805 Ok(())
806 } else {
807 bail!("CRON EXPR day-of-month/month combination can never occur")
808 }
809 }
810 }
811
812 #[derive(Debug, Clone)]
813 struct CronField {
814 values: Vec<u32>,
815 is_wildcard: bool,
816 }
817
818 impl CronField {
819 fn parse(raw: &str, min: u32, max: u32, names: CronNameMap, field_name: &str) -> Result<Self> {
820 let trimmed = raw.trim();
821 if trimmed.is_empty() {
822 bail!("CRON {field_name} field must not be empty");
823 }
824 let mut values = Vec::new();
825 let is_wildcard = trimmed == "*";
826 for part in trimmed.split(',') {
827 let part = part.trim();
828 if part.is_empty() {
829 bail!("CRON {field_name} field contains an empty list item");
830 }
831 let (base, step) = if let Some((base, step)) = part.split_once('/') {
832 let step = step
833 .trim()
834 .parse::<u32>()
835 .with_context(|| format!("Failed to parse CRON {field_name} step"))?;
836 if step == 0 {
837 bail!("CRON {field_name} step must be >= 1");
838 }
839 (base.trim(), step)
840 } else {
841 (part, 1)
842 };
843
844 let range = if base == "*" {
845 (min, max)
846 } else if let Some((start, end)) = base.split_once('-') {
847 let start = parse_cron_atom(start.trim(), min, max, names, field_name)?;
848 let end = parse_cron_atom(end.trim(), min, max, names, field_name)?;
849 if start > end {
850 bail!("CRON {field_name} range start must be <= end");
851 }
852 (start, end)
853 } else {
854 let start = parse_cron_atom(base, min, max, names, field_name)?;
855 if part.contains('/') {
856 (start, max)
857 } else {
858 (start, start)
859 }
860 };
861
862 let mut current = range.0;
863 while current <= range.1 {
864 if !values.contains(&current) {
865 values.push(current);
866 }
867 let Some(next) = current.checked_add(step) else {
868 break;
869 };
870 if next <= current {
871 break;
872 }
873 current = next;
874 }
875 }
876 values.sort_unstable();
877 Ok(Self {
878 values,
879 is_wildcard,
880 })
881 }
882
883 fn normalized_day_of_week(mut self) -> Self {
884 for value in &mut self.values {
885 if *value == 7 {
886 *value = 0;
887 }
888 }
889 self.values.sort_unstable();
890 self.values.dedup();
891 self
892 }
893
894 fn contains(&self, value: u32) -> bool {
895 self.values.binary_search(&value).is_ok()
896 }
897
898 fn values(&self) -> &[u32] {
899 &self.values
900 }
901 }
902
903 #[derive(Debug, Clone, Copy)]
904 struct CronNameMap(&'static [(&'static str, u32)]);
905
906 impl CronNameMap {
907 const fn none() -> Self {
908 Self(&[])
909 }
910
911 const fn month() -> Self {
912 Self(&[
913 ("JAN", 1),
914 ("FEB", 2),
915 ("MAR", 3),
916 ("APR", 4),
917 ("MAY", 5),
918 ("JUN", 6),
919 ("JUL", 7),
920 ("AUG", 8),
921 ("SEP", 9),
922 ("OCT", 10),
923 ("NOV", 11),
924 ("DEC", 12),
925 ])
926 }
927
928 const fn weekday() -> Self {
929 Self(&[
930 ("SUN", 0),
931 ("MON", 1),
932 ("TUE", 2),
933 ("WED", 3),
934 ("THU", 4),
935 ("FRI", 5),
936 ("SAT", 6),
937 ])
938 }
939
940 fn lookup(self, token: &str) -> Option<u32> {
941 let needle = token.trim().to_ascii_uppercase();
942 self.0
943 .iter()
944 .find_map(|(name, value)| (*name == needle).then_some(*value))
945 }
946 }
947
948 fn parse_cron_atom(
949 raw: &str,
950 min: u32,
951 max: u32,
952 names: CronNameMap,
953 field_name: &str,
954 ) -> Result<u32> {
955 let value = names
956 .lookup(raw)
957 .or_else(|| raw.parse::<u32>().ok())
958 .ok_or_else(|| anyhow::anyhow!("Invalid CRON {field_name} value '{raw}'"))?;
959 if !(min..=max).contains(&value) {
960 bail!("CRON {field_name} value {value} is out of range {min}-{max}");
961 }
962 Ok(value)
963 }
964
965 fn weekday_to_cron(day: Weekday) -> u32 {
966 match day {
967 Weekday::Sun => 0,
968 Weekday::Mon => 1,
969 Weekday::Tue => 2,
970 Weekday::Wed => 3,
971 Weekday::Thu => 4,
972 Weekday::Fri => 5,
973 Weekday::Sat => 6,
974 }
975 }
976
977 fn days_in_month(year: i32, month: u32) -> u32 {
978 match month {
979 1 | 3 | 5 | 7 | 8 | 10 | 12 => 31,
980 4 | 6 | 9 | 11 => 30,
981 2 => {
982 let leap = (year % 4 == 0 && year % 100 != 0) || year % 400 == 0;
983 if leap { 29 } else { 28 }
984 }
985 _ => 0,
986 }
987 }
988
989 #[derive(Debug, Clone)]
990 pub struct AutomationManager {
991 execution_scope: Option<String>,
992 automations_dir: PathBuf,
993 runs_dir: PathBuf,
994 triggers_dir: PathBuf,
995 }
996
997 impl AutomationManager {
998 fn open_lock(&self, name: &str) -> Result<fd_lock::RwLock<fs::File>> {
999 let path = self
1000 .automations_dir
1001 .parent()
1002 .context("automation root")?
1003 .join(name);
1004 let mut options = fs::OpenOptions::new();
1005 options.create(true).truncate(false).read(true).write(true);
1006 #[cfg(unix)]
1007 {
1008 use std::os::unix::fs::OpenOptionsExt as _;
1009 options
1010 .mode(0o600)
1011 .custom_flags(libc::O_NOFOLLOW | libc::O_CLOEXEC);
1012 }
1013 #[cfg(windows)]
1014 {
1015 use std::os::windows::fs::OpenOptionsExt as _;
1016 options.custom_flags(0x0020_0000); // FILE_FLAG_OPEN_REPARSE_POINT
1017 }
1018 let file = options
1019 .open(&path)
1020 .with_context(|| format!("open {}", path.display()))?;
1021 let metadata = file.metadata()?;
1022 if !metadata.is_file() {
1023 bail!("Automation lock must be a regular file");
1024 }
1025 #[cfg(unix)]
1026 {
1027 use std::os::unix::fs::MetadataExt as _;
1028 if metadata.nlink() != 1 {
1029 bail!("Automation lock must not have hard links");
1030 }
1031 }
1032 #[cfg(windows)]
1033 {
1034 use std::os::windows::fs::MetadataExt as _;
1035 if metadata.file_attributes() & 0x400 != 0 {
1036 bail!("Automation lock must not be a reparse point");
1037 }
1038 }
1039 Ok(fd_lock::RwLock::new(file))
1040 }
1041
1042 fn with_transaction<T>(&self, operation: impl FnOnce() -> Result<T>) -> Result<T> {
1043 let mut lock = self.open_lock("state.lock")?;
1044 let _guard = lock.write().context("lock automation state")?;
1045 operation()
1046 }
1047
1048 /// Short read/modify/write transaction shared with scheduler admission.
1049 /// Returning None leaves an absent record absent; it does not delete one.
1050 pub(crate) fn edit_automation(
1051 &self,
1052 id: &str,
1053 edit: impl FnOnce(Option<AutomationRecord>) -> Result<Option<AutomationRecord>>,
1054 ) -> Result<Option<AutomationRecord>> {
1055 self.with_transaction(|| {
1056 let current = if self.automation_path(id)?.try_exists()? {
1057 Some(self.get_automation(id)?)
1058 } else {
1059 None
1060 };
1061 let edited = edit(current)?;
1062 if let Some(record) = &edited {
1063 if record.id != id {
1064 bail!("Automation transaction cannot replace its identity");
1065 }
1066 self.save_automation_unlocked(record)?;
1067 }
1068 Ok(edited)
1069 })
1070 }
1071
1072 pub fn open(root: PathBuf) -> Result<Self> {
1073 let automations_dir = root.join("automations");
1074 let runs_dir = root.join("runs");
1075 let triggers_dir = root.join("triggers");
1076 fs::create_dir_all(&automations_dir)
1077 .with_context(|| format!("Failed to create {}", automations_dir.display()))?;
1078 fs::create_dir_all(&runs_dir)
1079 .with_context(|| format!("Failed to create {}", runs_dir.display()))?;
1080 fs::create_dir_all(&triggers_dir)
1081 .with_context(|| format!("Failed to create {}", triggers_dir.display()))?;
1082 Ok(Self {
1083 execution_scope: None,
1084 automations_dir,
1085 runs_dir,
1086 triggers_dir,
1087 })
1088 }
1089
1090 #[cfg(test)]
1091 pub(crate) fn open_for_test(root: PathBuf) -> Result<Self> {
1092 let mut manager = Self::open(root)?;
1093 manager.execution_scope = Some(crate::task_manager::test_execution_scope("test"));
1094 Ok(manager)
1095 }
1096
1097 pub(crate) fn bind_task_manager(
1098 &mut self,
1099 tasks: &crate::task_manager::TaskManager,
1100 ) -> Result<()> {
1101 if self
1102 .execution_scope
1103 .as_deref()
1104 .is_some_and(|scope| scope != tasks.execution_scope())
1105 {
1106 bail!("Automation service belongs to another Runtime scope");
1107 }
1108 self.execution_scope = Some(tasks.execution_scope().to_string());
1109 Ok(())
1110 }
1111
1112 pub(crate) fn execution_scope(&self) -> Option<&str> {
1113 self.execution_scope.as_deref()
1114 }
1115
1116 fn eligible_scope(&self, scope: Option<&str>) -> bool {
1117 scope.is_some() && scope == self.execution_scope()
1118 }
1119
1120 /// Explicit control may bind an unbound definition; saved admissions never
1121 /// read this field back from the definition during recovery.
1122 fn adopt_for_run(&self, automation: &mut AutomationRecord) -> Result<()> {
1123 let scope = self
1124 .execution_scope()
1125 .context("Automation execution ownership is unverified")?;
1126 if let Some(bound) = &automation.execution_scope {
1127 if bound != scope {
1128 bail!("Automation belongs to another Runtime execution scope");
1129 }
1130 } else {
1131 automation.execution_scope = Some(scope.to_string());
1132 automation.schema_version = CURRENT_AUTOMATION_SCHEMA_VERSION;
1133 automation.updated_at = Utc::now();
1134 if automation.status == AutomationStatus::Active {
1135 let schedule = AutomationSchedule::parse_rrule(&automation.rrule)?;
1136 automation.next_run_at =
1137 match schedule.next_after_with_anchor(Utc::now(), automation.created_at) {
1138 Ok(next) => Some(next),
1139 Err(_) if matches!(schedule, AutomationSchedule::Once { .. }) => {
1140 automation.status = AutomationStatus::Paused;
1141 None
1142 }
1143 Err(error) => return Err(error),
1144 };
1145 }
1146 self.save_automation_unlocked(automation)?;
1147 }
1148 Ok(())
1149 }
1150
1151 pub fn default_location() -> Result<Self> {
1152 Self::open(default_automations_dir())
1153 }
1154
1155 fn automation_path(&self, id: &str) -> Result<PathBuf> {
1156 ensure_safe_storage_id("automation id", id)?;
1157 Ok(self.automations_dir.join(format!("{id}.json")))
1158 }
1159
1160 fn runs_dir_for(&self, automation_id: &str) -> Result<PathBuf> {
1161 ensure_safe_storage_id("automation id", automation_id)?;
1162 Ok(self.runs_dir.join(automation_id))
1163 }
1164
1165 fn trigger_path(&self, trigger_id: &str) -> Result<PathBuf> {
1166 ensure_safe_storage_id("trigger id", trigger_id)?;
1167 Ok(self.triggers_dir.join(format!("{trigger_id}.json")))
1168 }
1169
1170 /// Current run file name: `{sortable-created-at}-{run_id}.json`. The
1171 /// fixed-width timestamp prefix makes directory listings sort
1172 /// chronologically without reading file contents (see [`Self::list_runs`]).
1173 fn run_path(&self, run: &AutomationRunRecord) -> Result<PathBuf> {
1174 ensure_safe_storage_id("run id", &run.id)?;
1175 Ok(self.runs_dir_for(&run.automation_id)?.join(format!(
1176 "{}-{}.json",
1177 run_file_stamp(run.created_at),
1178 run.id
1179 )))
1180 }
1181
1182 /// Pre-sortable-name run file: `{run_id}.json` (run ids are UUIDs, so
1183 /// these carry no ordering hint and must be read to learn `created_at`).
1184 fn legacy_run_path(&self, automation_id: &str, run_id: &str) -> Result<PathBuf> {
1185 ensure_safe_storage_id("run id", run_id)?;
1186 Ok(self
1187 .runs_dir_for(automation_id)?
1188 .join(format!("{run_id}.json")))
1189 }
1190
1191 pub fn create_automation(&self, req: CreateAutomationRequest) -> Result<AutomationRecord> {
1192 validate_name_and_prompt(&req.name, &req.prompt)?;
1193 let schedule = AutomationSchedule::parse_rrule(&req.rrule)?;
1194 let now = Utc::now();
1195 let status = req.status.unwrap_or(AutomationStatus::Active);
1196 let next_run_at = if matches!(status, AutomationStatus::Active) {
1197 Some(schedule.next_after_with_anchor(now, now)?)
1198 } else {
1199 None
1200 };
1201
1202 let record = AutomationRecord {
1203 schema_version: CURRENT_AUTOMATION_SCHEMA_VERSION,
1204 execution_scope: self.execution_scope.clone(),
1205 id: Uuid::new_v4().to_string(),
1206 name: req.name.trim().to_string(),
1207 prompt: req.prompt.trim().to_string(),
1208 rrule: req.rrule.trim().to_ascii_uppercase(),
1209 cwds: req.cwds,
1210 model: normalize_optional_string(req.model),
1211 model_provider: normalize_optional_string(req.model_provider),
1212 model_provider_id: normalize_optional_string(req.model_provider_id),
1213 mode: normalize_optional_string(req.mode),
1214 allow_shell: req.allow_shell,
1215 trust_mode: req.trust_mode,
1216 auto_approve: req.auto_approve,
1217 delivery_mode: req.delivery_mode,
1218 status,
1219 created_at: now,
1220 updated_at: now,
1221 next_run_at,
1222 last_run_at: None,
1223 };
1224
1225 self.save_automation(&record)?;
1226 Ok(record)
1227 }
1228
1229 pub fn get_automation(&self, id: &str) -> Result<AutomationRecord> {
1230 let path = self.automation_path(id)?;
1231 read_automation_file(&path)
1232 }
1233
1234 pub fn save_automation(&self, record: &AutomationRecord) -> Result<()> {
1235 self.with_transaction(|| self.save_automation_unlocked(record))
1236 }
1237
1238 fn save_automation_unlocked(&self, record: &AutomationRecord) -> Result<()> {
1239 if record.model_provider.is_some() || record.model_provider_id.is_some() {
1240 if record
1241 .model
1242 .as_deref()
1243 .is_none_or(|model| model.trim().is_empty())
1244 {
1245 bail!("A pinned automation provider requires an explicit model");
1246 }
1247 let mut record = record.clone();
1248 record.schema_version = record.schema_version.max(CURRENT_AUTOMATION_SCHEMA_VERSION);
1249 return write_json_atomic(&self.automation_path(&record.id)?, &record);
1250 }
1251 write_json_atomic(&self.automation_path(&record.id)?, record)
1252 }
1253
1254 pub fn list_automations(&self) -> Result<Vec<AutomationRecord>> {
1255 let mut out = Vec::new();
1256 for entry in fs::read_dir(&self.automations_dir)
1257 .with_context(|| format!("Failed to read {}", self.automations_dir.display()))?
1258 {
1259 let entry = entry?;
1260 let path = entry.path();
1261 if path.extension().is_none_or(|ext| ext != "json") {
1262 continue;
1263 }
1264 let record = read_automation_file(&path)?;
1265 out.push(record);
1266 }
1267 out.sort_by_key(|r| std::cmp::Reverse(r.updated_at));
1268 Ok(out)
1269 }
1270
1271 pub fn update_automation(
1272 &self,
1273 id: &str,
1274 req: UpdateAutomationRequest,
1275 ) -> Result<AutomationRecord> {
1276 self.with_transaction(|| self.update_automation_unlocked(id, req))
1277 }
1278
1279 fn update_automation_unlocked(
1280 &self,
1281 id: &str,
1282 req: UpdateAutomationRequest,
1283 ) -> Result<AutomationRecord> {
1284 let mut existing = self.get_automation(id)?;
1285 let adopting = existing.execution_scope.is_none()
1286 && self.execution_scope.is_some()
1287 && req.status != Some(AutomationStatus::Paused);
1288 if adopting {
1289 existing.execution_scope = self.execution_scope.clone();
1290 }
1291 let schedule_changed = adopting || req.rrule.is_some() || req.status.is_some();
1292
1293 if let Some(name) = req.name {
1294 if name.trim().is_empty() {
1295 bail!("Automation name cannot be empty");
1296 }
1297 existing.name = name.trim().to_string();
1298 }
1299 if let Some(prompt) = req.prompt {
1300 if prompt.trim().is_empty() {
1301 bail!("Automation prompt cannot be empty");
1302 }
1303 existing.prompt = prompt.trim().to_string();
1304 }
1305 if let Some(rrule) = req.rrule {
1306 let normalized = rrule.trim().to_ascii_uppercase();
1307 AutomationSchedule::parse_rrule(&normalized)?;
1308 existing.rrule = normalized;
1309 }
1310 if let Some(cwds) = req.cwds {
1311 existing.cwds = cwds;
1312 }
1313 if let Some(model) = req.model {
1314 existing.model = normalize_optional_string(Some(model));
1315 }
1316 if let Some(provider) = req.model_provider {
1317 existing.model_provider = normalize_optional_string(Some(provider));
1318 }
1319 if let Some(provider_id) = req.model_provider_id {
1320 existing.model_provider_id = normalize_optional_string(Some(provider_id));
1321 }
1322 if let Some(mode) = req.mode {
1323 existing.mode = normalize_optional_string(Some(mode));
1324 }
1325 if let Some(allow_shell) = req.allow_shell {
1326 existing.allow_shell = Some(allow_shell);
1327 }
1328 if let Some(trust_mode) = req.trust_mode {
1329 existing.trust_mode = Some(trust_mode);
1330 }
1331 if let Some(auto_approve) = req.auto_approve {
1332 existing.auto_approve = Some(auto_approve);
1333 }
1334 if let Some(delivery_mode) = req.delivery_mode {
1335 existing.delivery_mode = Some(delivery_mode);
1336 }
1337 if let Some(status) = req.status {
1338 existing.status = status;
1339 }
1340 // Evaluate the final status once: editing a schedule and pausing it is
1341 // one atomic update, and must not first schedule an active run.
1342 if schedule_changed {
1343 if matches!(existing.status, AutomationStatus::Paused) {
1344 existing.next_run_at = None;
1345 } else {
1346 let schedule = AutomationSchedule::parse_rrule(&existing.rrule)?;
1347 existing.next_run_at =
1348 Some(schedule.next_after_with_anchor(Utc::now(), existing.created_at)?);
1349 }
1350 }
1351
1352 if existing.execution_scope.is_some()
1353 || existing.model_provider.is_some()
1354 || existing.model_provider_id.is_some()
1355 {
1356 existing.schema_version = CURRENT_AUTOMATION_SCHEMA_VERSION;
1357 }
1358
1359 existing.updated_at = Utc::now();
1360 self.save_automation_unlocked(&existing)?;
1361 Ok(existing)
1362 }
1363
1364 pub fn pause_automation(&self, id: &str) -> Result<AutomationRecord> {
1365 self.update_automation(
1366 id,
1367 UpdateAutomationRequest {
1368 status: Some(AutomationStatus::Paused),
1369 ..UpdateAutomationRequest::default()
1370 },
1371 )
1372 }
1373
1374 pub fn resume_automation(&self, id: &str) -> Result<AutomationRecord> {
1375 self.update_automation(
1376 id,
1377 UpdateAutomationRequest {
1378 status: Some(AutomationStatus::Active),
1379 ..UpdateAutomationRequest::default()
1380 },
1381 )
1382 }
1383
1384 pub fn delete_automation(&self, id: &str) -> Result<AutomationRecord> {
1385 self.with_transaction(|| {
1386 let existing = self.get_automation(id)?;
1387 let path = self.automation_path(id)?;
1388 fs::remove_file(&path)
1389 .with_context(|| format!("Failed to delete automation {}", path.display()))?;
1390 // A claimed occurrence has already crossed the admission boundary.
1391 // Keep its binding through deletion so recovery cannot lose or repeat it.
1392 for run in self.list_runs_with_visibility(id, None, true)? {
1393 if !matches!(
1394 run.status,
1395 AutomationRunStatus::Queued | AutomationRunStatus::Running
1396 ) {
1397 self.delete_run(&run)?;
1398 }
1399 }
1400 let runs_dir = self.runs_dir_for(id)?;
1401 if runs_dir.try_exists()? && fs::read_dir(&runs_dir)?.next().is_none() {
1402 fs::remove_dir(&runs_dir).with_context(|| {
1403 format!(
1404 "Failed to remove empty run directory {}",
1405 runs_dir.display()
1406 )
1407 })?;
1408 }
1409 Ok(existing)
1410 })
1411 }
1412
1413 pub fn list_runs(
1414 &self,
1415 automation_id: &str,
1416 limit: Option<usize>,
1417 ) -> Result<Vec<AutomationRunRecord>> {
1418 self.list_runs_with_visibility(automation_id, limit, false)
1419 }
1420
1421 fn list_runs_with_visibility(
1422 &self,
1423 automation_id: &str,
1424 limit: Option<usize>,
1425 include_suppressed: bool,
1426 ) -> Result<Vec<AutomationRunRecord>> {
1427 let dir = self.runs_dir_for(automation_id)?;
1428 if !dir.exists() {
1429 return Ok(Vec::new());
1430 }
1431
1432 // Split the listing into sortable-name files (newest-first by file
1433 // name alone, so reads stop after the newest `limit`) and legacy
1434 // `{uuid}.json` files, which must all be read to learn `created_at`.
1435 let mut sortable = Vec::new();
1436 let mut legacy = Vec::new();
1437 for entry in
1438 fs::read_dir(&dir).with_context(|| format!("Failed to read {}", dir.display()))?
1439 {
1440 let entry = entry?;
1441 let path = entry.path();
1442 if path.extension().is_none_or(|ext| ext != "json") {
1443 continue;
1444 }
1445 if path
1446 .file_stem()
1447 .and_then(|stem| stem.to_str())
1448 .is_some_and(has_sortable_run_stem)
1449 {
1450 sortable.push(path);
1451 } else {
1452 legacy.push(path);
1453 }
1454 }
1455
1456 // A sortable receipt supersedes its legacy copy even when hidden or
1457 // older than the requested window. Fence identity before visibility.
1458 let sortable_ids: std::collections::BTreeSet<_> = sortable
1459 .iter()
1460 .filter_map(|path| path.file_stem()?.to_str()?.get(RUN_STAMP_LEN + 1..))
1461 .collect();
1462 legacy.retain(|path| {
1463 !path
1464 .file_stem()
1465 .and_then(|stem| stem.to_str())
1466 .is_some_and(|id| sortable_ids.contains(id))
1467 });
1468 sortable.sort_by(|a, b| b.file_name().cmp(&a.file_name()));
1469 let visible = |run: &AutomationRunRecord| {
1470 include_suppressed
1471 || !run
1472 .dispatch
1473 .as_ref()
1474 .is_some_and(|dispatch| dispatch.suppress_report)
1475 };
1476 let mut out = Vec::new();
1477 for path in sortable {
1478 if limit.is_some_and(|limit| out.len() >= limit) {
1479 break;
1480 }
1481 let run = read_run_file(&path)?;
1482 if visible(&run) {
1483 out.push(run);
1484 }
1485 }
1486 for path in legacy {
1487 let run = read_run_file(&path)?;
1488 if visible(&run) {
1489 out.push(run);
1490 }
1491 }
1492
1493 out.sort_by_key(|r| std::cmp::Reverse(r.created_at));
1494 // A crash between the sortable-name write and the legacy-file removal
1495 // in `save_run` can leave one run under both names; keep the sortable
1496 // copy (chained first above, so it survives the stable sort).
1497 out.dedup_by(|a, b| a.id == b.id);
1498 if let Some(limit) = limit {
1499 out.truncate(limit);
1500 }
1501 Ok(out)
1502 }
1503
1504 /// Re-read specific runs by id without a full history pass. Run file
1505 /// names end in `-{run_id}.json` (or are legacy `{run_id}.json`), so a
1506 /// directory listing locates them and only those files are read. Used
1507 /// by the activity-band scan to keep watching runs this session saw go
1508 /// live even after newer runs push them past the newest-run window.
1509 pub fn get_runs_by_ids(
1510 &self,
1511 automation_id: &str,
1512 run_ids: &std::collections::BTreeSet<String>,
1513 ) -> Result<Vec<AutomationRunRecord>> {
1514 if run_ids.is_empty() {
1515 return Ok(Vec::new());
1516 }
1517 let dir = self.runs_dir_for(automation_id)?;
1518 if !dir.exists() {
1519 return Ok(Vec::new());
1520 }
1521 let mut paths = BTreeMap::<String, (bool, PathBuf)>::new();
1522 for entry in
1523 fs::read_dir(&dir).with_context(|| format!("Failed to read {}", dir.display()))?
1524 {
1525 let path = entry?.path();
1526 if path.extension().and_then(|ext| ext.to_str()) != Some("json") {
1527 continue;
1528 }
1529 let Some(stem) = path.file_stem().and_then(|stem| stem.to_str()) else {
1530 continue;
1531 };
1532 let sortable = has_sortable_run_stem(stem);
1533 let id = if sortable {
1534 &stem[RUN_STAMP_LEN + 1..]
1535 } else {
1536 stem
1537 };
1538 if run_ids.contains(id)
1539 && paths
1540 .get(id)
1541 .is_none_or(|(current, _)| !current && sortable)
1542 {
1543 paths.insert(id.to_string(), (sortable, path));
1544 }
1545 }
1546 let mut out = Vec::new();
1547 for (_, path) in paths.into_values() {
1548 let run = read_run_file(&path)?;
1549 if !run
1550 .dispatch
1551 .as_ref()
1552 .is_some_and(|dispatch| dispatch.suppress_report)
1553 {
1554 out.push(run);
1555 }
1556 }
1557 Ok(out)
1558 }
1559
1560 fn save_run(&self, run: &AutomationRunRecord) -> Result<()> {
1561 let dir = self.runs_dir_for(&run.automation_id)?;
1562 fs::create_dir_all(&dir).with_context(|| format!("Failed to create {}", dir.display()))?;
1563 let path = self.run_path(run)?;
1564 write_json_atomic(&path, run)?;
1565 // Rewrites of a legacy-named run migrate it to the sortable name; drop
1566 // the old file so the run never exists twice.
1567 let legacy = self.legacy_run_path(&run.automation_id, &run.id)?;
1568 if legacy != path && legacy.exists() {
1569 fs::remove_file(&legacy)
1570 .with_context(|| format!("Failed to remove legacy run {}", legacy.display()))?;
1571 }
1572 Ok(())
1573 }
1574
1575 fn delete_run(&self, run: &AutomationRunRecord) -> Result<()> {
1576 let sortable = self.run_path(run)?;
1577 if sortable.exists() {
1578 fs::remove_file(&sortable)
1579 .with_context(|| format!("Failed to delete run {}", sortable.display()))?;
1580 }
1581 let legacy = self.legacy_run_path(&run.automation_id, &run.id)?;
1582 if legacy.exists() {
1583 fs::remove_file(&legacy)
1584 .with_context(|| format!("Failed to delete run {}", legacy.display()))?;
1585 }
1586 Ok(())
1587 }
1588
1589 /// Definitions this build can read, in `list_automations` order.
1590 ///
1591 /// One corrupt, unreadable, or newer-schema file is quarantined in place:
1592 /// its bytes stay on disk and every pass logs the path, but it cannot
1593 /// starve collection of the healthy definitions behind it, and it is never
1594 /// rewritten or adopted by a runtime that does not understand it.
1595 fn readable_automations(&self) -> Result<Vec<AutomationRecord>> {
1596 let mut out = Vec::new();
1597 for entry in fs::read_dir(&self.automations_dir)
1598 .with_context(|| format!("Failed to read {}", self.automations_dir.display()))?
1599 {
1600 let entry = entry?;
1601 let path = entry.path();
1602 if path.extension().is_none_or(|ext| ext != "json") {
1603 continue;
1604 }
1605 match read_automation_file(&path) {
1606 Ok(record) => out.push(record),
1607 Err(error) => {
1608 tracing::warn!("Skipping damaged automation file: {error:#}");
1609 }
1610 }
1611 }
1612 out.sort_by_key(|r| std::cmp::Reverse(r.updated_at));
1613 Ok(out)
1614 }
1615
1616 /// List proposals only. Every proposal is revalidated and durably claimed
1617 /// immediately before dispatch, not when an earlier batch item is awaited.
1618 fn collect_due_runs(
1619 &self,
1620 now: DateTime<Utc>,
1621 ) -> Result<Vec<(AutomationRecord, AutomationRunRecord)>> {
1622 self.with_transaction(|| {
1623 let mut due = Vec::new();
1624 for mut automation in self.readable_automations()? {
1625 if automation.status != AutomationStatus::Active
1626 || !self.eligible_scope(automation.execution_scope.as_deref())
1627 {
1628 continue;
1629 }
1630 // An owned definition whose schedule cannot be evaluated is
1631 // quarantined like a damaged file: left untouched (a newer
1632 // build may understand it), diagnosed every pass, and never
1633 // allowed to take down the rest of the collection.
1634 let schedule = match AutomationSchedule::parse_rrule(&automation.rrule) {
1635 Ok(schedule) => schedule,
1636 Err(error) => {
1637 tracing::warn!(
1638 "Skipping automation {} with unevaluable schedule {:?}: {error:#}",
1639 automation.id,
1640 automation.rrule
1641 );
1642 continue;
1643 }
1644 };
1645 let Some(due_at) = automation.next_run_at else {
1646 automation.next_run_at =
1647 match schedule.next_after_with_anchor(now, automation.created_at) {
1648 Ok(next) => Some(next),
1649 Err(error)
1650 if matches!(schedule, AutomationSchedule::Once { .. })
1651 && error
1652 .to_string()
1653 .contains("Once schedule has no future run") =>
1654 {
1655 automation.status = AutomationStatus::Paused;
1656 None
1657 }
1658 Err(error) => {
1659 tracing::warn!(
1660 "Skipping automation {} whose schedule cannot produce a slot: {error:#}",
1661 automation.id
1662 );
1663 continue;
1664 }
1665 };
1666 automation.updated_at = now;
1667 self.save_automation_unlocked(&automation)?;
1668 continue;
1669 };
1670 if due_at <= now {
1671 due.push((
1672 automation.clone(),
1673 new_run_record(&automation.id, due_at, now),
1674 ));
1675 }
1676 }
1677 Ok(due)
1678 })
1679 }
1680
1681 fn claim_scheduled_run(
1682 &self,
1683 observed: &AutomationRecord,
1684 mut run: AutomationRunRecord,
1685 task_data_dir: &Path,
1686 ) -> Result<Option<AutomationRunRecord>> {
1687 self.with_transaction(|| {
1688 if !self.automation_path(&observed.id)?.try_exists()? {
1689 return Ok(None);
1690 }
1691 let mut current = self.get_automation(&observed.id)?;
1692 if !self.eligible_scope(current.execution_scope.as_deref())
1693 || current != *observed
1694 || current.status != AutomationStatus::Active
1695 || current.next_run_at != Some(run.scheduled_for)
1696 {
1697 return Ok(None);
1698 }
1699 let schedule = AutomationSchedule::parse_rrule(&current.rrule)?;
1700 // Include the complete history, including legacy occurrence ids.
1701 let history = self.list_runs_with_visibility(&current.id, None, true)?;
1702 if history
1703 .iter()
1704 .any(|existing| existing.scheduled_for == run.scheduled_for)
1705 {
1706 self.advance_automation_after_slot(
1707 &mut current,
1708 &schedule,
1709 run.scheduled_for,
1710 Utc::now(),
1711 )?;
1712 return Ok(None);
1713 }
1714 // Keep the owed slot until the earlier run settles. Its eventual
1715 // catch-up coalesces the backlog without overlapping executions
1716 // or inventing cancellation receipts. Explicit run-now requests
1717 // remain operator intent and stay ungated.
1718 if history.iter().any(|existing| {
1719 matches!(
1720 existing.status,
1721 AutomationRunStatus::Queued | AutomationRunStatus::Running
1722 )
1723 }) {
1724 return Ok(None);
1725 }
1726 bind_run_dispatch(&mut run, &current, task_data_dir, true)?;
1727 self.save_run(&run)?;
1728 // The durable claim is the point of no return. Pause/delete after
1729 // this point affects future occurrences, not this admitted work.
1730 // No task can start before the binding above is durable.
1731 self.advance_automation_after_slot(
1732 &mut current,
1733 &schedule,
1734 run.scheduled_for,
1735 Utc::now(),
1736 )?;
1737 Ok(Some(run))
1738 })
1739 }
1740
1741 /// Repair only a torn claim/advance transaction. A replacement definition,
1742 /// pause, or Operate kick has a different generation or slot and wins.
1743 fn recover_schedule_advance(&self, run: &AutomationRunRecord) -> Result<()> {
1744 let Some(admitted) = run
1745 .dispatch
1746 .as_ref()
1747 .and_then(|dispatch| dispatch.schedule.as_ref())
1748 else {
1749 return Ok(());
1750 };
1751 self.with_transaction(|| {
1752 if !self.automation_path(&run.automation_id)?.try_exists()? {
1753 return Ok(());
1754 }
1755 let mut current = self.get_automation(&run.automation_id)?;
1756 if current.status == AutomationStatus::Active
1757 && current.updated_at == admitted.updated_at
1758 && current.rrule == admitted.rrule
1759 && current.next_run_at == Some(run.scheduled_for)
1760 {
1761 let schedule = AutomationSchedule::parse_rrule(&current.rrule)?;
1762 self.advance_automation_after_slot(
1763 &mut current,
1764 &schedule,
1765 run.scheduled_for,
1766 Utc::now(),
1767 )?;
1768 }
1769 Ok(())
1770 })
1771 }
1772
1773 /// Completion publishes only this occurrence's receipt. It never advances
1774 /// a definition that may have been edited during the enqueue await.
1775 fn finish_scheduled_run(&self, run: &AutomationRunRecord, now: DateTime<Utc>) -> Result<()> {
1776 self.with_transaction(|| {
1777 self.save_run(run)?;
1778 if matches!(
1779 run.status,
1780 AutomationRunStatus::Completed
1781 | AutomationRunStatus::Failed
1782 | AutomationRunStatus::Canceled
1783 ) && !run
1784 .dispatch
1785 .as_ref()
1786 .is_some_and(|dispatch| dispatch.suppress_report)
1787 && self.automation_path(&run.automation_id)?.try_exists()?
1788 {
1789 let mut current = self.get_automation(&run.automation_id)?;
1790 let ended = run.ended_at.unwrap_or(now);
1791 current.last_run_at = Some(
1792 current
1793 .last_run_at
1794 .map_or(ended, |previous| previous.max(ended)),
1795 );
1796 // Receipt metadata is not a new schedule generation. Never
1797 // recompute the current definition from this run's old slot.
1798 self.save_automation_unlocked(&current)?;
1799 }
1800 Ok(())
1801 })
1802 }
1803
1804 fn advance_automation_after_slot(
1805 &self,
1806 automation: &mut AutomationRecord,
1807 schedule: &AutomationSchedule,
1808 slot: DateTime<Utc>,
1809 now: DateTime<Utc>,
1810 ) -> Result<()> {
1811 automation.updated_at = now;
1812 automation.next_run_at = schedule.next_unskipped_slot(slot, now, automation.created_at)?;
1813 if automation.next_run_at.is_none() {
1814 automation.status = AutomationStatus::Paused;
1815 }
1816 self.save_automation_unlocked(automation)
1817 }
1818
1819 /// Active receipts are independent of the definition's lifetime and of
1820 /// presentation limits on recent history.
1821 fn collect_pending_runs(&self) -> Result<Vec<AutomationRunRecord>> {
1822 let mut pending = Vec::new();
1823 for entry in fs::read_dir(&self.runs_dir)? {
1824 let entry = entry?;
1825 if !entry.file_type()?.is_dir() {
1826 continue;
1827 }
1828 let Some(id) = entry.file_name().to_str().map(str::to_owned) else {
1829 continue;
1830 };
1831 // A damaged receipt quarantines only its own automation: the bytes
1832 // stay on disk and the diagnostic is logged every pass, but one
1833 // corrupt file must not block recovery of every other pending run.
1834 let runs = match self.list_runs_with_visibility(&id, None, true) {
1835 Ok(runs) => runs,
1836 Err(error) => {
1837 tracing::warn!("Skipping damaged run history for automation {id}: {error:#}");
1838 continue;
1839 }
1840 };
1841 for run in runs {
1842 if matches!(
1843 run.status,
1844 AutomationRunStatus::Queued | AutomationRunStatus::Running
1845 ) && run.task_id.is_some()
1846 {
1847 pending.push(run);
1848 }
1849 }
1850 }
1851 Ok(pending)
1852 }
1853
1854 // ── Delayed-trigger storage methods ──────────────────────────────────
1855
1856 /// Persist a new delayed trigger and return the record.
1857 pub fn create_trigger(&self, req: CreateDelayedTriggerRequest) -> Result<DelayedTriggerRecord> {
1858 let now = Utc::now();
1859 if req.fire_at <= now {
1860 bail!(
1861 "fire_at must be in the future (got {}, now is {})",
1862 req.fire_at.to_rfc3339(),
1863 now.to_rfc3339()
1864 );
1865 }
1866 if req.message.trim().is_empty() {
1867 bail!("Trigger message must not be empty");
1868 }
1869 let record = DelayedTriggerRecord {
1870 schema_version: CURRENT_TRIGGER_SCHEMA_VERSION,
1871 execution_scope: self.execution_scope.clone(),
1872 trigger_id: format!("trig_{}", Uuid::new_v4().simple()),
1873 fire_at: req.fire_at,
1874 message: req.message.trim().to_string(),
1875 workspace: req.workspace,
1876 owner_session_id: req.owner_session_id,
1877 status: DelayedTriggerStatus::Pending,
1878 created_at: now,
1879 fired_at: None,
1880 task_id: None,
1881 thread_id: None,
1882 error: None,
1883 parent_trigger_id: req.parent_trigger_id,
1884 dispatch: None,
1885 };
1886 self.save_trigger(&record)?;
1887 Ok(record)
1888 }
1889
1890 /// Load a trigger by id.
1891 pub fn get_trigger(&self, trigger_id: &str) -> Result<DelayedTriggerRecord> {
1892 let path = self.trigger_path(trigger_id)?;
1893 let raw = fs::read_to_string(&path)
1894 .with_context(|| format!("Trigger '{trigger_id}' not found"))?;
1895 let record: DelayedTriggerRecord = serde_json::from_str(&raw)
1896 .with_context(|| format!("Failed to parse trigger '{trigger_id}'"))?;
1897 if record.schema_version > CURRENT_TRIGGER_SCHEMA_VERSION {
1898 bail!(
1899 "Trigger schema v{} is newer than supported v{}",
1900 record.schema_version,
1901 CURRENT_TRIGGER_SCHEMA_VERSION
1902 );
1903 }
1904 Ok(record)
1905 }
1906
1907 /// Load a trigger only when it belongs to the given session.
1908 ///
1909 /// Foreign, ownerless legacy, unreadable, and absent records share the same
1910 /// result so trigger existence cannot be disclosed across sessions.
1911 pub fn get_trigger_for_owner(
1912 &self,
1913 trigger_id: &str,
1914 owner_session_id: &str,
1915 ) -> Result<DelayedTriggerRecord> {
1916 self.get_trigger(trigger_id)
1917 .ok()
1918 .filter(|record| record.owner_session_id.as_deref() == Some(owner_session_id))
1919 .ok_or_else(|| anyhow::anyhow!("Trigger '{trigger_id}' not found"))
1920 }
1921
1922 /// Atomically persist a trigger record.
1923 pub fn save_trigger(&self, record: &DelayedTriggerRecord) -> Result<()> {
1924 self.with_transaction(|| self.save_trigger_unlocked(record))
1925 }
1926
1927 fn save_trigger_unlocked(&self, record: &DelayedTriggerRecord) -> Result<()> {
1928 let path = self.trigger_path(&record.trigger_id)?;
1929 write_json_atomic(&path, record)
1930 }
1931
1932 /// List triggers, newest first. Pass `status_filter` to restrict results.
1933 pub fn list_triggers(
1934 &self,
1935 status_filter: Option<DelayedTriggerStatus>,
1936 limit: Option<usize>,
1937 ) -> Result<Vec<DelayedTriggerRecord>> {
1938 let mut out = Vec::new();
1939 if !self.triggers_dir.exists() {
1940 return Ok(out);
1941 }
1942 for entry in fs::read_dir(&self.triggers_dir)
1943 .with_context(|| format!("Failed to read {}", self.triggers_dir.display()))?
1944 {
1945 let entry = entry?;
1946 let path = entry.path();
1947 if path.extension().is_none_or(|ext| ext != "json") {
1948 continue;
1949 }
1950 match fs::read_to_string(&path)
1951 .ok()
1952 .and_then(|raw| serde_json::from_str::<DelayedTriggerRecord>(&raw).ok())
1953 {
1954 Some(record) => {
1955 if let Some(filter) = status_filter
1956 && record.status != filter
1957 {
1958 continue;
1959 }
1960 out.push(record);
1961 }
1962 None => {
1963 tracing::warn!("Skipping unreadable trigger file {}", path.display());
1964 }
1965 }
1966 }
1967 out.sort_by_key(|r| std::cmp::Reverse(r.created_at));
1968 if let Some(limit) = limit {
1969 out.truncate(limit);
1970 }
1971 Ok(out)
1972 }
1973
1974 /// List session-owned triggers, applying ownership before sorting and limit.
1975 pub fn list_triggers_for_owner(
1976 &self,
1977 status_filter: Option<DelayedTriggerStatus>,
1978 limit: Option<usize>,
1979 owner_session_id: &str,
1980 ) -> Result<Vec<DelayedTriggerRecord>> {
1981 let mut records = self.list_triggers(status_filter, None)?;
1982 records.retain(|record| record.owner_session_id.as_deref() == Some(owner_session_id));
1983 if let Some(limit) = limit {
1984 records.truncate(limit);
1985 }
1986 Ok(records)
1987 }
1988
1989 /// Cancel a pending trigger owned by the given session.
1990 pub fn cancel_trigger_for_owner(
1991 &self,
1992 trigger_id: &str,
1993 owner_session_id: &str,
1994 ) -> Result<DelayedTriggerRecord> {
1995 self.with_transaction(|| {
1996 let mut record = self.get_trigger_for_owner(trigger_id, owner_session_id)?;
1997 if record.status != DelayedTriggerStatus::Pending || record.dispatch.is_some() {
1998 bail!(
1999 "Trigger '{trigger_id}' cannot be canceled after admission (status: {:?})",
2000 record.status
2001 );
2002 }
2003 record.status = DelayedTriggerStatus::Canceled;
2004 self.save_trigger_unlocked(&record)?;
2005 Ok(record)
2006 })
2007 }
2008
2009 /// Return due proposals and unfinished durable trigger admissions.
2010 pub fn collect_due_triggers(&self, now: DateTime<Utc>) -> Result<Vec<DelayedTriggerRecord>> {
2011 Ok(self
2012 .list_triggers(None, None)?
2013 .into_iter()
2014 .filter(|trigger| {
2015 trigger.status == DelayedTriggerStatus::Dispatching
2016 || (trigger.status == DelayedTriggerStatus::Pending
2017 && trigger.owner_session_id.is_some()
2018 && trigger.fire_at <= now)
2019 })
2020 .collect())
2021 }
2022 }
2023
2024 fn new_run_record(
2025 automation_id: &str,
2026 scheduled_for: DateTime<Utc>,
2027 created_at: DateTime<Utc>,
2028 ) -> AutomationRunRecord {
2029 AutomationRunRecord {
2030 schema_version: CURRENT_RUN_SCHEMA_VERSION,
2031 id: Uuid::new_v4().to_string(),
2032 automation_id: automation_id.to_string(),
2033 scheduled_for,
2034 status: AutomationRunStatus::Queued,
2035 created_at,
2036 started_at: None,
2037 ended_at: None,
2038 task_id: None,
2039 thread_id: None,
2040 turn_id: None,
2041 error: None,
2042 dispatch: None,
2043 }
2044 }
2045
2046 /// The posture an automation's `auto_approve` bit stands for, in the wire
2047 /// spelling a task request carries.
2048 ///
2049 /// A scheduled run fires with no session to inherit from, so the only authority
2050 /// it has is its own record — and stating it explicitly keeps the scheduled
2051 /// task from re-deriving a posture out of a legacy bit on every later change to
2052 /// what that bit means.
2053 fn automation_posture(auto_approve: bool) -> String {
2054 crate::runtime_policy::approval_wire(crate::core::authority::posture_from_auto_approve(
2055 auto_approve,
2056 ))
2057 .to_string()
2058 }
2059
2060 fn automation_task_request(automation: &AutomationRecord) -> NewTaskRequest {
2061 NewTaskRequest {
2062 prompt: automation.prompt.clone(),
2063 name: Some(automation.name.clone()),
2064 model: automation.model.clone(),
2065 model_provider: automation.model_provider.clone(),
2066 model_provider_id: automation.model_provider_id.clone(),
2067 workspace: automation.cwds.first().cloned(),
2068 mode: Some(automation.task_mode()),
2069 allow_shell: Some(automation.task_allow_shell()),
2070 trust_mode: Some(automation.task_trust_mode()),
2071 auto_approve: Some(automation.task_auto_approve()),
2072 permission_posture: Some(automation_posture(automation.task_auto_approve())),
2073 owner_session_id: None,
2074 }
2075 }
2076
2077 fn bind_run_dispatch(
2078 run: &mut AutomationRunRecord,
2079 automation: &AutomationRecord,
2080 task_data_dir: &Path,
2081 scheduled: bool,
2082 ) -> Result<()> {
2083 if automation.execution_scope.is_none() {
2084 bail!("Automation execution ownership is unverified");
2085 }
2086 run.schema_version = CURRENT_RUN_SCHEMA_VERSION;
2087 run.task_id = Some(crate::task_manager::TaskManager::new_task_id());
2088 run.dispatch = Some(AutomationDispatch {
2089 execution_scope: automation.execution_scope.clone(),
2090 request: automation_task_request(automation),
2091 task_data_dir: task_data_dir
2092 .canonicalize()
2093 .context("resolve task store before automation admission")?,
2094 accepted: false,
2095 delivery_mode: automation.delivery_mode(),
2096 suppress_report: false,
2097 schedule: scheduled.then(|| AdmittedSchedule {
2098 updated_at: automation.updated_at,
2099 rrule: automation.rrule.clone(),
2100 }),
2101 });
2102 Ok(())
2103 }
2104
2105 fn check_dispatch_store(dispatch: &AutomationDispatch, tasks: &SharedTaskManager) -> Result<()> {
2106 if tasks.data_dir().canonicalize()? != dispatch.task_data_dir.canonicalize()? {
2107 bail!("Automation admission belongs to a different task store; it cannot be replayed here");
2108 }
2109 Ok(())
2110 }
2111
2112 async fn dispatch_bound_task(
2113 dispatch: &mut AutomationDispatch,
2114 task_id: &str,
2115 tasks: &SharedTaskManager,
2116 ) -> Result<crate::task_manager::TaskRecord> {
2117 check_dispatch_store(dispatch, tasks)?;
2118 if dispatch.execution_scope.as_deref() != Some(tasks.execution_scope()) {
2119 bail!(
2120 "Automation admission execution ownership is unverified or belongs to another Runtime"
2121 );
2122 }
2123 let task = if dispatch.accepted {
2124 tasks
2125 .read_bound_task(task_id)?
2126 .context("Accepted automation task is missing; refusing to replay it")?
2127 } else {
2128 tasks
2129 .recover_task_admission(dispatch.request.clone(), task_id.to_owned())
2130 .await?
2131 };
2132 crate::task_manager::validate_bound_task_request(&task, &dispatch.request)?;
2133 dispatch.accepted = true;
2134 dispatch.suppress_report = dispatch.delivery_mode == AutomationDeliveryMode::Watcher
2135 && task.status == TaskStatus::Completed
2136 && task
2137 .result_summary
2138 .as_deref()
2139 .is_some_and(|summary| summary.trim() == AUTOMATION_WATCHER_NO_REPORT_SENTINEL);
2140 Ok(task)
2141 }
2142
2143 /// Caller owns dispatch.lock and has already persisted this exact binding.
2144 async fn enqueue_run_task(run: &mut AutomationRunRecord, tasks: &SharedTaskManager) {
2145 let result = match (&mut run.dispatch, &run.task_id) {
2146 (Some(dispatch), Some(task_id)) => dispatch_bound_task(dispatch, task_id, tasks).await,
2147 _ => Err(anyhow::anyhow!(
2148 "Automation run has no durable task binding"
2149 )),
2150 };
2151 match result {
2152 Ok(task) => {
2153 run.error = None;
2154 apply_task_status(run, &task);
2155 }
2156 Err(error) => {
2157 // Keep the same pending identity after uncertain admission. A later
2158 // tick first looks for its canonical task; no new id is allocated.
2159 run.error = Some(format!(
2160 "Automation task admission needs recovery: {error:#}"
2161 ));
2162 }
2163 }
2164 }
2165
2166 fn dispatch_lock_busy(error: &std::io::Error) -> bool {
2167 error.kind() == std::io::ErrorKind::WouldBlock || matches!(error.raw_os_error(), Some(32 | 33))
2168 }
2169
2170 pub async fn run_now_shared(
2171 automations: &SharedAutomationManager,
2172 automation_id: &str,
2173 task_manager: &SharedTaskManager,
2174 ) -> Result<AutomationRunRecord> {
2175 automations.lock().await.bind_task_manager(task_manager)?;
2176 let task_manager = Arc::clone(task_manager);
2177 let task_data_dir = task_manager.data_dir();
2178 run_now_with(
2179 automations,
2180 automation_id,
2181 &task_data_dir,
2182 move |_, mut run| async move {
2183 enqueue_run_task(&mut run, &task_manager).await;
2184 run
2185 },
2186 )
2187 .await
2188 }
2189
2190 /// Keep the process-wide dispatch claim over the await, never the manager mutex.
2191 async fn run_now_with<F, Fut>(
2192 automations: &SharedAutomationManager,
2193 automation_id: &str,
2194 task_data_dir: &Path,
2195 enqueue: F,
2196 ) -> Result<AutomationRunRecord>
2197 where
2198 F: FnOnce(AutomationRecord, AutomationRunRecord) -> Fut,
2199 Fut: Future<Output = AutomationRunRecord>,
2200 {
2201 let mut lock = automations.lock().await.open_lock("dispatch.lock")?;
2202 let _dispatch = lock
2203 .try_write()
2204 .context("Automation dispatcher is busy; retry the run")?;
2205 let (automation, run) = {
2206 let manager = automations.lock().await;
2207 manager.with_transaction(|| {
2208 let mut automation = manager.get_automation(automation_id)?;
2209 manager.adopt_for_run(&mut automation)?;
2210 let now = Utc::now();
2211 let mut run = new_run_record(&automation.id, now, now);
2212 bind_run_dispatch(&mut run, &automation, task_data_dir, false)?;
2213 manager.save_run(&run)?;
2214 Ok((automation, run))
2215 })?
2216 };
2217 let run = enqueue(automation, run).await;
2218 automations
2219 .lock()
2220 .await
2221 .finish_scheduled_run(&run, Utc::now())?;
2222 Ok(run)
2223 }
2224
2225 /// Run a scheduler store operation on the blocking pool (#6149), holding the
2226 /// manager lock exactly as the inline call did: the whole-store scans (every
2227 /// automation or run record read and parsed) and the per-record claim,
2228 /// recovery, and receipt transactions, which wait on the cross-process
2229 /// `state.lock` and read/write JSON. Known limit: opening `dispatch.lock`
2230 /// (one small file open whose non-blocking guard must outlive the awaits)
2231 /// and the in-memory scope checks stay inline.
2232 async fn with_manager_blocking<T, F>(automations: &SharedAutomationManager, scan: F) -> Result<T>
2233 where
2234 T: Send + 'static,
2235 F: FnOnce(&AutomationManager) -> Result<T> + Send + 'static,
2236 {
2237 let manager = Arc::clone(automations).lock_owned().await;
2238 tokio::task::spawn_blocking(move || scan(&manager))
2239 .await
2240 .context("automation store scan task failed")?
2241 }
2242
2243 async fn scheduler_tick_shared(
2244 automations: &SharedAutomationManager,
2245 task_manager: &SharedTaskManager,
2246 ) -> Result<()> {
2247 automations.lock().await.bind_task_manager(task_manager)?;
2248 let tasks = Arc::clone(task_manager);
2249 scheduler_tick_with(automations, &tasks.data_dir(), move |mut run| {
2250 let tasks = Arc::clone(&tasks);
2251 async move {
2252 enqueue_run_task(&mut run, &tasks).await;
2253 run
2254 }
2255 })
2256 .await
2257 }
2258
2259 async fn scheduler_tick_with<F, Fut>(
2260 automations: &SharedAutomationManager,
2261 task_data_dir: &Path,
2262 mut enqueue: F,
2263 ) -> Result<()>
2264 where
2265 F: FnMut(AutomationRunRecord) -> Fut,
2266 Fut: Future<Output = AutomationRunRecord>,
2267 {
2268 let mut lock = automations.lock().await.open_lock("dispatch.lock")?;
2269 let _dispatch = match lock.try_write() {
2270 Ok(guard) => guard,
2271 Err(error) if dispatch_lock_busy(&error) => return Ok(()),
2272 Err(error) => return Err(error).context("claim automation dispatch"),
2273 };
2274 // Repair admitted work before collecting a new occurrence, including claims
2275 // whose definitions were edited/deleted or whose final enqueue save tore.
2276 let pending =
2277 with_manager_blocking(automations, AutomationManager::collect_pending_runs).await?;
2278 for run in pending.into_iter().filter(|run| {
2279 run.dispatch
2280 .as_ref()
2281 .is_some_and(|dispatch| !dispatch.accepted)
2282 }) {
2283 if !automations.lock().await.eligible_scope(
2284 run.dispatch
2285 .as_ref()
2286 .and_then(|dispatch| dispatch.execution_scope.as_deref()),
2287 ) {
2288 continue;
2289 }
2290 // A single damaged admission is quarantined to its diagnostic; it must
2291 // not take down recovery of every pending run behind it.
2292 let recovered = with_manager_blocking(automations, {
2293 let run = run.clone();
2294 move |manager| manager.recover_schedule_advance(&run)
2295 })
2296 .await;
2297 if let Err(error) = recovered {
2298 tracing::warn!(
2299 "automation schedule recovery failed for run {}: {error:#}",
2300 run.id
2301 );
2302 continue;
2303 }
2304 let run = enqueue(run).await;
2305 if let Err(error) = with_manager_blocking(automations, {
2306 let run = run.clone();
2307 move |manager| manager.finish_scheduled_run(&run, Utc::now())
2308 })
2309 .await
2310 {
2311 tracing::warn!(
2312 "automation run {} receipt could not be persisted: {error:#}",
2313 run.id
2314 );
2315 }
2316 }
2317 let now = Utc::now();
2318 let due =
2319 with_manager_blocking(automations, move |manager| manager.collect_due_runs(now)).await?;
2320 for (observed, proposed) in due {
2321 if !automations
2322 .lock()
2323 .await
2324 .eligible_scope(observed.execution_scope.as_deref())
2325 {
2326 continue;
2327 }
2328 let claimed = with_manager_blocking(automations, {
2329 let observed = observed.clone();
2330 let task_data_dir = task_data_dir.to_path_buf();
2331 move |manager| manager.claim_scheduled_run(&observed, proposed, &task_data_dir)
2332 })
2333 .await;
2334 let run = match claimed {
2335 Ok(run) => run,
2336 Err(error) => {
2337 // One automation's claim failure (for example a corrupt
2338 // receipt in its own dedup history) quarantines that
2339 // automation, not the tick: later due work still dispatches.
2340 tracing::warn!(
2341 "automation {} occurrence claim failed: {error:#}",
2342 observed.id
2343 );
2344 continue;
2345 }
2346 };
2347 let Some(run) = run else {
2348 continue;
2349 };
2350 let run = enqueue(run).await;
2351 if let Err(error) = with_manager_blocking(automations, {
2352 let run = run.clone();
2353 move |manager| manager.finish_scheduled_run(&run, now)
2354 })
2355 .await
2356 {
2357 tracing::warn!(
2358 "automation run {} receipt could not be persisted: {error:#}",
2359 run.id
2360 );
2361 }
2362 }
2363 Ok(())
2364 }
2365
2366 async fn fire_due_triggers_shared(
2367 automations: &SharedAutomationManager,
2368 task_manager: &SharedTaskManager,
2369 ) -> Result<()> {
2370 automations.lock().await.bind_task_manager(task_manager)?;
2371 let tasks = Arc::clone(task_manager);
2372 fire_due_triggers_with(automations, &tasks.data_dir(), move |trigger| {
2373 let tasks = Arc::clone(&tasks);
2374 async move { enqueue_trigger_task(trigger, &tasks).await }
2375 })
2376 .await
2377 }
2378
2379 async fn enqueue_trigger_task(
2380 mut trigger: DelayedTriggerRecord,
2381 tasks: &SharedTaskManager,
2382 ) -> Result<DelayedTriggerRecord> {
2383 let result = dispatch_bound_task(
2384 trigger.dispatch.as_mut().context("trigger dispatch")?,
2385 trigger.task_id.as_deref().context("trigger task binding")?,
2386 tasks,
2387 )
2388 .await;
2389 match result {
2390 Ok(task) => {
2391 trigger.status = DelayedTriggerStatus::Fired;
2392 trigger.fired_at = Some(Utc::now());
2393 trigger.thread_id = task.thread_id.clone();
2394 trigger.error = None;
2395 }
2396 Err(error) => {
2397 trigger.error = Some(format!(
2398 "Delayed trigger admission needs recovery: {error:#}"
2399 ))
2400 }
2401 }
2402 Ok(trigger)
2403 }
2404
2405 async fn fire_due_triggers_with<F, Fut>(
2406 automations: &SharedAutomationManager,
2407 task_data_dir: &Path,
2408 mut enqueue: F,
2409 ) -> Result<()>
2410 where
2411 F: FnMut(DelayedTriggerRecord) -> Fut,
2412 Fut: Future<Output = Result<DelayedTriggerRecord>>,
2413 {
2414 let mut lock = automations.lock().await.open_lock("dispatch.lock")?;
2415 let _dispatch = match lock.try_write() {
2416 Ok(guard) => guard,
2417 Err(error) if dispatch_lock_busy(&error) => return Ok(()),
2418 Err(error) => return Err(error).context("claim delayed-trigger dispatch"),
2419 };
2420 let now = Utc::now();
2421 let candidates = with_manager_blocking(automations, move |manager| {
2422 manager.collect_due_triggers(now)
2423 })
2424 .await?;
2425 for candidate in candidates {
2426 let scope = if candidate.status == DelayedTriggerStatus::Dispatching {
2427 candidate
2428 .dispatch
2429 .as_ref()
2430 .and_then(|d| d.execution_scope.as_deref())
2431 } else {
2432 candidate.execution_scope.as_deref()
2433 };
2434 if !automations.lock().await.eligible_scope(scope) {
2435 continue;
2436 }
2437 if !(candidate.status == DelayedTriggerStatus::Dispatching
2438 || (candidate.status == DelayedTriggerStatus::Pending && candidate.fire_at <= now))
2439 {
2440 continue;
2441 }
2442 let claim = with_manager_blocking(automations, {
2443 let trigger_id = candidate.trigger_id.clone();
2444 let task_data_dir = task_data_dir.to_path_buf();
2445 move |manager| {
2446 manager.with_transaction(|| {
2447 let mut current = manager.get_trigger(&trigger_id)?;
2448 if current.status == DelayedTriggerStatus::Dispatching {
2449 if !manager.eligible_scope(
2450 current
2451 .dispatch
2452 .as_ref()
2453 .and_then(|d| d.execution_scope.as_deref()),
2454 ) {
2455 return Ok(None);
2456 }
2457 if current.dispatch.is_none() || current.task_id.is_none() {
2458 bail!("Claimed delayed trigger has no durable task binding");
2459 }
2460 return Ok(Some(current));
2461 }
2462 if !manager.eligible_scope(current.execution_scope.as_deref())
2463 || current.status != DelayedTriggerStatus::Pending
2464 || current.fire_at > now
2465 || current.owner_session_id.is_none()
2466 {
2467 return Ok(None);
2468 }
2469 current.schema_version = CURRENT_TRIGGER_SCHEMA_VERSION;
2470 current.status = DelayedTriggerStatus::Dispatching;
2471 current.task_id = Some(crate::task_manager::TaskManager::new_task_id());
2472 current.dispatch = Some(AutomationDispatch {
2473 execution_scope: current.execution_scope.clone(),
2474 request: NewTaskRequest {
2475 prompt: current.message.clone(),
2476 name: None,
2477 model: None,
2478 model_provider: None,
2479 model_provider_id: None,
2480 workspace: current.workspace.clone(),
2481 mode: Some("agent".into()),
2482 allow_shell: Some(false),
2483 trust_mode: Some(false),
2484 auto_approve: Some(false),
2485 permission_posture: Some(automation_posture(false)),
2486 owner_session_id: current.owner_session_id.clone(),
2487 },
2488 task_data_dir: task_data_dir.canonicalize()?,
2489 accepted: false,
2490 delivery_mode: AutomationDeliveryMode::Task,
2491 suppress_report: false,
2492 schedule: None,
2493 });
2494 manager.save_trigger_unlocked(&current)?;
2495 Ok(Some(current))
2496 })
2497 }
2498 })
2499 .await;
2500 let claimed = match claim {
2501 Ok(claimed) => claimed,
2502 Err(error) => {
2503 // One damaged trigger record quarantines to a diagnostic;
2504 // the remaining due triggers still fire this pass.
2505 tracing::warn!(
2506 "delayed trigger {} claim failed: {error:#}",
2507 candidate.trigger_id
2508 );
2509 continue;
2510 }
2511 };
2512 let Some(trigger) = claimed else {
2513 continue;
2514 };
2515 let trigger = match enqueue(trigger).await {
2516 Ok(trigger) => trigger,
2517 Err(error) => {
2518 tracing::warn!(
2519 "delayed trigger {} enqueue failed: {error:#}",
2520 candidate.trigger_id
2521 );
2522 continue;
2523 }
2524 };
2525 if let Err(error) = with_manager_blocking(automations, {
2526 let trigger = trigger.clone();
2527 move |manager| manager.save_trigger(&trigger)
2528 })
2529 .await
2530 {
2531 tracing::warn!(
2532 "delayed trigger {} receipt could not be persisted: {error:#}",
2533 trigger.trigger_id
2534 );
2535 }
2536 }
2537 Ok(())
2538 }
2539
2540 /// Fold a durable task's state back into its automation run. Returns whether
2541 /// the run changed and needs persisting.
2542 fn apply_task_status(
2543 run: &mut AutomationRunRecord,
2544 task: &crate::task_manager::TaskRecord,
2545 ) -> bool {
2546 let mut changed = run.thread_id != task.thread_id || run.turn_id != task.turn_id;
2547 run.thread_id = task.thread_id.clone();
2548 run.turn_id = task.turn_id.clone();
2549 match task.status {
2550 TaskStatus::Queued => {
2551 if !matches!(run.status, AutomationRunStatus::Queued) {
2552 run.status = AutomationRunStatus::Queued;
2553 changed = true;
2554 }
2555 }
2556 TaskStatus::Running => {
2557 if !matches!(run.status, AutomationRunStatus::Running) {
2558 run.status = AutomationRunStatus::Running;
2559 changed = true;
2560 }
2561 if run.started_at.is_none() {
2562 run.started_at = Some(task.started_at.unwrap_or_else(Utc::now));
2563 changed = true;
2564 }
2565 }
2566 TaskStatus::Completed => {
2567 run.status = AutomationRunStatus::Completed;
2568 run.started_at = run.started_at.or(task.started_at);
2569 run.ended_at = task.ended_at.or(Some(Utc::now()));
2570 run.error = None;
2571 changed = true;
2572 }
2573 TaskStatus::Failed => {
2574 run.status = AutomationRunStatus::Failed;
2575 run.started_at = run.started_at.or(task.started_at);
2576 run.ended_at = task.ended_at.or(Some(Utc::now()));
2577 run.error = task.error.clone();
2578 changed = true;
2579 }
2580 TaskStatus::Canceled => {
2581 run.status = AutomationRunStatus::Canceled;
2582 run.started_at = run.started_at.or(task.started_at);
2583 run.ended_at = task.ended_at.or(Some(Utc::now()));
2584 // #6162: a cancellation is not silent. Keep the task's own error
2585 // when it recorded one, otherwise name the terminal reason so the
2586 // settled receipt can say who or what canceled the run.
2587 run.error = task.error.clone().or_else(|| {
2588 task.terminal_reason
2589 .as_deref()
2590 .map(cancellation_reason_text)
2591 });
2592 changed = true;
2593 }
2594 }
2595 changed
2596 }
2597
2598 /// Human-readable cancellation detail for a run whose task ended without an
2599 /// error of its own. The task manager's terminal reasons are stable strings
2600 /// (`TaskTerminalReason::as_str`); anything unknown is passed through.
2601 fn cancellation_reason_text(terminal_reason: &str) -> String {
2602 // A cooperative cancel is the one path that arrives without an error of
2603 // its own; cancel-timeout and shutdown already carry the task manager's
2604 // receipt message, so those arms are a fallback for records that lost it.
2605 match terminal_reason {
2606 "canceled" => "canceled by request".to_string(),
2607 "cancel_timeout" => "canceled; the task did not stop within the cancel timeout".to_string(),
2608 "shutdown" => "canceled by shutdown".to_string(),
2609 other => format!("canceled ({other})"),
2610 }
2611 }
2612
2613 async fn reconcile_run_statuses_shared(
2614 automations: &SharedAutomationManager,
2615 task_manager: &SharedTaskManager,
2616 ) -> Result<()> {
2617 automations.lock().await.bind_task_manager(task_manager)?;
2618 let mut lock = automations.lock().await.open_lock("dispatch.lock")?;
2619 let _dispatch = match lock.try_write() {
2620 Ok(guard) => guard,
2621 Err(error) if dispatch_lock_busy(&error) => return Ok(()),
2622 Err(error) => return Err(error).context("claim automation reconciliation"),
2623 };
2624 let pending =
2625 with_manager_blocking(automations, AutomationManager::collect_pending_runs).await?;
2626 for mut run in pending {
2627 // Shared storage is not shared ownership. A receipt admitted under
2628 // another execution scope is reconciled by that scope's owner; this
2629 // process must not stamp errors onto it or rewrite it from a bound
2630 // task record it cannot see.
2631 if run
2632 .dispatch
2633 .as_ref()
2634 .and_then(|dispatch| dispatch.execution_scope.as_deref())
2635 != Some(task_manager.execution_scope())
2636 {
2637 continue;
2638 }
2639 let Some(task_id) = run.task_id.clone() else {
2640 continue;
2641 };
2642 // The bound task record is a JSON read under the task store: keep it
2643 // off the Tokio worker like the scans above (#6149).
2644 let lookup = tokio::task::spawn_blocking({
2645 let dispatch = run.dispatch.clone();
2646 let task_id = task_id.clone();
2647 let task_manager = Arc::clone(task_manager);
2648 move || {
2649 if let Some(dispatch) = &dispatch {
2650 check_dispatch_store(dispatch, &task_manager)?;
2651 }
2652 let task = task_manager.read_bound_task(&task_id)?;
2653 if let Some(task) = &task
2654 && let Some(dispatch) = &dispatch
2655 {
2656 crate::task_manager::validate_bound_task_request(task, &dispatch.request)?;
2657 }
2658 Ok::<_, anyhow::Error>(task)
2659 }
2660 })
2661 .await
2662 .context("automation task lookup failed")
2663 .and_then(|lookup| lookup);
2664 let task = match lookup {
2665 Ok(Some(task)) => task,
2666 Ok(None) => {
2667 if run
2668 .dispatch
2669 .as_ref()
2670 .is_some_and(|dispatch| dispatch.accepted)
2671 {
2672 // The admission was durably accepted but the bound task
2673 // record is gone: this occurrence can never be replayed
2674 // or reconciled. Settle it terminally instead of retrying
2675 // the same lookup every pass and starving the schedule
2676 // behind it.
2677 run.status = AutomationRunStatus::Failed;
2678 run.ended_at = Some(run.ended_at.unwrap_or_else(Utc::now));
2679 run.error = Some(format!(
2680 "Bound automation task {task_id} is missing after durable acceptance"
2681 ));
2682 } else {
2683 // Unaccepted admissions are the scheduler's recovery path:
2684 // the next tick reuses the durable binding or recreates
2685 // the task; reconcile only records the uncertainty.
2686 run.error = Some(format!(
2687 "Automation reconciliation unavailable: bound task {task_id} is missing"
2688 ));
2689 }
2690 if let Err(error) = with_manager_blocking(automations, {
2691 let run = run.clone();
2692 move |manager| manager.finish_scheduled_run(&run, Utc::now())
2693 })
2694 .await
2695 {
2696 tracing::warn!(
2697 "automation run {} receipt could not be persisted: {error:#}",
2698 run.id
2699 );
2700 }
2701 continue;
2702 }
2703 Err(error) => {
2704 run.error = Some(format!("Automation reconciliation unavailable: {error:#}"));
2705 if let Err(error) = with_manager_blocking(automations, {
2706 let run = run.clone();
2707 move |manager| manager.finish_scheduled_run(&run, Utc::now())
2708 })
2709 .await
2710 {
2711 tracing::warn!(
2712 "automation run {} receipt could not be persisted: {error:#}",
2713 run.id
2714 );
2715 }
2716 continue;
2717 }
2718 };
2719 // The bound task itself belongs to another Runtime — hands off.
2720 if task.execution_scope.as_deref() != Some(task_manager.execution_scope()) {
2721 continue;
2722 }
2723 let watcher = run
2724 .dispatch
2725 .as_ref()
2726 .map(|dispatch| dispatch.delivery_mode)
2727 .or_else(|| {
2728 automations.try_lock().ok().and_then(|manager| {
2729 manager
2730 .get_automation(&run.automation_id)
2731 .ok()
2732 .map(|automation| automation.delivery_mode())
2733 })
2734 })
2735 == Some(AutomationDeliveryMode::Watcher);
2736 let watcher_noop = watcher
2737 && task.status == TaskStatus::Completed
2738 && task
2739 .result_summary
2740 .as_deref()
2741 .is_some_and(|summary| summary.trim() == AUTOMATION_WATCHER_NO_REPORT_SENTINEL);
2742 if !apply_task_status(&mut run, &task) {
2743 continue;
2744 }
2745 let dispatch = run.dispatch.get_or_insert_with(|| AutomationDispatch {
2746 execution_scope: task.execution_scope.clone(),
2747 request: NewTaskRequest::from_task(&task),
2748 task_data_dir: task_manager.data_dir(),
2749 accepted: true,
2750 delivery_mode: if watcher {
2751 AutomationDeliveryMode::Watcher
2752 } else {
2753 AutomationDeliveryMode::Task
2754 },
2755 suppress_report: false,
2756 schedule: None,
2757 });
2758 run.schema_version = CURRENT_RUN_SCHEMA_VERSION;
2759 dispatch.accepted = true;
2760 dispatch.suppress_report = watcher_noop;
2761 if let Err(error) = with_manager_blocking(automations, {
2762 let run = run.clone();
2763 move |manager| manager.finish_scheduled_run(&run, Utc::now())
2764 })
2765 .await
2766 {
2767 tracing::warn!(
2768 "automation run {} receipt could not be persisted: {error:#}",
2769 run.id
2770 );
2771 }
2772 }
2773 Ok(())
2774 }
2775
2776 /// Fixed-width, lexically-sortable UTC stamp for run file names, e.g.
2777 /// `20260705T142530123Z` (millisecond precision; the run id suffix breaks
2778 /// same-millisecond ties deterministically).
2779 const RUN_STAMP_FORMAT: &str = "%Y%m%dT%H%M%S%3fZ";
2780 const RUN_STAMP_LEN: usize = "20260705T142530123Z".len();
2781
2782 fn run_file_stamp(created_at: DateTime<Utc>) -> String {
2783 created_at.format(RUN_STAMP_FORMAT).to_string()
2784 }
2785
2786 /// Shape check for `{stamp}-{run_id}` file stems. Ordering trusts the file
2787 /// name only for pruning; the parsed record's `created_at` stays
2788 /// authoritative for the final sort.
2789 fn has_sortable_run_stem(stem: &str) -> bool {
2790 let Some((stamp, rest)) = stem.split_at_checked(RUN_STAMP_LEN) else {
2791 return false;
2792 };
2793 if !rest.starts_with('-') || rest.len() < 2 {
2794 return false;
2795 }
2796 stamp.char_indices().all(|(idx, ch)| match idx {
2797 8 => ch == 'T',
2798 18 => ch == 'Z',
2799 _ => ch.is_ascii_digit(),
2800 })
2801 }
2802
2803 fn read_automation_file(path: &Path) -> Result<AutomationRecord> {
2804 let raw = fs::read_to_string(path)
2805 .with_context(|| format!("Failed to read automation {}", path.display()))?;
2806 let record: AutomationRecord = serde_json::from_str(&raw)
2807 .with_context(|| format!("Failed to parse automation {}", path.display()))?;
2808 if record.schema_version > CURRENT_AUTOMATION_SCHEMA_VERSION {
2809 bail!(
2810 "Automation schema v{} is newer than supported v{}",
2811 record.schema_version,
2812 CURRENT_AUTOMATION_SCHEMA_VERSION
2813 );
2814 }
2815 Ok(record)
2816 }
2817
2818 fn read_run_file(path: &Path) -> Result<AutomationRunRecord> {
2819 let raw =
2820 fs::read_to_string(path).with_context(|| format!("Failed to read {}", path.display()))?;
2821 let run: AutomationRunRecord = serde_json::from_str(&raw)
2822 .with_context(|| format!("Failed to parse {}", path.display()))?;
2823 if run.schema_version > CURRENT_RUN_SCHEMA_VERSION {
2824 bail!(
2825 "Automation run schema v{} is newer than supported v{}",
2826 run.schema_version,
2827 CURRENT_RUN_SCHEMA_VERSION
2828 );
2829 }
2830 Ok(run)
2831 }
2832
2833 fn ensure_safe_storage_id(kind: &str, value: &str) -> Result<()> {
2834 let mut components = Path::new(value).components();
2835 let Some(component) = components.next() else {
2836 bail!("{kind} must not be empty");
2837 };
2838 if components.next().is_some() || !matches!(component, std::path::Component::Normal(_)) {
2839 bail!("{kind} must be a single path component");
2840 }
2841 Ok(())
2842 }
2843
2844 fn validate_name_and_prompt(name: &str, prompt: &str) -> Result<()> {
2845 if name.trim().is_empty() {
2846 bail!("Automation name is required");
2847 }
2848 if prompt.trim().is_empty() {
2849 bail!("Automation prompt is required");
2850 }
2851 Ok(())
2852 }
2853
2854 fn normalize_optional_string(value: Option<String>) -> Option<String> {
2855 value
2856 .map(|value| value.trim().to_string())
2857 .filter(|value| !value.is_empty())
2858 }
2859
2860 fn write_json_atomic<T: Serialize>(path: &Path, value: &T) -> Result<()> {
2861 if let Some(parent) = path.parent() {
2862 fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
2863 }
2864 crate::utils::write_atomic(path, &serde_json::to_vec_pretty(value)?)
2865 .with_context(|| format!("write {}", path.display()))
2866 }
2867
2868 pub fn default_automations_dir() -> PathBuf {
2869 // Most-specific override: an explicit automations dir.
2870 for var in ["CODEWHALE_AUTOMATIONS_DIR", "DEEPSEEK_AUTOMATIONS_DIR"] {
2871 if let Ok(path) = std::env::var(var) {
2872 let trimmed = path.trim();
2873 if !trimmed.is_empty() {
2874 return PathBuf::from(trimmed);
2875 }
2876 }
2877 }
2878 // $CODEWHALE_HOME is a hard override of the base data directory
2879 // (docs/CONFIGURATION.md): when SET, automations live under it and we do
2880 // NOT fall back to the legacy ~/.deepseek path — silent fallback would
2881 // defeat the isolation the override promises. Check the env var directly
2882 // (not codewhale_home()'s Ok/Err, which succeeds for the default home too).
2883 if let Some(home) = codewhale_paths::codewhale_home_override().ok().flatten() {
2884 return home.join("automations");
2885 }
2886 codewhale_paths::user_home()
2887 .map(|home| {
2888 let primary = home.join(".codewhale").join("automations");
2889 let legacy = home.join(".deepseek").join("automations");
2890 if primary.exists() || !legacy.exists() {
2891 return primary;
2892 }
2893 legacy
2894 })
2895 .unwrap_or_else(|| PathBuf::from(".codewhale").join("automations"))
2896 }
2897
2898 pub type SharedAutomationManager = Arc<Mutex<AutomationManager>>;
2899
2900 #[derive(Debug, Clone)]
2901 pub struct AutomationSchedulerConfig {
2902 pub tick_interval_secs: u64,
2903 }
2904
2905 impl Default for AutomationSchedulerConfig {
2906 fn default() -> Self {
2907 Self {
2908 tick_interval_secs: 15,
2909 }
2910 }
2911 }
2912
2913 pub fn spawn_scheduler(
2914 automations: SharedAutomationManager,
2915 task_manager: SharedTaskManager,
2916 cancel: CancellationToken,
2917 config: AutomationSchedulerConfig,
2918 ) -> tokio::task::JoinHandle<()> {
2919 spawn_supervised(
2920 "automation-scheduler",
2921 std::panic::Location::caller(),
2922 async move {
2923 let interval = config.tick_interval_secs.max(5);
2924 loop {
2925 if cancel.is_cancelled() {
2926 break;
2927 }
2928
2929 // Lock scope lives inside the shared helpers: the manager
2930 // mutex is dropped across every task-manager await so API and
2931 // tool callers are never queued behind enqueue/status latency.
2932 if let Err(err) = scheduler_tick_shared(&automations, &task_manager).await {
2933 tracing::warn!("automation scheduler tick failed: {err}");
2934 }
2935 if let Err(err) = reconcile_run_statuses_shared(&automations, &task_manager).await {
2936 tracing::warn!("automation reconcile failed: {err}");
2937 }
2938 if let Err(err) = fire_due_triggers_shared(&automations, &task_manager).await {
2939 tracing::warn!("delayed trigger tick failed: {err}");
2940 }
2941
2942 tokio::select! {
2943 _ = cancel.cancelled() => break,
2944 _ = sleep(std::time::Duration::from_secs(interval)) => {}
2945 }
2946 }
2947 },
2948 )
2949 }
2950
2951 #[cfg(test)]
2952 mod tests {
2953 use super::*;
2954 use async_trait::async_trait;
2955 use chrono::{FixedOffset, LocalResult, NaiveDate};
2956 use tokio::sync::mpsc;
2957
2958 use crate::task_manager::{
2959 ExecutionTask, TaskExecutionEvent, TaskExecutionResult, TaskExecutor, TaskManager,
2960 TaskManagerConfig,
2961 };
2962
2963 struct AutomationRecordingExecutor(PathBuf);
2964
2965 #[async_trait]
2966 impl TaskExecutor for AutomationRecordingExecutor {
2967 async fn execute(
2968 &self,
2969 task: ExecutionTask,
2970 _events: mpsc::Sender<TaskExecutionEvent>,
2971 _cancel: CancellationToken,
2972 ) -> TaskExecutionResult {
2973 use std::io::Write as _;
2974 let recorded = (|| -> Result<()> {
2975 let id = task
2976 .thread_request()
2977 .task_id
2978 .context("fixture task identity")?;
2979 let mut file = fs::OpenOptions::new()
2980 .create(true)
2981 .append(true)
2982 .open(&self.0)?;
2983 writeln!(file, "{id}")?;
2984 file.sync_all()?;
2985 Ok(())
2986 })();
2987 match recorded {
2988 Ok(()) => TaskExecutionResult {
2989 status: TaskStatus::Completed,
2990 result_text: Some("automation fixture completed".into()),
2991 error: None,
2992 terminal_reason: crate::task_manager::TaskTerminalReason::Completed,
2993 },
2994 Err(error) => TaskExecutionResult {
2995 status: TaskStatus::Failed,
2996 result_text: None,
2997 error: Some(error.to_string()),
2998 terminal_reason: crate::task_manager::TaskTerminalReason::Failed,
2999 },
3000 }
3001 }
3002 }
3003
3004 fn fixture_executions(path: &Path) -> Vec<String> {
3005 match fs::read_to_string(path) {
3006 Ok(text) => text.lines().map(str::to_owned).collect(),
3007 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Vec::new(),
3008 Err(error) => panic!("read independent execution receipts: {error}"),
3009 }
3010 }
3011
3012 async fn fixture_tasks(root: &Path, receipts: &Path) -> Result<SharedTaskManager> {
3013 TaskManager::start_with_executor(
3014 automation_task_config(root.to_path_buf()),
3015 Arc::new(AutomationRecordingExecutor(receipts.to_path_buf())),
3016 )
3017 .await
3018 }
3019
3020 fn fixture_due_automation(
3021 manager: &AutomationManager,
3022 id: &str,
3023 order: i64,
3024 ) -> AutomationRecord {
3025 let mut automation = automation_record_with_settings(None, None, None, None);
3026 automation.id = id.to_owned();
3027 automation.prompt = format!("automation fixture {id}");
3028 automation.created_at = Utc::now() - Duration::hours(2);
3029 automation.updated_at = Utc::now() - Duration::minutes(order);
3030 automation.next_run_at = Some(Utc::now() - Duration::minutes(1));
3031 manager
3032 .save_automation(&automation)
3033 .expect("save due fixture");
3034 automation
3035 }
3036
3037 #[tokio::test]
3038 async fn interrupted_dispatch_reuses_binding_and_preserves_replacement_schedule() -> Result<()>
3039 {
3040 for accepted_before_cut in [false, true] {
3041 let root = tempfile::tempdir()?;
3042 let receipts = root.path().join("executions");
3043 let tasks = fixture_tasks(&root.path().join("tasks"), &receipts).await?;
3044 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
3045 let automation = fixture_due_automation(&manager, "recover", 1);
3046 let shared = Arc::new(Mutex::new(manager));
3047 let (at_cut, reached_cut) = tokio::sync::oneshot::channel();
3048 let tick = tokio::spawn({
3049 let shared = shared.clone();
3050 let tasks = tasks.clone();
3051 let mut at_cut = Some(at_cut);
3052 async move {
3053 scheduler_tick_with(&shared, &tasks.data_dir(), move |mut run| {
3054 let tasks = tasks.clone();
3055 let at_cut = at_cut.take();
3056 async move {
3057 if accepted_before_cut {
3058 enqueue_run_task(&mut run, &tasks).await;
3059 }
3060 if let Some(at_cut) = at_cut {
3061 let _ = at_cut.send(run.clone());
3062 }
3063 std::future::pending::<()>().await;
3064 run
3065 }
3066 })
3067 .await
3068 }
3069 });
3070 let observed =
3071 tokio::time::timeout(std::time::Duration::from_secs(5), reached_cut).await??;
3072 let bound_id = observed.task_id.clone().context("claimed task id")?;
3073 let persisted = shared.lock().await.list_runs(&automation.id, None)?;
3074 assert_eq!(
3075 persisted.len(),
3076 1,
3077 "intent exists before either interruption boundary"
3078 );
3079 assert_eq!(persisted[0].task_id.as_deref(), Some(bound_id.as_str()));
3080 assert!(
3081 !persisted[0].dispatch.as_ref().unwrap().accepted,
3082 "final acceptance save has not run"
3083 );
3084 if accepted_before_cut {
3085 let task = crate::task_manager::wait_for_terminal_state(
3086 &tasks,
3087 &bound_id,
3088 std::time::Duration::from_secs(5),
3089 )
3090 .await?;
3091 assert_eq!(task.status, TaskStatus::Completed);
3092 assert!(
3093 task.result_summary
3094 .as_deref()
3095 .unwrap_or_default()
3096 .contains("automation fixture completed")
3097 );
3098 assert_eq!(fixture_executions(&receipts), vec![bound_id.clone()]);
3099 } else {
3100 assert!(tasks.read_bound_task(&bound_id)?.is_none());
3101 assert!(fixture_executions(&receipts).is_empty());
3102 }
3103 // Simulate a dropped request/process before the final run receipt.
3104 tick.abort();
3105 assert!(tick.await.unwrap_err().is_cancelled());
3106 let reopened = AutomationManager::open_for_test(root.path().join("automations"))?;
3107 let future = Utc::now() + Duration::hours(4);
3108 let edited = reopened.update_automation(
3109 &automation.id,
3110 UpdateAutomationRequest {
3111 rrule: Some(format!("FREQ=ONCE;AT={}", future.to_rfc3339())),
3112 prompt: Some("future revised prompt".into()),
3113 ..Default::default()
3114 },
3115 )?;
3116 let restarted = Arc::new(Mutex::new(reopened));
3117 scheduler_tick_shared(&restarted, &tasks).await?;
3118 let task = crate::task_manager::wait_for_terminal_state(
3119 &tasks,
3120 &bound_id,
3121 std::time::Duration::from_secs(5),
3122 )
3123 .await?;
3124 assert_eq!(task.status, TaskStatus::Completed);
3125 assert_eq!(
3126 task.prompt, automation.prompt,
3127 "claim keeps its original request"
3128 );
3129 reconcile_run_statuses_shared(&restarted, &tasks).await?;
3130 scheduler_tick_shared(&restarted, &tasks).await?;
3131 assert_eq!(
3132 fixture_executions(&receipts),
3133 vec![bound_id.clone()],
3134 "same occurrence executes exactly once"
3135 );
3136 let manager = restarted.lock().await;
3137 let runs = manager.list_runs(&automation.id, None)?;
3138 assert_eq!(runs.len(), 1);
3139 assert_eq!(runs[0].task_id.as_deref(), Some(bound_id.as_str()));
3140 assert_eq!(runs[0].status, AutomationRunStatus::Completed);
3141 assert!(runs[0].dispatch.as_ref().unwrap().accepted);
3142 let current = manager.get_automation(&automation.id)?;
3143 assert_eq!(current.rrule, edited.rrule);
3144 assert_eq!(
3145 current.next_run_at, edited.next_run_at,
3146 "old completion must not clear future ONCE"
3147 );
3148 assert_eq!(current.status, AutomationStatus::Active);
3149 assert_eq!(current.updated_at, edited.updated_at);
3150 assert_eq!(current.prompt, edited.prompt);
3151 assert!(current.last_run_at.is_some());
3152 tasks.shutdown();
3153 }
3154 Ok(())
3155 }
3156
3157 #[tokio::test]
3158 async fn pause_and_delete_before_claim_stop_collected_work_but_preserve_admitted_work()
3159 -> Result<()> {
3160 let root = tempfile::tempdir()?;
3161 let receipts = root.path().join("executions");
3162 let tasks = fixture_tasks(&root.path().join("tasks"), &receipts).await?;
3163 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
3164 let first = fixture_due_automation(&manager, "first", 1);
3165 let paused = fixture_due_automation(&manager, "paused", 2);
3166 let deleted = fixture_due_automation(&manager, "deleted", 3);
3167 let shared = Arc::new(Mutex::new(manager));
3168 let (entered, reached) = tokio::sync::oneshot::channel();
3169 let (release, released) = tokio::sync::oneshot::channel();
3170 let tick = tokio::spawn({
3171 let shared = shared.clone();
3172 let tasks = tasks.clone();
3173 let mut barrier = Some((entered, released));
3174 async move {
3175 scheduler_tick_with(&shared, &tasks.data_dir(), move |mut run| {
3176 let barrier = barrier.take();
3177 let tasks = tasks.clone();
3178 async move {
3179 if let Some((entered, released)) = barrier {
3180 let _ = entered.send(run.clone());
3181 let _ = released.await;
3182 }
3183 enqueue_run_task(&mut run, &tasks).await;
3184 run
3185 }
3186 })
3187 .await
3188 }
3189 });
3190 let admitted = tokio::time::timeout(std::time::Duration::from_secs(5), reached).await??;
3191 assert_eq!(
3192 admitted.automation_id, first.id,
3193 "ordered first claim really entered"
3194 );
3195 let other = AutomationManager::open_for_test(root.path().join("automations"))?;
3196 other.pause_automation(&paused.id)?;
3197 other.delete_automation(&deleted.id)?;
3198 // Pause after durable admission affects future runs; this run remains owned.
3199 other.pause_automation(&first.id)?;
3200 release
3201 .send(())
3202 .map_err(|_| anyhow::anyhow!("release fixture"))?;
3203 tokio::time::timeout(std::time::Duration::from_secs(5), tick).await???;
3204 let id = admitted.task_id.context("admitted task identity")?;
3205 let task = crate::task_manager::wait_for_terminal_state(
3206 &tasks,
3207 &id,
3208 std::time::Duration::from_secs(5),
3209 )
3210 .await?;
3211 assert_eq!(task.status, TaskStatus::Completed);
3212 reconcile_run_statuses_shared(&shared, &tasks).await?;
3213 assert_eq!(fixture_executions(&receipts), vec![id]);
3214 assert!(other.list_runs(&paused.id, None)?.is_empty());
3215 assert!(other.list_runs(&deleted.id, None)?.is_empty());
3216 assert_eq!(
3217 other.get_automation(&first.id)?.status,
3218 AutomationStatus::Paused
3219 );
3220 assert!(other.get_automation(&first.id)?.next_run_at.is_none());
3221 tasks.shutdown();
3222 Ok(())
3223 }
3224
3225 #[tokio::test]
3226 async fn reconciliation_reaches_old_runs_and_runs_of_deleted_definitions() -> Result<()> {
3227 for delete_definition in [false, true] {
3228 let root = tempfile::tempdir()?;
3229 let receipts = root.path().join("executions");
3230 let tasks = fixture_tasks(&root.path().join("tasks"), &receipts).await?;
3231 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
3232 let automation = fixture_due_automation(&manager, "history", 1);
3233 let shared = Arc::new(Mutex::new(manager));
3234 let run = run_now_shared(&shared, &automation.id, &tasks).await?;
3235 let id = run.task_id.clone().context("real task id")?;
3236 let task = crate::task_manager::wait_for_terminal_state(
3237 &tasks,
3238 &id,
3239 std::time::Duration::from_secs(5),
3240 )
3241 .await?;
3242 assert_eq!(task.status, TaskStatus::Completed);
3243 {
3244 let manager = shared.lock().await;
3245 for offset in 1..=110 {
3246 let mut newer = new_run_record(
3247 &automation.id,
3248 Utc::now(),
3249 run.created_at + Duration::seconds(offset),
3250 );
3251 newer.status = AutomationRunStatus::Completed;
3252 newer.ended_at = Some(newer.created_at);
3253 manager.save_run(&newer)?;
3254 }
3255 assert!(
3256 manager
3257 .list_runs(&automation.id, Some(100))?
3258 .iter()
3259 .all(|candidate| candidate.id != run.id)
3260 );
3261 if delete_definition {
3262 manager.delete_automation(&automation.id)?;
3263 }
3264 }
3265 reconcile_run_statuses_shared(&shared, &tasks).await?;
3266 let manager = shared.lock().await;
3267 let found =
3268 manager.get_runs_by_ids(&automation.id, &[run.id.clone()].into_iter().collect())?;
3269 assert_eq!(found.len(), 1, "unfinished receipt survives deletion");
3270 assert_eq!(found[0].status, AutomationRunStatus::Completed);
3271 assert_eq!(found[0].task_id.as_deref(), Some(id.as_str()));
3272 assert_eq!(fixture_executions(&receipts), vec![id]);
3273 if delete_definition {
3274 assert!(manager.get_automation(&automation.id).is_err());
3275 } else {
3276 assert!(
3277 manager
3278 .get_automation(&automation.id)?
3279 .last_run_at
3280 .is_some()
3281 );
3282 }
3283 tasks.shutdown();
3284 }
3285 Ok(())
3286 }
3287
3288 #[tokio::test]
3289 async fn foreign_store_and_missing_accepted_task_preserve_binding_without_fallback()
3290 -> Result<()> {
3291 let root = tempfile::tempdir()?;
3292 let receipts = root.path().join("executions");
3293 let tasks = fixture_tasks(&root.path().join("tasks"), &receipts).await?;
3294 let other_receipts = root.path().join("foreign-executions");
3295 let other_tasks =
3296 fixture_tasks(&root.path().join("foreign-tasks"), &other_receipts).await?;
3297 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
3298 let automation = fixture_due_automation(&manager, "bound-store", 1);
3299 let proposed = manager.collect_due_runs(Utc::now())?.remove(0);
3300 let mut run = manager
3301 .claim_scheduled_run(&proposed.0, proposed.1, &tasks.data_dir())?
3302 .context("claim")?;
3303 let id = run.task_id.clone().context("bound task")?;
3304 let shared = Arc::new(Mutex::new(manager));
3305 scheduler_tick_shared(&shared, &other_tasks).await?;
3306 let rows = shared.lock().await.list_runs(&automation.id, None)?;
3307 assert_eq!(rows[0].task_id.as_deref(), Some(id.as_str()));
3308 assert!(
3309 rows[0]
3310 .error
3311 .as_deref()
3312 .unwrap_or_default()
3313 .contains("different task store")
3314 );
3315 assert!(fixture_executions(&other_receipts).is_empty());
3316 enqueue_run_task(&mut run, &tasks).await;
3317 let task = crate::task_manager::wait_for_terminal_state(
3318 &tasks,
3319 &id,
3320 std::time::Duration::from_secs(5),
3321 )
3322 .await?;
3323 assert_eq!(task.status, TaskStatus::Completed);
3324 assert!(run.dispatch.as_ref().unwrap().accepted);
3325 shared.lock().await.finish_scheduled_run(&run, Utc::now())?;
3326 fs::remove_file(tasks.data_dir().join("tasks").join(format!("{id}.json")))?;
3327 scheduler_tick_shared(&shared, &tasks).await?;
3328 reconcile_run_statuses_shared(&shared, &tasks).await?;
3329 let rows = shared.lock().await.list_runs(&automation.id, None)?;
3330 assert_eq!(rows[0].task_id.as_deref(), Some(id.as_str()));
3331 assert!(rows[0].dispatch.as_ref().unwrap().accepted);
3332 assert!(
3333 rows[0]
3334 .error
3335 .as_deref()
3336 .unwrap_or_default()
3337 .contains("missing")
3338 );
3339 assert_eq!(fixture_executions(&receipts), vec![id]);
3340 assert!(fixture_executions(&other_receipts).is_empty());
3341 tasks.shutdown();
3342 other_tasks.shutdown();
3343 Ok(())
3344 }
3345
3346 #[tokio::test]
3347 async fn canceled_collected_trigger_is_not_admitted_after_an_earlier_enqueue_wait() -> Result<()>
3348 {
3349 let root = tempfile::tempdir()?;
3350 let receipts = root.path().join("executions");
3351 let tasks = fixture_tasks(&root.path().join("tasks"), &receipts).await?;
3352 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
3353 let mut triggers = Vec::new();
3354 for index in 0..2 {
3355 let mut trigger = manager.create_trigger(CreateDelayedTriggerRequest {
3356 fire_at: Utc::now() + Duration::hours(1),
3357 message: format!("trigger fixture {index}"),
3358 workspace: None,
3359 owner_session_id: Some("fixture-owner".into()),
3360 parent_trigger_id: None,
3361 })?;
3362 trigger.fire_at = Utc::now() - Duration::minutes(1);
3363 trigger.created_at = Utc::now() - Duration::minutes(index);
3364 manager.save_trigger(&trigger)?;
3365 triggers.push(trigger);
3366 }
3367 let shared = Arc::new(Mutex::new(manager));
3368 let (entered, reached) = tokio::sync::oneshot::channel();
3369 let (release, released) = tokio::sync::oneshot::channel();
3370 let tick = tokio::spawn({
3371 let shared = shared.clone();
3372 let tasks = tasks.clone();
3373 let mut barrier = Some((entered, released));
3374 async move {
3375 fire_due_triggers_with(&shared, &tasks.data_dir(), move |trigger| {
3376 let tasks = tasks.clone();
3377 let barrier = barrier.take();
3378 async move {
3379 if let Some((entered, released)) = barrier {
3380 let _ = entered.send(trigger.clone());
3381 let _ = released.await;
3382 }
3383 enqueue_trigger_task(trigger, &tasks).await
3384 }
3385 })
3386 .await
3387 }
3388 });
3389 let admitted = tokio::time::timeout(std::time::Duration::from_secs(5), reached).await??;
3390 assert_eq!(admitted.trigger_id, triggers[0].trigger_id);
3391 assert_eq!(admitted.status, DelayedTriggerStatus::Dispatching);
3392 let other = AutomationManager::open_for_test(root.path().join("automations"))?;
3393 other.cancel_trigger_for_owner(&triggers[1].trigger_id, "fixture-owner")?;
3394 assert!(
3395 other
3396 .cancel_trigger_for_owner(&admitted.trigger_id, "fixture-owner")
3397 .is_err(),
3398 "claimed work has crossed the admission boundary"
3399 );
3400 release
3401 .send(())
3402 .map_err(|_| anyhow::anyhow!("release trigger fixture"))?;
3403 tokio::time::timeout(std::time::Duration::from_secs(5), tick).await???;
3404 let id = admitted.task_id.context("trigger task identity")?;
3405 let task = crate::task_manager::wait_for_terminal_state(
3406 &tasks,
3407 &id,
3408 std::time::Duration::from_secs(5),
3409 )
3410 .await?;
3411 assert_eq!(task.status, TaskStatus::Completed);
3412 assert_eq!(task.owner_session_id.as_deref(), Some("fixture-owner"));
3413 assert!(!task.allow_shell && !task.trust_mode && !task.auto_approve);
3414 assert_eq!(fixture_executions(&receipts), vec![id]);
3415 assert_eq!(
3416 other.get_trigger(&triggers[0].trigger_id)?.status,
3417 DelayedTriggerStatus::Fired
3418 );
3419 let canceled = other.get_trigger(&triggers[1].trigger_id)?;
3420 assert_eq!(canceled.status, DelayedTriggerStatus::Canceled);
3421 assert!(canceled.task_id.is_none());
3422 tasks.shutdown();
3423 Ok(())
3424 }
3425
3426 async fn wait_fixture_path(path: &Path) -> Result<()> {
3427 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(15);
3428 while !path.try_exists()? {
3429 if std::time::Instant::now() >= deadline {
3430 bail!("fixture barrier timed out: {}", path.display());
3431 }
3432 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
3433 }
3434 Ok(())
3435 }
3436
3437 struct SchedulerFixtureChild(std::process::Child);
3438
3439 impl Drop for SchedulerFixtureChild {
3440 fn drop(&mut self) {
3441 let _ = self.0.kill();
3442 let _ = self.0.wait();
3443 }
3444 }
3445
3446 fn spawn_scheduler_fixture(root: &Path, role: &str) -> Result<SchedulerFixtureChild> {
3447 let log = fs::File::create(root.join(format!("{role}.log")))?;
3448 let home = root.join(format!("{role}-home"));
3449 fs::create_dir_all(&home)?;
3450 Ok(SchedulerFixtureChild(
3451 std::process::Command::new(std::env::current_exe()?)
3452 .args([
3453 "--ignored",
3454 "--exact",
3455 "automation_manager::tests::scheduler_process_fixture",
3456 "--nocapture",
3457 "--test-threads=1",
3458 ])
3459 .env("CW_AUTOMATION_PROCESS_FIXTURE", root)
3460 .env("CW_AUTOMATION_PROCESS_ROLE", role)
3461 .env("CODEWHALE_HOME", &home)
3462 .env("HOME", &home)
3463 .env("USERPROFILE", &home)
3464 .stdin(std::process::Stdio::null())
3465 .stdout(std::process::Stdio::from(log.try_clone()?))
3466 .stderr(std::process::Stdio::from(log))
3467 .spawn()?,
3468 ))
3469 }
3470
3471 async fn finish_scheduler_fixture(child: &mut SchedulerFixtureChild, log: &Path) -> Result<()> {
3472 let deadline = std::time::Instant::now() + std::time::Duration::from_secs(15);
3473 loop {
3474 if let Some(status) = child.0.try_wait()? {
3475 if !status.success() {
3476 bail!("fixture process failed: {}", fs::read_to_string(log)?);
3477 }
3478 return Ok(());
3479 }
3480 if std::time::Instant::now() >= deadline {
3481 bail!("fixture process timed out: {}", log.display());
3482 }
3483 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
3484 }
3485 }
3486
3487 /// Only the parent test launches this entry point with explicit temporary
3488 /// stores. Missing fixture parameters fail; this helper cannot pass by skip.
3489 #[tokio::test]
3490 #[ignore = "subprocess entry point for independent scheduler contention fixture"]
3491 async fn scheduler_process_fixture() -> Result<()> {
3492 let root = PathBuf::from(
3493 std::env::var_os("CW_AUTOMATION_PROCESS_FIXTURE").context("required fixture root")?,
3494 );
3495 let role = std::env::var("CW_AUTOMATION_PROCESS_ROLE")?;
3496 if !matches!(role.as_str(), "first" | "second") {
3497 bail!("invalid fixture role");
3498 }
3499 let scope = if role == "first" {
3500 "test"
3501 } else {
3502 "foreign-scheduler"
3503 };
3504 let tasks = TaskManager::start_with_executor_in_scope(
3505 automation_task_config(root.join("tasks")),
3506 Arc::new(AutomationRecordingExecutor(root.join("executions"))),
3507 scope,
3508 )
3509 .await?;
3510 let mut service = AutomationManager::open(root.join("automations"))?;
3511 service.bind_task_manager(&tasks)?;
3512 let shared = Arc::new(Mutex::new(service));
3513 crate::utils::write_atomic(&root.join(format!("{role}-ready")), b"ready")?;
3514 wait_fixture_path(&root.join(format!("{role}-go"))).await?;
3515 let observations = Arc::new(std::sync::Mutex::new(Vec::new()));
3516 scheduler_tick_with(&shared, &tasks.data_dir(), {
3517 let tasks = tasks.clone();
3518 let root = root.clone();
3519 let role = role.clone();
3520 let observations = observations.clone();
3521 move |mut run| {
3522 let tasks = tasks.clone();
3523 let root = root.clone();
3524 let role = role.clone();
3525 let observations = observations.clone();
3526 async move {
3527 let id = run.task_id.clone().expect("durably claimed fixture id");
3528 observations.lock().unwrap().push(id.clone());
3529 crate::utils::write_atomic(
3530 &root.join(format!("{role}-entered")),
3531 id.as_bytes(),
3532 )
3533 .expect("record dispatch entry");
3534 if role == "first" {
3535 crate::utils::write_atomic(&root.join("first-held"), id.as_bytes())
3536 .expect("record claim barrier");
3537 wait_fixture_path(&root.join("release-first"))
3538 .await
3539 .expect("release owner");
3540 }
3541 enqueue_run_task(&mut run, &tasks).await;
3542 run
3543 }
3544 }
3545 })
3546 .await?;
3547 let ids = observations.lock().unwrap().clone();
3548 for id in &ids {
3549 let task = crate::task_manager::wait_for_terminal_state(
3550 &tasks,
3551 id,
3552 std::time::Duration::from_secs(5),
3553 )
3554 .await?;
3555 assert_eq!(task.status, TaskStatus::Completed);
3556 }
3557 reconcile_run_statuses_shared(&shared, &tasks).await?;
3558 crate::utils::write_atomic(
3559 &root.join(format!("{role}-done")),
3560 serde_json::to_vec(&ids)?.as_slice(),
3561 )?;
3562 tasks.shutdown();
3563 Ok(())
3564 }
3565
3566 #[tokio::test]
3567 async fn two_scheduler_processes_preserve_the_bound_scope_and_execute_one_task() -> Result<()> {
3568 let root = tempfile::tempdir()?;
3569 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
3570 let automation = fixture_due_automation(&manager, "cross-process", 1);
3571 let mut first = spawn_scheduler_fixture(root.path(), "first")?;
3572 let mut second = spawn_scheduler_fixture(root.path(), "second")?;
3573 // Separate Runtime scopes share storage. The foreign scheduler cannot
3574 // adopt the first scope's immutable dispatch even before acceptance.
3575 wait_fixture_path(&root.path().join("first-ready")).await?;
3576 wait_fixture_path(&root.path().join("second-ready")).await?;
3577 crate::utils::write_atomic(&root.path().join("first-go"), b"go")?;
3578 wait_fixture_path(&root.path().join("first-held")).await?;
3579 let id = fs::read_to_string(root.path().join("first-held"))?;
3580 crate::utils::write_atomic(&root.path().join("second-go"), b"go")?;
3581 wait_fixture_path(&root.path().join("second-done")).await?;
3582 finish_scheduler_fixture(&mut second, &root.path().join("second.log")).await?;
3583 assert!(
3584 !root.path().join("second-entered").exists(),
3585 "a live dispatcher still owns the unaccepted intent"
3586 );
3587 assert!(
3588 fixture_executions(&root.path().join("executions")).is_empty(),
3589 "owner is blocked before actual enqueue"
3590 );
3591 crate::utils::write_atomic(&root.path().join("release-first"), b"release")?;
3592 finish_scheduler_fixture(&mut first, &root.path().join("first.log")).await?;
3593 assert_eq!(
3594 fixture_executions(&root.path().join("executions")),
3595 vec![id.clone()]
3596 );
3597 let entered: Vec<String> =
3598 serde_json::from_slice(&fs::read(root.path().join("first-done"))?)?;
3599 assert_eq!(
3600 entered,
3601 vec![id.clone()],
3602 "first process traversed the real enqueue path"
3603 );
3604 let runs = manager.list_runs(&automation.id, None)?;
3605 assert_eq!(runs.len(), 1);
3606 assert_eq!(runs[0].task_id.as_deref(), Some(id.as_str()));
3607 assert_eq!(runs[0].status, AutomationRunStatus::Completed);
3608 assert!(runs[0].dispatch.as_ref().unwrap().accepted);
3609 let task: crate::task_manager::TaskRecord = serde_json::from_slice(&fs::read(
3610 root.path().join("tasks/tasks").join(format!("{id}.json")),
3611 )?)?;
3612 assert_eq!(task.id, id);
3613 assert_eq!(task.status, TaskStatus::Completed);
3614 assert!(
3615 task.result_summary
3616 .as_deref()
3617 .unwrap_or_default()
3618 .contains("automation fixture completed")
3619 );
3620 Ok(())
3621 }
3622
3623 struct AutomationNoopExecutor;
3624 struct AutomationWatcherNoopExecutor;
3625
3626 /// A deterministic America/New_York-compatible zone for the 2026 DST
3627 /// boundary tests. Keeping the transition table local avoids mutating the
3628 /// process-wide `TZ` setting while the test binary runs in parallel.
3629 #[derive(Debug, Clone, Copy)]
3630 struct Eastern2026;
3631
3632 impl Eastern2026 {
3633 fn standard_offset() -> FixedOffset {
3634 FixedOffset::west_opt(5 * 60 * 60).expect("valid standard offset")
3635 }
3636
3637 fn daylight_offset() -> FixedOffset {
3638 FixedOffset::west_opt(4 * 60 * 60).expect("valid daylight offset")
3639 }
3640
3641 fn time(month: u32, day: u32, hour: u32) -> NaiveDateTime {
3642 NaiveDate::from_ymd_opt(2026, month, day)
3643 .expect("valid transition date")
3644 .and_hms_opt(hour, 0, 0)
3645 .expect("valid transition time")
3646 }
3647 }
3648
3649 impl TimeZone for Eastern2026 {
3650 type Offset = FixedOffset;
3651
3652 fn from_offset(_offset: &Self::Offset) -> Self {
3653 Self
3654 }
3655
3656 fn offset_from_local_date(&self, local: &NaiveDate) -> LocalResult<Self::Offset> {
3657 self.offset_from_local_datetime(
3658 &local
3659 .and_hms_opt(12, 0, 0)
3660 .expect("valid local date midpoint"),
3661 )
3662 }
3663
3664 fn offset_from_local_datetime(&self, local: &NaiveDateTime) -> LocalResult<Self::Offset> {
3665 let gap_start = Self::time(3, 8, 2);
3666 let gap_end = Self::time(3, 8, 3);
3667 let fold_start = Self::time(11, 1, 1);
3668 let fold_end = Self::time(11, 1, 2);
3669
3670 if *local >= gap_start && *local < gap_end {
3671 LocalResult::None
3672 } else if *local >= fold_start && *local < fold_end {
3673 LocalResult::Ambiguous(Self::daylight_offset(), Self::standard_offset())
3674 } else if *local >= gap_end && *local < fold_start {
3675 LocalResult::Single(Self::daylight_offset())
3676 } else {
3677 LocalResult::Single(Self::standard_offset())
3678 }
3679 }
3680
3681 fn offset_from_utc_date(&self, utc: &NaiveDate) -> Self::Offset {
3682 self.offset_from_utc_datetime(
3683 &utc.and_hms_opt(12, 0, 0).expect("valid UTC date midpoint"),
3684 )
3685 }
3686
3687 fn offset_from_utc_datetime(&self, utc: &NaiveDateTime) -> Self::Offset {
3688 let daylight_start = Self::time(3, 8, 7);
3689 let daylight_end = Self::time(11, 1, 6);
3690 if *utc >= daylight_start && *utc < daylight_end {
3691 Self::daylight_offset()
3692 } else {
3693 Self::standard_offset()
3694 }
3695 }
3696 }
3697
3698 #[async_trait]
3699 impl TaskExecutor for AutomationNoopExecutor {
3700 async fn execute(
3701 &self,
3702 _task: ExecutionTask,
3703 _events: mpsc::Sender<TaskExecutionEvent>,
3704 _cancel: CancellationToken,
3705 ) -> TaskExecutionResult {
3706 TaskExecutionResult {
3707 status: TaskStatus::Completed,
3708 result_text: Some("done".to_string()),
3709 error: None,
3710 terminal_reason: crate::task_manager::TaskTerminalReason::Completed,
3711 }
3712 }
3713 }
3714
3715 #[async_trait]
3716 impl TaskExecutor for AutomationWatcherNoopExecutor {
3717 async fn execute(
3718 &self,
3719 _task: ExecutionTask,
3720 _events: mpsc::Sender<TaskExecutionEvent>,
3721 _cancel: CancellationToken,
3722 ) -> TaskExecutionResult {
3723 TaskExecutionResult {
3724 status: TaskStatus::Completed,
3725 result_text: Some(AUTOMATION_WATCHER_NO_REPORT_SENTINEL.to_string()),
3726 error: None,
3727 terminal_reason: crate::task_manager::TaskTerminalReason::Completed,
3728 }
3729 }
3730 }
3731
3732 fn automation_task_config(root: PathBuf) -> TaskManagerConfig {
3733 TaskManagerConfig {
3734 data_dir: root,
3735 worker_count: 1,
3736 default_workspace: PathBuf::from("."),
3737 default_model: "deepseek-v4-flash".to_string(),
3738 default_mode: "plan".to_string(),
3739 allow_shell: true,
3740 trust_mode: true,
3741 execution_limits: crate::task_manager::TaskExecutionLimits::default(),
3742 }
3743 }
3744
3745 fn automation_record_with_settings(
3746 mode: Option<&str>,
3747 allow_shell: Option<bool>,
3748 trust_mode: Option<bool>,
3749 auto_approve: Option<bool>,
3750 ) -> AutomationRecord {
3751 let now = Utc::now();
3752 AutomationRecord {
3753 schema_version: CURRENT_AUTOMATION_SCHEMA_VERSION,
3754 execution_scope: Some(crate::task_manager::test_execution_scope("test")),
3755 id: Uuid::new_v4().to_string(),
3756 name: "Test automation".to_string(),
3757 prompt: "Run the automation".to_string(),
3758 rrule: "FREQ=HOURLY;INTERVAL=1".to_string(),
3759 cwds: Vec::new(),
3760 model: None,
3761 model_provider: None,
3762 model_provider_id: None,
3763 mode: mode.map(ToString::to_string),
3764 allow_shell,
3765 trust_mode,
3766 auto_approve,
3767 delivery_mode: None,
3768 status: AutomationStatus::Active,
3769 created_at: now,
3770 updated_at: now,
3771 next_run_at: None,
3772 last_run_at: None,
3773 }
3774 }
3775
3776 fn queued_run_for(automation: &AutomationRecord) -> AutomationRunRecord {
3777 let now = Utc::now();
3778 AutomationRunRecord {
3779 schema_version: CURRENT_RUN_SCHEMA_VERSION,
3780 id: Uuid::new_v4().to_string(),
3781 automation_id: automation.id.clone(),
3782 scheduled_for: now,
3783 status: AutomationRunStatus::Queued,
3784 created_at: now,
3785 started_at: None,
3786 ended_at: None,
3787 task_id: None,
3788 thread_id: None,
3789 turn_id: None,
3790 error: None,
3791 dispatch: None,
3792 }
3793 }
3794
3795 fn eastern_datetime(year: i32, month: u32, day: u32, hour: u32, minute: u32) -> DateTime<Utc> {
3796 Eastern2026
3797 .with_ymd_and_hms(year, month, day, hour, minute, 0)
3798 .single()
3799 .expect("unambiguous Eastern wall time")
3800 .with_timezone(&Utc)
3801 }
3802
3803 fn anchored_automation(
3804 created_at: DateTime<Utc>,
3805 status: AutomationStatus,
3806 ) -> AutomationRecord {
3807 let mut record = automation_record_with_settings(None, None, None, None);
3808 record.rrule = "FREQ=HOURLY;INTERVAL=7;BYMINUTE=17".to_string();
3809 record.status = status;
3810 record.created_at = created_at;
3811 record.updated_at = created_at;
3812 record.next_run_at = None;
3813 record
3814 }
3815
3816 fn local_naive_to_utc(naive: NaiveDateTime) -> DateTime<Utc> {
3817 Local
3818 .from_local_datetime(&naive)
3819 .earliest()
3820 .expect("valid unambiguous local time")
3821 .with_timezone(&Utc)
3822 }
3823
3824 #[test]
3825 fn parses_hourly_rrule() {
3826 let parsed =
3827 AutomationSchedule::parse_rrule("FREQ=HOURLY;INTERVAL=2;BYDAY=MO,TU").expect("parse");
3828 match parsed {
3829 AutomationSchedule::Hourly {
3830 interval_hours,
3831 byday,
3832 ..
3833 } => {
3834 assert_eq!(interval_hours, 2);
3835 assert_eq!(byday.expect("byday").len(), 2);
3836 }
3837 _ => panic!("expected hourly"),
3838 }
3839 }
3840
3841 #[test]
3842 fn parses_once_rrule() {
3843 let parsed =
3844 AutomationSchedule::parse_rrule("FREQ=ONCE;AT=2026-08-03T14:30").expect("parse");
3845 match parsed {
3846 AutomationSchedule::Once { at } => {
3847 assert_eq!(
3848 at,
3849 local_naive_to_utc(
3850 NaiveDateTime::parse_from_str("2026-08-03T14:30", "%Y-%m-%dT%H:%M")
3851 .expect("naive")
3852 )
3853 );
3854 }
3855 _ => panic!("expected once"),
3856 }
3857 }
3858
3859 #[test]
3860 fn parses_hourly_clock_anchor() {
3861 let parsed =
3862 AutomationSchedule::parse_rrule("FREQ=HOURLY;INTERVAL=24;BYHOUR=8;BYMINUTE=30")
3863 .expect("parse anchored hourly schedule");
3864
3865 assert!(matches!(
3866 parsed,
3867 AutomationSchedule::Hourly {
3868 anchor_hour: Some(8),
3869 anchor_minute: Some(30),
3870 ..
3871 }
3872 ));
3873
3874 let minute_only = AutomationSchedule::parse_rrule("FREQ=HOURLY;INTERVAL=1;BYMINUTE=15")
3875 .expect("parse minute-only anchor");
3876 assert!(matches!(
3877 minute_only,
3878 AutomationSchedule::Hourly {
3879 anchor_hour: None,
3880 anchor_minute: Some(15),
3881 ..
3882 }
3883 ));
3884 }
3885
3886 #[test]
3887 fn anchored_hourly_schedule_keeps_wall_time_across_spring_forward() {
3888 let schedule =
3889 AutomationSchedule::parse_rrule("FREQ=HOURLY;INTERVAL=24;BYHOUR=8;BYMINUTE=30")
3890 .expect("parse");
3891 let created_at = eastern_datetime(2026, 3, 6, 7, 0);
3892 let after = eastern_datetime(2026, 3, 7, 9, 0);
3893
3894 let next = schedule
3895 .next_after_in_timezone(after, created_at, &Eastern2026)
3896 .expect("next run");
3897
3898 assert_eq!(next, eastern_datetime(2026, 3, 8, 8, 30));
3899 }
3900
3901 #[test]
3902 fn anchored_hourly_schedule_keeps_wall_time_across_fall_back() {
3903 let schedule =
3904 AutomationSchedule::parse_rrule("FREQ=HOURLY;INTERVAL=24;BYHOUR=8;BYMINUTE=30")
3905 .expect("parse");
3906 let created_at = eastern_datetime(2026, 10, 30, 7, 0);
3907 let after = eastern_datetime(2026, 10, 31, 9, 0);
3908
3909 let next = schedule
3910 .next_after_in_timezone(after, created_at, &Eastern2026)
3911 .expect("next run");
3912
3913 assert_eq!(next, eastern_datetime(2026, 11, 1, 8, 30));
3914 }
3915
3916 #[test]
3917 fn anchored_hourly_schedule_skips_nonexistent_wall_time() {
3918 let schedule =
3919 AutomationSchedule::parse_rrule("FREQ=HOURLY;INTERVAL=24;BYHOUR=2;BYMINUTE=30")
3920 .expect("parse");
3921 let created_at = eastern_datetime(2026, 3, 7, 1, 0);
3922 let after = eastern_datetime(2026, 3, 7, 3, 0);
3923
3924 let next = schedule
3925 .next_after_in_timezone(after, created_at, &Eastern2026)
3926 .expect("next run after spring-forward gap");
3927
3928 assert_eq!(next, eastern_datetime(2026, 3, 9, 2, 30));
3929 }
3930
3931 #[test]
3932 fn anchored_hourly_schedule_uses_first_ambiguous_wall_time_once() {
3933 let schedule =
3934 AutomationSchedule::parse_rrule("FREQ=HOURLY;INTERVAL=24;BYHOUR=1;BYMINUTE=30")
3935 .expect("parse");
3936 let created_at = eastern_datetime(2026, 10, 31, 0, 0);
3937 let after = eastern_datetime(2026, 10, 31, 2, 0);
3938 let first_fold_occurrence = Eastern2026
3939 .with_ymd_and_hms(2026, 11, 1, 1, 30, 0)
3940 .earliest()
3941 .expect("first fold occurrence")
3942 .with_timezone(&Utc);
3943
3944 let next = schedule
3945 .next_after_in_timezone(after, created_at, &Eastern2026)
3946 .expect("next run at fall-back fold");
3947 assert_eq!(next, first_fold_occurrence);
3948
3949 let during_second_fold = Eastern2026
3950 .with_ymd_and_hms(2026, 11, 1, 1, 15, 0)
3951 .latest()
3952 .expect("second fold occurrence")
3953 .with_timezone(&Utc);
3954 let after_fold = schedule
3955 .next_after_in_timezone(during_second_fold, created_at, &Eastern2026)
3956 .expect("next run after fold");
3957 assert_eq!(after_fold, eastern_datetime(2026, 11, 2, 1, 30));
3958 }
3959
3960 #[test]
3961 fn anchored_hourly_schedule_reuses_persisted_anchor_after_restart_and_resume() {
3962 let rrule = "FREQ=HOURLY;INTERVAL=24;BYHOUR=8;BYMINUTE=30";
3963 let created_at = eastern_datetime(2026, 3, 6, 7, 0);
3964 let schedule = AutomationSchedule::parse_rrule(rrule).expect("parse");
3965 let before_restart = schedule
3966 .next_after_in_timezone(
3967 eastern_datetime(2026, 3, 7, 12, 0),
3968 created_at,
3969 &Eastern2026,
3970 )
3971 .expect("next before restart");
3972 assert_eq!(before_restart, eastern_datetime(2026, 3, 8, 8, 30));
3973
3974 // Reparsing models a process restart; the persisted creation timestamp
3975 // remains the recurrence anchor when the record is loaded or resumed.
3976 let restarted = AutomationSchedule::parse_rrule(rrule).expect("reparse after restart");
3977 let after_restart = restarted
3978 .next_after_in_timezone(
3979 eastern_datetime(2026, 3, 8, 10, 0),
3980 created_at,
3981 &Eastern2026,
3982 )
3983 .expect("next after restart");
3984 assert_eq!(after_restart, eastern_datetime(2026, 3, 9, 8, 30));
3985
3986 let after_resume = restarted
3987 .next_after_in_timezone(
3988 eastern_datetime(2026, 3, 10, 12, 0),
3989 created_at,
3990 &Eastern2026,
3991 )
3992 .expect("next after resume");
3993 assert_eq!(after_resume, eastern_datetime(2026, 3, 11, 8, 30));
3994 }
3995
3996 #[test]
3997 fn scheduler_restart_uses_persisted_creation_anchor() {
3998 let tempdir = tempfile::tempdir().expect("tempdir");
3999 let now = Utc::now();
4000 let created_at = now - Duration::hours(51);
4001 let automation = anchored_automation(created_at, AutomationStatus::Active);
4002 let schedule = AutomationSchedule::parse_rrule(&automation.rrule).expect("parse");
4003 let expected = schedule
4004 .next_after_with_anchor(now, created_at)
4005 .expect("persisted-anchor schedule");
4006 let reset_anchor = schedule
4007 .next_after_with_anchor(now, now)
4008 .expect("reset-anchor schedule");
4009 assert_ne!(expected, reset_anchor, "fixture must detect anchor resets");
4010
4011 let manager =
4012 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4013 manager.save_automation(&automation).expect("save");
4014 drop(manager);
4015
4016 let restarted =
4017 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("reopen");
4018 assert!(
4019 restarted
4020 .collect_due_runs(now)
4021 .expect("restart tick")
4022 .is_empty(),
4023 "an uninitialized future slot must not enqueue immediately"
4024 );
4025 let reloaded = restarted
4026 .get_automation(&automation.id)
4027 .expect("reloaded automation");
4028 assert_eq!(reloaded.next_run_at, Some(expected));
4029 }
4030
4031 #[test]
4032 fn resume_uses_persisted_creation_anchor() {
4033 let tempdir = tempfile::tempdir().expect("tempdir");
4034 let manager =
4035 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4036 let before = Utc::now();
4037 let created_at = before - Duration::hours(51);
4038 let automation = anchored_automation(created_at, AutomationStatus::Paused);
4039 let schedule = AutomationSchedule::parse_rrule(&automation.rrule).expect("parse");
4040 manager.save_automation(&automation).expect("save");
4041
4042 let expected_before = schedule
4043 .next_after_with_anchor(before, created_at)
4044 .expect("next before resume");
4045 let reset_anchor = schedule
4046 .next_after_with_anchor(before, before)
4047 .expect("reset-anchor schedule");
4048 assert_ne!(
4049 expected_before, reset_anchor,
4050 "fixture must detect anchor resets"
4051 );
4052
4053 let resumed = manager
4054 .resume_automation(&automation.id)
4055 .expect("resume automation");
4056 let after = Utc::now();
4057 let expected_after = schedule
4058 .next_after_with_anchor(after, created_at)
4059 .expect("next after resume");
4060 let actual = resumed.next_run_at.expect("resumed next run");
4061 assert!(
4062 actual == expected_before || actual == expected_after,
4063 "resume must keep the persisted creation anchor"
4064 );
4065 }
4066
4067 #[test]
4068 fn anchored_hourly_schedule_applies_byday_on_calendar_slots() {
4069 let schedule = AutomationSchedule::parse_rrule(
4070 "FREQ=HOURLY;INTERVAL=24;BYDAY=MO,TU,WE,TH,FR;BYHOUR=8;BYMINUTE=30",
4071 )
4072 .expect("parse");
4073 let created_at = eastern_datetime(2026, 3, 6, 7, 0);
4074
4075 let next = schedule
4076 .next_after_in_timezone(eastern_datetime(2026, 3, 6, 9, 0), created_at, &Eastern2026)
4077 .expect("next weekday run");
4078
4079 assert_eq!(next, eastern_datetime(2026, 3, 9, 8, 30));
4080 }
4081
4082 #[test]
4083 fn parses_weekly_rrule() {
4084 let parsed =
4085 AutomationSchedule::parse_rrule("FREQ=WEEKLY;BYDAY=MO,WE;BYHOUR=9;BYMINUTE=30")
4086 .expect("parse");
4087 match parsed {
4088 AutomationSchedule::Weekly {
4089 byday,
4090 byhour,
4091 byminute,
4092 } => {
4093 assert_eq!(byday.len(), 2);
4094 assert_eq!(byhour, 9);
4095 assert_eq!(byminute, 30);
4096 }
4097 _ => panic!("expected weekly"),
4098 }
4099 }
4100
4101 #[test]
4102 fn parses_cron_rrule_and_computes_next_minute_slot() {
4103 let schedule =
4104 AutomationSchedule::parse_rrule("FREQ=CRON;EXPR=*/17 * * * *").expect("parse");
4105 let after = Utc
4106 .with_ymd_and_hms(2026, 8, 3, 9, 17, 1)
4107 .single()
4108 .expect("after");
4109 let next = schedule
4110 .next_after_in_timezone(after, after, &Utc)
4111 .expect("next cron run");
4112 assert_eq!(
4113 next,
4114 Utc.with_ymd_and_hms(2026, 8, 3, 9, 34, 0)
4115 .single()
4116 .expect("next")
4117 );
4118 }
4119
4120 #[test]
4121 fn cron_weekday_schedule_uses_standard_five_field_local_time() {
4122 let schedule =
4123 AutomationSchedule::parse_rrule("FREQ=CRON;EXPR=3 9 * * MON-FRI").expect("parse");
4124 let after = Utc
4125 .with_ymd_and_hms(2026, 8, 7, 9, 4, 0)
4126 .single()
4127 .expect("after");
4128 let next = schedule
4129 .next_after_in_timezone(after, after, &Utc)
4130 .expect("next weekday cron run");
4131 assert_eq!(
4132 next,
4133 Utc.with_ymd_and_hms(2026, 8, 10, 9, 3, 0)
4134 .single()
4135 .expect("next")
4136 );
4137 }
4138
4139 #[test]
4140 fn cron_rejects_impossible_date() {
4141 let err = AutomationSchedule::parse_rrule("FREQ=CRON;EXPR=0 9 31 2 *")
4142 .expect_err("impossible february date must fail");
4143 assert!(err.to_string().contains("can never occur"));
4144 }
4145
4146 #[test]
4147 fn rejects_invalid_rrule_fields() {
4148 let err =
4149 AutomationSchedule::parse_rrule("FREQ=WEEKLY;BYSECOND=5").expect_err("should fail");
4150 assert!(err.to_string().contains("Unsupported RRULE field"));
4151 }
4152
4153 #[test]
4154 fn automation_model_round_trips_through_create_and_update() {
4155 let tempdir = tempfile::tempdir().expect("tempdir");
4156 let manager =
4157 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4158
4159 let created = manager
4160 .create_automation(CreateAutomationRequest {
4161 name: "Pinned model".to_string(),
4162 prompt: "prompt".to_string(),
4163 rrule: "FREQ=HOURLY;INTERVAL=1".to_string(),
4164 cwds: Vec::new(),
4165 model: Some(" scheduled-model ".to_string()),
4166 model_provider: None,
4167 model_provider_id: None,
4168 mode: None,
4169 allow_shell: None,
4170 trust_mode: None,
4171 auto_approve: None,
4172 delivery_mode: None,
4173 status: Some(AutomationStatus::Active),
4174 })
4175 .expect("create");
4176 assert_eq!(created.model.as_deref(), Some("scheduled-model"));
4177 assert_eq!(
4178 manager
4179 .get_automation(&created.id)
4180 .expect("reload")
4181 .model
4182 .as_deref(),
4183 Some("scheduled-model")
4184 );
4185
4186 let updated = manager
4187 .update_automation(
4188 &created.id,
4189 UpdateAutomationRequest {
4190 model: Some("replacement-model".to_string()),
4191 ..UpdateAutomationRequest::default()
4192 },
4193 )
4194 .expect("update");
4195 assert_eq!(updated.model.as_deref(), Some("replacement-model"));
4196 }
4197
4198 #[test]
4199 fn deletes_definition_and_settled_runs_but_retains_unfinished_receipts() {
4200 let tempdir = tempfile::tempdir().expect("tempdir");
4201 let manager =
4202 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4203
4204 let created = manager
4205 .create_automation(CreateAutomationRequest {
4206 name: "Delete me".to_string(),
4207 prompt: "prompt".to_string(),
4208 rrule: "FREQ=HOURLY;INTERVAL=1".to_string(),
4209 cwds: Vec::new(),
4210 model: None,
4211 model_provider: None,
4212 model_provider_id: None,
4213 mode: None,
4214 allow_shell: None,
4215 trust_mode: None,
4216 auto_approve: None,
4217 delivery_mode: None,
4218 status: Some(AutomationStatus::Active),
4219 })
4220 .expect("create");
4221
4222 let run = AutomationRunRecord {
4223 schema_version: CURRENT_RUN_SCHEMA_VERSION,
4224 id: Uuid::new_v4().to_string(),
4225 automation_id: created.id.clone(),
4226 scheduled_for: Utc::now(),
4227 status: AutomationRunStatus::Queued,
4228 created_at: Utc::now(),
4229 started_at: None,
4230 ended_at: None,
4231 task_id: None,
4232 thread_id: None,
4233 turn_id: None,
4234 error: None,
4235 dispatch: None,
4236 };
4237 manager.save_run(&run).expect("save run");
4238 let settled = AutomationRunRecord {
4239 id: Uuid::new_v4().to_string(),
4240 status: AutomationRunStatus::Completed,
4241 ended_at: Some(Utc::now()),
4242 ..run.clone()
4243 };
4244 manager.save_run(&settled).expect("save settled run");
4245 assert!(
4246 manager
4247 .runs_dir_for(&created.id)
4248 .expect("runs dir")
4249 .exists()
4250 );
4251
4252 manager
4253 .delete_automation(&created.id)
4254 .expect("delete automation");
4255
4256 assert!(manager.get_automation(&created.id).is_err());
4257 let remaining = manager
4258 .list_runs(&created.id, None)
4259 .expect("retained receipts");
4260 assert_eq!(remaining.len(), 1);
4261 assert_eq!(remaining[0].id, run.id);
4262 }
4263
4264 #[test]
4265 fn automation_storage_rejects_traversal_ids() {
4266 let tempdir = tempfile::tempdir().expect("tempdir");
4267 let manager =
4268 AutomationManager::open_for_test(tempdir.path().join("root")).expect("manager");
4269 let escaped_file = tempdir.path().join("escape.json");
4270 let escaped_runs = tempdir.path().join("escape-runs");
4271
4272 let err = manager
4273 .get_automation("../escape")
4274 .expect_err("traversal automation ids must be rejected");
4275 assert!(err.to_string().contains("single path component"));
4276 assert!(!escaped_file.exists());
4277
4278 let err = manager
4279 .list_runs("../escape-runs", None)
4280 .expect_err("traversal run dirs must be rejected");
4281 assert!(err.to_string().contains("single path component"));
4282 assert!(!escaped_runs.exists());
4283
4284 let run = AutomationRunRecord {
4285 schema_version: CURRENT_RUN_SCHEMA_VERSION,
4286 id: "../escape-run".to_string(),
4287 automation_id: Uuid::new_v4().to_string(),
4288 scheduled_for: Utc::now(),
4289 status: AutomationRunStatus::Queued,
4290 created_at: Utc::now(),
4291 started_at: None,
4292 ended_at: None,
4293 task_id: None,
4294 thread_id: None,
4295 turn_id: None,
4296 error: None,
4297 dispatch: None,
4298 };
4299 let err = manager
4300 .save_run(&run)
4301 .expect_err("traversal run ids must be rejected");
4302 assert!(err.to_string().contains("single path component"));
4303 assert!(!tempdir.path().join("escape-run.json").exists());
4304 }
4305
4306 #[test]
4307 fn automation_task_settings_default_for_legacy_records() {
4308 let now = Utc::now().to_rfc3339();
4309 let record: AutomationRecord = serde_json::from_value(serde_json::json!({
4310 "schema_version": CURRENT_AUTOMATION_SCHEMA_VERSION,
4311 "id": Uuid::new_v4().to_string(),
4312 "name": "Legacy automation",
4313 "prompt": "Run legacy automation",
4314 "rrule": "FREQ=HOURLY;INTERVAL=1",
4315 "cwds": [],
4316 "status": "active",
4317 "created_at": now,
4318 "updated_at": now
4319 }))
4320 .expect("legacy automation record should deserialize");
4321
4322 assert_eq!(record.mode, None);
4323 assert_eq!(record.task_mode(), "agent");
4324 assert!(!record.task_allow_shell());
4325 assert!(!record.task_trust_mode());
4326 assert!(!record.task_auto_approve());
4327 assert_eq!(record.delivery_mode(), AutomationDeliveryMode::Task);
4328 }
4329
4330 #[test]
4331 fn automation_requests_pin_the_posture_its_bit_stands_for() {
4332 let now = Utc::now().to_rfc3339();
4333 let legacy: AutomationRecord = serde_json::from_value(serde_json::json!({
4334 "schema_version": CURRENT_AUTOMATION_SCHEMA_VERSION,
4335 "id": Uuid::new_v4().to_string(),
4336 "name": "Nightly",
4337 "prompt": "Run the nightly sweep",
4338 "rrule": "FREQ=DAILY",
4339 "cwds": [],
4340 "status": "active",
4341 "created_at": now,
4342 "updated_at": now
4343 }))
4344 .expect("legacy automation record");
4345
4346 // No session exists when a scheduler fires, so the request states the
4347 // posture the record's own authority means instead of leaving the
4348 // runtime to re-derive it from the bit.
4349 assert!(legacy.auto_approve.is_none());
4350 assert_eq!(
4351 automation_task_request(&legacy)
4352 .permission_posture
4353 .as_deref(),
4354 Some("ask")
4355 );
4356
4357 let elevated: AutomationRecord = serde_json::from_value(serde_json::json!({
4358 "schema_version": CURRENT_AUTOMATION_SCHEMA_VERSION,
4359 "id": Uuid::new_v4().to_string(),
4360 "name": "Nightly, unattended",
4361 "prompt": "Run the nightly sweep",
4362 "rrule": "FREQ=DAILY",
4363 "cwds": [],
4364 "auto_approve": true,
4365 "status": "active",
4366 "created_at": now,
4367 "updated_at": now
4368 }))
4369 .expect("elevated automation record");
4370 assert_eq!(
4371 automation_task_request(&elevated)
4372 .permission_posture
4373 .as_deref(),
4374 Some("full_access")
4375 );
4376 }
4377
4378 #[tokio::test]
4379 async fn automation_enqueue_uses_default_and_explicit_task_settings() -> Result<()> {
4380 let tempdir = tempfile::tempdir().expect("tempdir");
4381 let task_manager = TaskManager::start_with_executor(
4382 automation_task_config(tempdir.path().join("tasks")),
4383 std::sync::Arc::new(AutomationNoopExecutor),
4384 )
4385 .await?;
4386
4387 let default_automation = automation_record_with_settings(None, None, None, None);
4388 let mut default_run = queued_run_for(&default_automation);
4389 bind_run_dispatch(
4390 &mut default_run,
4391 &default_automation,
4392 &task_manager.data_dir(),
4393 false,
4394 )?;
4395 enqueue_run_task(&mut default_run, &task_manager).await;
4396 let default_task = task_manager
4397 .get_task(default_run.task_id.as_deref().expect("task id"))
4398 .await?;
4399 assert_eq!(default_task.model, "deepseek-v4-flash");
4400 assert_eq!(default_task.mode, "agent");
4401 assert!(!default_task.allow_shell);
4402 assert!(!default_task.trust_mode);
4403 assert!(!default_task.auto_approve);
4404
4405 let mut explicit_automation =
4406 automation_record_with_settings(Some("plan"), Some(true), Some(true), Some(true));
4407 explicit_automation.model = Some("scheduled-model".to_string());
4408 let mut explicit_run = queued_run_for(&explicit_automation);
4409 bind_run_dispatch(
4410 &mut explicit_run,
4411 &explicit_automation,
4412 &task_manager.data_dir(),
4413 false,
4414 )?;
4415 enqueue_run_task(&mut explicit_run, &task_manager).await;
4416 let explicit_task = task_manager
4417 .get_task(explicit_run.task_id.as_deref().expect("task id"))
4418 .await?;
4419 assert_eq!(explicit_task.model, "scheduled-model");
4420 assert_eq!(explicit_task.mode, "plan");
4421 assert!(explicit_task.allow_shell);
4422 assert!(explicit_task.trust_mode);
4423 assert!(explicit_task.auto_approve);
4424
4425 task_manager.shutdown();
4426 Ok(())
4427 }
4428
4429 #[tokio::test]
4430 async fn automation_provider_pin_survives_active_route_change_and_legacy_inherits() -> Result<()>
4431 {
4432 let _env = crate::test_support::lock_test_env();
4433 let root = tempfile::tempdir()?;
4434 let mut config: crate::config::Config = toml::from_str(
4435 r#"
4436 provider = "first"
4437 [providers.first]
4438 kind = "openai-compatible"
4439 base_url = "http://127.0.0.1:9/first/v1"
4440 api_key = "fixture-first"
4441 model = "private-model"
4442 [providers.second]
4443 kind = "openai-compatible"
4444 base_url = "http://127.0.0.1:9/second/v1"
4445 api_key = "fixture-second"
4446 model = "private-model"
4447 "#,
4448 )?;
4449 let automations = AutomationManager::open_for_test(root.path().join("schedules"))?;
4450 let mut record = automation_record_with_settings(None, None, None, None);
4451 record.schema_version = 1;
4452 record.model = Some("private-model".to_string());
4453 record.model_provider = Some("custom".to_string());
4454 record.model_provider_id = Some("first".to_string());
4455 automations.save_automation(&record)?;
4456 let record = automations.get_automation(&record.id)?;
4457 assert_eq!(
4458 record.schema_version, CURRENT_AUTOMATION_SCHEMA_VERSION,
4459 "a pin cannot be ignored by an older reader"
4460 );
4461
4462 let runtime = crate::runtime_threads::RuntimeThreadManager::open(
4463 config.clone(),
4464 root.path().to_path_buf(),
4465 crate::runtime_threads::RuntimeThreadManagerConfig::from_task_data_dir(
4466 root.path().join("runtime"),
4467 ),
4468 )?;
4469 config.provider = Some("second".to_string());
4470 runtime.reload_config(config).await?;
4471 let tasks = TaskManager::start_with_executor(
4472 automation_task_config(root.path().join("tasks")),
4473 std::sync::Arc::new(AutomationNoopExecutor),
4474 )
4475 .await?;
4476 let mut run = queued_run_for(&record);
4477 bind_run_dispatch(&mut run, &record, &tasks.data_dir(), false)?;
4478 enqueue_run_task(&mut run, &tasks).await;
4479 let task = tasks
4480 .get_task(run.task_id.as_deref().expect("task id"))
4481 .await?;
4482 let task: crate::task_manager::TaskRecord =
4483 serde_json::from_slice(&serde_json::to_vec(&task)?)?;
4484 assert_eq!(task.schema_version, 4);
4485 let thread = runtime
4486 .create_thread(ExecutionTask::from(&task).thread_request())
4487 .await?;
4488 assert_eq!(thread.model, "private-model");
4489 assert_eq!(thread.model_provider.as_deref(), Some("custom"));
4490 assert_eq!(thread.model_provider_id.as_deref(), Some("first"));
4491
4492 let mut legacy = serde_json::to_value(&record)?;
4493 let object = legacy.as_object_mut().unwrap();
4494 object.remove("model_provider");
4495 object.remove("model_provider_id");
4496 object.insert("schema_version".to_string(), serde_json::json!(1));
4497 let legacy: AutomationRecord = serde_json::from_value(legacy)?;
4498 assert_eq!(legacy.model_provider, None);
4499 let mut run = queued_run_for(&legacy);
4500 bind_run_dispatch(&mut run, &legacy, &tasks.data_dir(), false)?;
4501 enqueue_run_task(&mut run, &tasks).await;
4502 let task = tasks
4503 .get_task(run.task_id.as_deref().expect("legacy task id"))
4504 .await?;
4505 let mut value = serde_json::to_value(&task)?;
4506 value.as_object_mut().unwrap().remove("model_provider");
4507 value.as_object_mut().unwrap().remove("model_provider_id");
4508 value["schema_version"] = serde_json::json!(2);
4509 let task: crate::task_manager::TaskRecord = serde_json::from_value(value)?;
4510 let thread = runtime
4511 .create_thread(ExecutionTask::from(&task).thread_request())
4512 .await?;
4513 assert_eq!(thread.model_provider_id.as_deref(), Some("second"));
4514 assert_eq!(thread.model, "private-model");
4515 tasks.shutdown();
4516 Ok(())
4517 }
4518
4519 #[tokio::test]
4520 async fn delayed_trigger_fires_task_with_same_owner_and_skips_legacy_ownerless() -> Result<()> {
4521 let tempdir = tempfile::tempdir().expect("tempdir");
4522 let task_manager = TaskManager::start_with_executor(
4523 automation_task_config(tempdir.path().join("tasks")),
4524 std::sync::Arc::new(AutomationNoopExecutor),
4525 )
4526 .await?;
4527 let manager = AutomationManager::open_for_test(tempdir.path().join("automations"))?;
4528
4529 let mut owned = manager.create_trigger(CreateDelayedTriggerRequest {
4530 fire_at: Utc::now() + Duration::hours(1),
4531 message: "owned delayed continuation".to_string(),
4532 workspace: None,
4533 owner_session_id: Some("session-a".to_string()),
4534 parent_trigger_id: None,
4535 })?;
4536 owned.fire_at = Utc::now() - Duration::minutes(1);
4537 manager.save_trigger(&owned)?;
4538
4539 let mut legacy = manager.create_trigger(CreateDelayedTriggerRequest {
4540 fire_at: Utc::now() + Duration::hours(1),
4541 message: "legacy delayed continuation".to_string(),
4542 workspace: None,
4543 owner_session_id: None,
4544 parent_trigger_id: None,
4545 })?;
4546 legacy.fire_at = Utc::now() - Duration::minutes(1);
4547 manager.save_trigger(&legacy)?;
4548
4549 let shared = Arc::new(Mutex::new(manager));
4550 fire_due_triggers_shared(&shared, &task_manager).await?;
4551
4552 let manager = shared.lock().await;
4553 let fired = manager.get_trigger(&owned.trigger_id)?;
4554 assert_eq!(fired.status, DelayedTriggerStatus::Fired);
4555 let task = task_manager
4556 .get_task(fired.task_id.as_deref().expect("fired task id"))
4557 .await?;
4558 assert_eq!(task.owner_session_id.as_deref(), Some("session-a"));
4559
4560 let legacy_after = manager.get_trigger(&legacy.trigger_id)?;
4561 assert_eq!(legacy_after.status, DelayedTriggerStatus::Pending);
4562 assert!(legacy_after.task_id.is_none());
4563 drop(manager);
4564 task_manager.shutdown();
4565 Ok(())
4566 }
4567
4568 #[test]
4569 fn legacy_delayed_trigger_deserializes_without_owner() -> Result<()> {
4570 let now = Utc::now();
4571 let record: DelayedTriggerRecord = serde_json::from_value(serde_json::json!({
4572 "schema_version": CURRENT_TRIGGER_SCHEMA_VERSION,
4573 "trigger_id": "trig_legacy",
4574 "fire_at": (now + Duration::hours(1)).to_rfc3339(),
4575 "message": "legacy trigger",
4576 "status": "pending",
4577 "created_at": now.to_rfc3339(),
4578 "fired_at": null,
4579 "task_id": null,
4580 "thread_id": null,
4581 "error": null,
4582 "parent_trigger_id": null
4583 }))?;
4584 assert_eq!(record.owner_session_id, None);
4585 Ok(())
4586 }
4587
4588 fn write_legacy_run_file(manager: &AutomationManager, run: &AutomationRunRecord) {
4589 let dir = manager.runs_dir_for(&run.automation_id).expect("runs dir");
4590 fs::create_dir_all(&dir).expect("create runs dir");
4591 fs::write(
4592 dir.join(format!("{}.json", run.id)),
4593 serde_json::to_string_pretty(run).expect("serialize run"),
4594 )
4595 .expect("write legacy run");
4596 }
4597
4598 fn run_created_at(
4599 automation: &AutomationRecord,
4600 created_at: DateTime<Utc>,
4601 ) -> AutomationRunRecord {
4602 let mut run = queued_run_for(automation);
4603 run.created_at = created_at;
4604 run.scheduled_for = created_at;
4605 run
4606 }
4607
4608 #[test]
4609 fn interrupted_watcher_migration_deduplicates_before_visibility() -> Result<()> {
4610 for suppressed in [false, true] {
4611 let root = tempfile::tempdir()?;
4612 let manager = AutomationManager::open_for_test(root.path().join("automations"))?;
4613 let automation = automation_record_with_settings(None, None, None, None);
4614 let mut run = queued_run_for(&automation);
4615 write_legacy_run_file(&manager, &run);
4616 bind_run_dispatch(&mut run, &automation, root.path(), true)?;
4617 run.status = AutomationRunStatus::Completed;
4618 let dispatch = run.dispatch.as_mut().unwrap();
4619 dispatch.accepted = true;
4620 dispatch.delivery_mode = AutomationDeliveryMode::Watcher;
4621 dispatch.suppress_report = suppressed;
4622 // Persist the post-write/pre-legacy-removal crash boundary.
4623 write_json_atomic(&manager.run_path(&run)?, &run)?;
4624 let ids = std::collections::BTreeSet::from([run.id.clone()]);
4625 for visible in [
4626 manager.list_runs(&automation.id, None)?,
4627 manager.list_runs(&automation.id, Some(1))?,
4628 manager.get_runs_by_ids(&automation.id, &ids)?,
4629 ] {
4630 assert_eq!(visible.len(), usize::from(!suppressed));
4631 if let Some(receipt) = visible.first() {
4632 assert_eq!(receipt.status, AutomationRunStatus::Completed);
4633 assert_eq!(receipt.task_id, run.task_id);
4634 }
4635 }
4636 let durable = manager.list_runs_with_visibility(&automation.id, None, true)?;
4637 assert_eq!(durable.len(), 1);
4638 assert_eq!(durable[0].status, AutomationRunStatus::Completed);
4639 assert_eq!(
4640 durable[0].dispatch.as_ref().unwrap().suppress_report,
4641 suppressed
4642 );
4643 assert!(manager.collect_pending_runs()?.is_empty());
4644 }
4645 Ok(())
4646 }
4647
4648 #[test]
4649 fn save_run_uses_sortable_names_and_migrates_legacy_files() {
4650 let tempdir = tempfile::tempdir().expect("tempdir");
4651 let manager =
4652 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4653 let automation = automation_record_with_settings(None, None, None, None);
4654 let run = queued_run_for(&automation);
4655
4656 write_legacy_run_file(&manager, &run);
4657 manager.save_run(&run).expect("save run");
4658
4659 let dir = manager.runs_dir_for(&automation.id).expect("runs dir");
4660 let names: Vec<String> = fs::read_dir(&dir)
4661 .expect("read dir")
4662 .map(|entry| {
4663 entry
4664 .expect("entry")
4665 .file_name()
4666 .to_string_lossy()
4667 .into_owned()
4668 })
4669 .collect();
4670 let expected = format!("{}-{}.json", run_file_stamp(run.created_at), run.id);
4671 assert_eq!(names, vec![expected.clone()]);
4672 assert!(has_sortable_run_stem(expected.trim_end_matches(".json")));
4673 // Legacy uuid stems are not mistaken for sortable names.
4674 assert!(!has_sortable_run_stem(&run.id));
4675 }
4676
4677 #[test]
4678 fn finish_scheduled_run_persists_run_when_automation_deleted_mid_enqueue() {
4679 let tempdir = tempfile::tempdir().expect("tempdir");
4680 let manager =
4681 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4682 let automation = automation_record_with_settings(None, None, None, None);
4683 manager.save_automation(&automation).expect("save");
4684 let run = queued_run_for(&automation);
4685
4686 // Simulate the automation being deleted while the enqueue await ran
4687 // outside the lock. The task already exists in the task manager at
4688 // this point, so the run record must still be persisted — an early
4689 // return here orphans a real running task.
4690 manager.delete_automation(&automation.id).expect("delete");
4691 manager
4692 .finish_scheduled_run(&run, Utc::now())
4693 .expect("finish");
4694
4695 let runs = manager.list_runs(&automation.id, None).expect("list runs");
4696 assert_eq!(
4697 runs.iter().map(|r| r.id.as_str()).collect::<Vec<_>>(),
4698 vec![run.id.as_str()],
4699 "run must be persisted even though its automation was deleted"
4700 );
4701 assert!(
4702 manager.get_automation(&automation.id).is_err(),
4703 "the deleted automation must not be resurrected"
4704 );
4705 }
4706
4707 #[test]
4708 fn once_schedule_fires_once_and_auto_completes() {
4709 let tempdir = tempfile::tempdir().expect("tempdir");
4710 let manager =
4711 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4712 let due_at = Utc::now() - Duration::minutes(1);
4713 let automation = AutomationRecord {
4714 rrule: due_at
4715 .format("FREQ=ONCE;AT=%Y-%m-%dT%H:%M:%S+00:00")
4716 .to_string(),
4717 next_run_at: Some(due_at),
4718 created_at: due_at - Duration::minutes(5),
4719 updated_at: due_at - Duration::minutes(5),
4720 ..automation_record_with_settings(None, None, None, None)
4721 };
4722 manager
4723 .save_automation(&automation)
4724 .expect("save automation");
4725
4726 let due = manager
4727 .collect_due_runs(Utc::now())
4728 .expect("collect due runs");
4729 assert_eq!(due.len(), 1);
4730 let (observed, proposed) = &due[0];
4731 let run = manager
4732 .claim_scheduled_run(observed, proposed.clone(), tempdir.path())
4733 .expect("claim one-shot occurrence")
4734 .expect("due occurrence admitted");
4735 assert_eq!(run.scheduled_for, due_at);
4736
4737 manager
4738 .finish_scheduled_run(&run, Utc::now())
4739 .expect("finish one-shot run");
4740 let updated = manager
4741 .get_automation(&automation.id)
4742 .expect("updated automation");
4743 assert_eq!(updated.status, AutomationStatus::Paused);
4744 assert_eq!(updated.next_run_at, None);
4745 assert!(
4746 manager
4747 .collect_due_runs(Utc::now() + Duration::hours(1))
4748 .expect("later tick")
4749 .is_empty()
4750 );
4751 }
4752
4753 #[test]
4754 fn get_runs_by_ids_finds_live_runs_past_the_newest_window() {
4755 let tempdir = tempfile::tempdir().expect("tempdir");
4756 let manager =
4757 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4758 let automation = automation_record_with_settings(None, None, None, None);
4759 let base = Utc::now();
4760
4761 // The OLDEST run is the long-running task; 25 newer runs stack on
4762 // top of it while it is still Running.
4763 let long_running = run_created_at(&automation, base - Duration::minutes(60));
4764 let mut long_running = long_running;
4765 long_running.status = AutomationRunStatus::Running;
4766 long_running.started_at = Some(base - Duration::minutes(60));
4767 manager.save_run(&long_running).expect("save live run");
4768 for i in 0..25 {
4769 let mut newer = run_created_at(&automation, base - Duration::minutes(30 - i as i64));
4770 newer.status = AutomationRunStatus::Completed;
4771 newer.ended_at = Some(base - Duration::minutes(29 - i as i64));
4772 manager.save_run(&newer).expect("save newer settled run");
4773 }
4774
4775 let window = manager
4776 .list_runs(&automation.id, Some(25))
4777 .expect("windowed list");
4778 assert!(
4779 window.iter().all(|run| run.id != long_running.id),
4780 "the live run sits past the newest-25 window"
4781 );
4782
4783 let wanted: std::collections::BTreeSet<String> =
4784 [long_running.id.clone()].into_iter().collect();
4785 let found = manager
4786 .get_runs_by_ids(&automation.id, &wanted)
4787 .expect("re-read live run");
4788 assert_eq!(found.len(), 1, "the run is found wherever it sits");
4789 assert_eq!(found[0].id, long_running.id);
4790 assert_eq!(found[0].status, AutomationRunStatus::Running);
4791 }
4792
4793 #[test]
4794 fn get_runs_by_ids_ignores_non_json_noise() {
4795 let tempdir = tempfile::tempdir().expect("tempdir");
4796 let manager =
4797 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4798 let automation = automation_record_with_settings(None, None, None, None);
4799 let run = queued_run_for(&automation);
4800 manager.save_run(&run).expect("save json run");
4801
4802 let dir = manager.runs_dir_for(&automation.id).expect("runs dir");
4803 let json = dir
4804 .read_dir()
4805 .expect("list")
4806 .filter_map(|entry| entry.ok())
4807 .map(|entry| entry.path())
4808 .find(|path| path.extension().and_then(|ext| ext.to_str()) == Some("json"))
4809 .expect("json run file");
4810 let stem = json
4811 .file_stem()
4812 .and_then(|stem| stem.to_str())
4813 .expect("stem");
4814 fs::write(dir.join(format!("{stem}.tmp")), "not-json").expect("write tmp noise");
4815 fs::write(dir.join(format!("{stem}.json.bak")), "not-json").expect("write bak noise");
4816
4817 let wanted: std::collections::BTreeSet<String> = [run.id.clone()].into_iter().collect();
4818 let found = manager
4819 .get_runs_by_ids(&automation.id, &wanted)
4820 .expect("noise must not poison the re-read");
4821 assert_eq!(found.len(), 1);
4822 assert_eq!(found[0].id, run.id);
4823 }
4824
4825 #[test]
4826 fn list_runs_merges_legacy_and_sortable_files_newest_first() {
4827 let tempdir = tempfile::tempdir().expect("tempdir");
4828 let manager =
4829 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4830 let automation = automation_record_with_settings(None, None, None, None);
4831 let base = Utc::now();
4832
4833 // Legacy files sit at both ends of the timeline to prove the merge is
4834 // by created_at, not by file-name era.
4835 let legacy_oldest = run_created_at(&automation, base - Duration::minutes(30));
4836 let legacy_newest = run_created_at(&automation, base + Duration::minutes(30));
4837 write_legacy_run_file(&manager, &legacy_oldest);
4838 write_legacy_run_file(&manager, &legacy_newest);
4839
4840 let sortable_old = run_created_at(&automation, base - Duration::minutes(20));
4841 let sortable_new = run_created_at(&automation, base + Duration::minutes(20));
4842 manager.save_run(&sortable_old).expect("save old");
4843 manager.save_run(&sortable_new).expect("save new");
4844
4845 let all = manager.list_runs(&automation.id, None).expect("list all");
4846 let ids: Vec<&str> = all.iter().map(|run| run.id.as_str()).collect();
4847 assert_eq!(
4848 ids,
4849 vec![
4850 legacy_newest.id.as_str(),
4851 sortable_new.id.as_str(),
4852 sortable_old.id.as_str(),
4853 legacy_oldest.id.as_str(),
4854 ]
4855 );
4856
4857 let top_two = manager.list_runs(&automation.id, Some(2)).expect("list 2");
4858 let top_ids: Vec<&str> = top_two.iter().map(|run| run.id.as_str()).collect();
4859 assert_eq!(
4860 top_ids,
4861 vec![legacy_newest.id.as_str(), sortable_new.id.as_str()]
4862 );
4863 }
4864
4865 #[test]
4866 fn list_runs_with_limit_skips_older_sortable_files_entirely() {
4867 let tempdir = tempfile::tempdir().expect("tempdir");
4868 let manager =
4869 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4870 let automation = automation_record_with_settings(None, None, None, None);
4871 let base = Utc::now();
4872
4873 let newest = run_created_at(&automation, base);
4874 manager.save_run(&newest).expect("save newest");
4875
4876 // A corrupt sortable-named file older than the newest run: bounded
4877 // listing must never open it, while an unbounded listing fails.
4878 let dir = manager.runs_dir_for(&automation.id).expect("runs dir");
4879 let stale_stamp = run_file_stamp(base - Duration::minutes(5));
4880 fs::write(
4881 dir.join(format!("{stale_stamp}-{}.json", Uuid::new_v4())),
4882 "{ not json",
4883 )
4884 .expect("write corrupt run");
4885
4886 let bounded = manager
4887 .list_runs(&automation.id, Some(1))
4888 .expect("bounded list must not read files beyond the limit");
4889 assert_eq!(bounded.len(), 1);
4890 assert_eq!(bounded[0].id, newest.id);
4891
4892 assert!(manager.list_runs(&automation.id, None).is_err());
4893 }
4894
4895 #[tokio::test]
4896 async fn list_automations_completes_during_slow_enqueue() {
4897 let tempdir = tempfile::tempdir().expect("tempdir");
4898 let manager =
4899 AutomationManager::open_for_test(tempdir.path().to_path_buf()).expect("manager");
4900 let created = manager
4901 .create_automation(CreateAutomationRequest {
4902 name: "Slow enqueue".to_string(),
4903 prompt: "prompt".to_string(),
4904 rrule: "FREQ=HOURLY;INTERVAL=1".to_string(),
4905 cwds: Vec::new(),
4906 model: None,
4907 model_provider: None,
4908 model_provider_id: None,
4909 mode: None,
4910 allow_shell: None,
4911 trust_mode: None,
4912 auto_approve: None,
4913 delivery_mode: None,
4914 status: Some(AutomationStatus::Active),
4915 })
4916 .expect("create");
4917 let shared: SharedAutomationManager = Arc::new(Mutex::new(manager));
4918
4919 let (entered_tx, entered_rx) = tokio::sync::oneshot::channel::<()>();
4920 let (release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
4921
4922 let run_task = tokio::spawn({
4923 let shared = Arc::clone(&shared);
4924 let automation_id = created.id.clone();
4925 let task_data_dir = tempdir.path().to_path_buf();
4926 async move {
4927 run_now_with(
4928 &shared,
4929 &automation_id,
4930 &task_data_dir,
4931 move |_, mut run| async move {
4932 // Delayed task-manager stub: stall the enqueue await until
4933 // the test has proven the manager mutex is free.
4934 let _ = entered_tx.send(());
4935 let _ = release_rx.await;
4936 run.status = AutomationRunStatus::Failed;
4937 run.ended_at = Some(Utc::now());
4938 run.error = Some("stubbed enqueue".to_string());
4939 run
4940 },
4941 )
4942 .await
4943 }
4944 });
4945
4946 entered_rx.await.expect("enqueue phase entered");
4947
4948 let listed = tokio::time::timeout(std::time::Duration::from_secs(2), async {
4949 shared.lock().await.list_automations()
4950 })
4951 .await
4952 .expect("list_automations must not block behind a slow enqueue")
4953 .expect("list automations");
4954 assert_eq!(listed.len(), 1);
4955
4956 release_tx.send(()).expect("release stub");
4957 let run = run_task.await.expect("join").expect("run now");
4958 assert!(matches!(run.status, AutomationRunStatus::Failed));
4959
4960 // The final run state was persisted after the lock was reacquired.
4961 let manager = shared.lock().await;
4962 let runs = manager.list_runs(&created.id, None).expect("list runs");
4963 assert_eq!(runs.len(), 1);
4964 assert_eq!(runs[0].id, run.id);
4965 assert!(matches!(runs[0].status, AutomationRunStatus::Failed));
4966 let automation = manager.get_automation(&created.id).expect("automation");
4967 assert!(automation.last_run_at.is_some());
4968 }
4969
4970 #[tokio::test]
4971 async fn watcher_noop_completion_hides_but_retains_consumed_occurrence() -> Result<()> {
4972 let tempdir = tempfile::tempdir().expect("tempdir");
4973 let task_manager = TaskManager::start_with_executor(
4974 automation_task_config(tempdir.path().join("tasks")),
4975 std::sync::Arc::new(AutomationWatcherNoopExecutor),
4976 )
4977 .await?;
4978 let mut automation = automation_record_with_settings(None, None, None, None);
4979 automation.delivery_mode = Some(AutomationDeliveryMode::Watcher);
4980 automation.next_run_at = Some(Utc::now() - Duration::seconds(1));
4981 let manager =
4982 AutomationManager::open_for_test(tempdir.path().join("automations")).expect("manager");
4983 manager
4984 .save_automation(&automation)
4985 .expect("save automation");
4986 let shared: SharedAutomationManager = Arc::new(Mutex::new(manager));
4987
4988 scheduler_tick_shared(&shared, &task_manager).await?;
4989 let initial = shared
4990 .lock()
4991 .await
4992 .list_runs_with_visibility(&automation.id, None, true)?;
4993 assert_eq!(initial.len(), 1);
4994 let bound_id = initial[0]
4995 .task_id
4996 .as_deref()
4997 .context("watcher task binding")?;
4998 let completed = crate::task_manager::wait_for_terminal_state(
4999 &task_manager,
5000 bound_id,
5001 std::time::Duration::from_secs(5),
5002 )
5003 .await?;
5004 assert_eq!(completed.status, TaskStatus::Completed);
5005 reconcile_run_statuses_shared(&shared, &task_manager).await?;
5006
5007 let manager = shared.lock().await;
5008 assert!(
5009 manager.list_runs(&automation.id, None)?.is_empty(),
5010 "watcher no-op must not leave a phantom run row"
5011 );
5012 let durable = manager.list_runs_with_visibility(&automation.id, None, true)?;
5013 assert_eq!(durable.len(), 1);
5014 assert_eq!(durable[0].status, AutomationRunStatus::Completed);
5015 assert!(durable[0].dispatch.as_ref().unwrap().suppress_report);
5016 let mut updated = manager.get_automation(&automation.id)?;
5017 assert!(
5018 updated.next_run_at.is_some(),
5019 "watcher should keep scheduling"
5020 );
5021 assert_eq!(
5022 updated.last_run_at, None,
5023 "no-op checks are not reportable runs"
5024 );
5025 // Revisit the consumed slot after a torn schedule write. A hidden
5026 // receipt still owns its occurrence and must prevent another task.
5027 updated.next_run_at = Some(durable[0].scheduled_for);
5028 manager.save_automation(&updated)?;
5029 drop(manager);
5030 scheduler_tick_shared(&shared, &task_manager).await?;
5031 assert_eq!(task_manager.list_tasks(None).await?.len(), 1);
5032 assert!(
5033 shared
5034 .lock()
5035 .await
5036 .list_runs(&automation.id, None)?
5037 .is_empty()
5038 );
5039 task_manager.shutdown();
5040 Ok(())
5041 }
5042
5043 #[test]
5044 fn default_automations_dir_honors_codewhale_home_as_hard_override() {
5045 let _lock = crate::test_support::lock_test_env();
5046 let tmp = tempfile::TempDir::new().unwrap();
5047 // SAFETY: serialised by lock_test_env.
5048 unsafe {
5049 std::env::remove_var("DEEPSEEK_AUTOMATIONS_DIR");
5050 std::env::set_var("CODEWHALE_HOME", tmp.path());
5051 }
5052 // $CODEWHALE_HOME IS the home dir (no ".codewhale" appended); the
5053 // legacy ~/.deepseek fallback is bypassed entirely.
5054 assert_eq!(default_automations_dir(), tmp.path().join("automations"));
5055 // SAFETY: cleanup under the same lock.
5056 unsafe {
5057 std::env::remove_var("CODEWHALE_HOME");
5058 }
5059 }
5060
5061 #[test]
5062 fn default_automations_dir_prefers_deepseek_automations_dir_over_codewhale_home() {
5063 let _lock = crate::test_support::lock_test_env();
5064 let tmp = tempfile::TempDir::new().unwrap();
5065 // SAFETY: serialised by lock_test_env.
5066 unsafe {
5067 std::env::set_var("DEEPSEEK_AUTOMATIONS_DIR", tmp.path());
5068 std::env::set_var("CODEWHALE_HOME", "/should/not/be/used");
5069 }
5070 // The most-specific override wins over the base-data-dir override.
5071 assert_eq!(default_automations_dir(), tmp.path());
5072 // SAFETY: cleanup under the same lock.
5073 unsafe {
5074 std::env::remove_var("DEEPSEEK_AUTOMATIONS_DIR");
5075 std::env::remove_var("CODEWHALE_HOME");
5076 }
5077 }
5078 mod ownership;
5079 mod recovery;
5080
5081 /// #6162: the projection's Canceled branch must record why. A cooperative
5082 /// cancel arrives with no error of its own and takes the derived text; a
5083 /// task that already carries the task manager's receipt keeps it; a legacy
5084 /// record with neither degrades to no detail instead of inventing one.
5085 fn canceled_task_record(
5086 error: Option<&str>,
5087 terminal_reason: Option<&str>,
5088 ) -> crate::task_manager::TaskRecord {
5089 serde_json::from_value(serde_json::json!({
5090 "id": "task_canceled",
5091 "prompt": "nightly report",
5092 "model": "deepseek-v4-pro",
5093 "workspace": "/tmp/automation-fixture",
5094 "mode": "agent",
5095 "allow_shell": false,
5096 "trust_mode": false,
5097 "status": "canceled",
5098 "created_at": "2026-09-14T10:00:00Z",
5099 "started_at": "2026-09-14T10:00:01Z",
5100 "ended_at": "2026-09-14T10:00:05Z",
5101 "duration_ms": 4000,
5102 "result_summary": null,
5103 "result_detail_path": null,
5104 "error": error,
5105 "terminal_reason": terminal_reason,
5106 "tool_calls": [],
5107 "timeline": [],
5108 }))
5109 .expect("a canceled task record fixture")
5110 }
5111
5112 #[test]
5113 fn a_cooperatively_canceled_task_names_the_cancellation_on_the_run() {
5114 let mut run = new_run_record("auto_1", Utc::now(), Utc::now());
5115 run.status = AutomationRunStatus::Running;
5116 let task = canceled_task_record(None, Some("canceled"));
5117 assert!(apply_task_status(&mut run, &task));
5118 assert_eq!(run.status, AutomationRunStatus::Canceled);
5119 assert_eq!(run.error.as_deref(), Some("canceled by request"));
5120 assert_eq!(run.started_at, task.started_at);
5121 assert_eq!(run.ended_at, task.ended_at);
5122 }
5123
5124 #[test]
5125 fn a_canceled_task_with_its_own_error_keeps_that_error_on_the_run() {
5126 let mut run = new_run_record("auto_1", Utc::now(), Utc::now());
5127 run.status = AutomationRunStatus::Running;
5128 let task = canceled_task_record(
5129 Some("Task canceled because the task manager shut down"),
5130 Some("shutdown"),
5131 );
5132 assert!(apply_task_status(&mut run, &task));
5133 assert_eq!(run.status, AutomationRunStatus::Canceled);
5134 assert_eq!(
5135 run.error.as_deref(),
5136 Some("Task canceled because the task manager shut down")
5137 );
5138 }
5139
5140 #[test]
5141 fn a_legacy_canceled_task_without_a_reason_records_no_detail() {
5142 let mut run = new_run_record("auto_1", Utc::now(), Utc::now());
5143 run.status = AutomationRunStatus::Running;
5144 let task = canceled_task_record(None, None);
5145 assert!(apply_task_status(&mut run, &task));
5146 assert_eq!(run.status, AutomationRunStatus::Canceled);
5147 assert_eq!(run.error, None);
5148 }
5149 }
5150
5150 lines RUST