返回 CodeWhale
task_manager.rs
根目录 / crates / tui / src / task_manager.rs
1 //! Persistent background task manager for Codewhale agent work.
2 //!
3 //! Tasks are durable across restarts and execute with a bounded worker pool.
4 //! Execution uses the shared runtime provider route and links every task to
5 //! runtime thread/turn records for unified timelines.
6
7 use std::collections::{HashMap, HashSet, VecDeque};
8 use std::fs;
9 use std::path::{Path, PathBuf};
10 use std::sync::Arc;
11 #[cfg(test)]
12 use std::time::Duration as StdDuration;
13 use std::time::{Duration, Instant};
14
15 use anyhow::{Context, Result, anyhow, bail};
16 use async_trait::async_trait;
17 use chrono::{DateTime, Utc};
18 use serde::{Deserialize, Serialize};
19 use serde_json::{Value, json};
20 use tokio::sync::{Mutex, Notify, mpsc};
21 use tokio::time::sleep;
22 use tokio_util::sync::CancellationToken;
23 use uuid::Uuid;
24
25 use crate::config::Config;
26 use crate::runtime_threads::{
27 CreateThreadRequest, RUNTIME_STORE_FAILURE_EVENT, RuntimeEventRecord, RuntimeProcessOwnerLock,
28 RuntimeThreadManager, RuntimeThreadManagerConfig, RuntimeTurnStatus,
29 SharedRuntimeThreadManager, StartTurnRequest,
30 };
31 use crate::utils::spawn_supervised;
32
33 const DEFAULT_WORKERS: usize = 2;
34 const MAX_WORKERS: usize = 8;
35 const TIMELINE_SUMMARY_LIMIT: usize = 240;
36 const TIMELINE_ENTRY_LIMIT: usize = 256;
37 const TIMELINE_HEAD_KEEP: usize = 8;
38 const ARTIFACT_THRESHOLD: usize = 1200;
39 const TASK_EVENT_CHANNEL_CAPACITY: usize = 256;
40 const EVENT_CURSOR_BATCH: usize = 256;
41 const EVENT_CATCHUP_POLL: Duration = Duration::from_millis(200);
42 // v4 binds execution to a trusted Runtime scope. Older executors must not
43 // ignore its eligibility or generation fence.
44 const CURRENT_TASK_SCHEMA_VERSION: u32 = 4;
45 const STORE_REFRESH_INTERVAL: Duration = Duration::from_millis(200);
46 /// How long a worker with nothing claimable naps between queue-fingerprint
47 /// checks. An in-process notify still wakes it at once; this only bounds how
48 /// soon it notices another process's queue write, so an idle session stats the
49 /// queue about once a second per worker instead of five times (#6728, #6573).
50 const STORE_SETTLED_NAP: Duration = Duration::from_secs(1);
51 /// Retry an empty claim while queue metadata is missing or too recent to
52 /// distinguish writes on a coarse-mtime filesystem. Once settled, idle workers
53 /// wait for admission notifications or queue changes instead of reloading every
54 /// task record on a timer (#6728).
55 const STORE_IDLE_POLL_INTERVAL: Duration = Duration::from_secs(2);
56 /// Ceiling for the retry delay after a claim fails (e.g. "Task store is busy").
57 /// A failed claim has already waited out `lock_store`'s 5s deadline, so work
58 /// written by another process can wait up to about 13s after contention; an
59 /// in-process submission or a queue change ends the backoff early.
60 const STORE_BUSY_BACKOFF_MAX: Duration = Duration::from_secs(8);
61 /// A read snapshot is only trusted once both stamps are at least this old, so
62 /// a write landing in the same coarse mtime tick as the load (1s or 2s on
63 /// some filesystems) cannot hide behind an unchanged stamp (#6573).
64 const READ_SNAPSHOT_SETTLE: Duration = Duration::from_secs(2);
65
66 const fn default_task_schema_version() -> u32 {
67 CURRENT_TASK_SCHEMA_VERSION
68 }
69
70 /// Durable task status.
71 #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
72 #[serde(rename_all = "snake_case")]
73 pub enum TaskStatus {
74 Queued,
75 Running,
76 Completed,
77 Failed,
78 Canceled,
79 }
80
81 /// What the manager actually did while handling a cancellation request.
82 ///
83 /// This is returned from the same state-lock transaction as the task record,
84 /// so callers never have to infer an outcome from a stale pre-cancel read.
85 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
86 pub enum TaskCancelDisposition {
87 Forced,
88 Requested,
89 AlreadyFinished,
90 }
91
92 #[derive(Debug, Clone)]
93 pub struct TaskCancellation {
94 pub task: TaskRecord,
95 pub disposition: TaskCancelDisposition,
96 }
97
98 impl TaskStatus {
99 #[cfg(test)]
100 #[must_use]
101 pub fn is_terminal(self) -> bool {
102 matches!(self, Self::Completed | Self::Failed | Self::Canceled)
103 }
104 }
105
106 /// Why a durable task left the running state. Stored on the task record so
107 /// receipts and status views can show a forced timeout separately from a
108 /// cooperative cancel.
109 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
110 pub enum TaskTerminalReason {
111 Completed,
112 Canceled,
113 CancelTimeout,
114 Shutdown,
115 WallTimeout,
116 IdleTimeout,
117 Failed,
118 }
119
120 impl TaskTerminalReason {
121 #[must_use]
122 pub const fn as_str(self) -> &'static str {
123 match self {
124 Self::Completed => "completed",
125 Self::Canceled => "canceled",
126 Self::CancelTimeout => "cancel_timeout",
127 Self::Shutdown => "shutdown",
128 Self::WallTimeout => "wall_timeout",
129 Self::IdleTimeout => "idle_timeout",
130 Self::Failed => "failed",
131 }
132 }
133
134 #[must_use]
135 pub const fn task_status(self) -> TaskStatus {
136 match self {
137 Self::Completed => TaskStatus::Completed,
138 Self::Canceled | Self::CancelTimeout | Self::Shutdown => TaskStatus::Canceled,
139 Self::WallTimeout | Self::IdleTimeout | Self::Failed => TaskStatus::Failed,
140 }
141 }
142
143 #[must_use]
144 pub fn receipt_message(self) -> String {
145 match self {
146 Self::Completed => "Task completed".to_string(),
147 Self::Canceled => "Task canceled".to_string(),
148 Self::CancelTimeout => {
149 "Task did not terminalize after cancellation; worker released".to_string()
150 }
151 Self::Shutdown => "Task canceled because the task manager shut down".to_string(),
152 Self::WallTimeout => {
153 "Task exceeded its wall-time deadline without completing".to_string()
154 }
155 Self::IdleTimeout => {
156 "Task made no model or tool progress before the idle deadline".to_string()
157 }
158 Self::Failed => "Task ended unexpectedly".to_string(),
159 }
160 }
161
162 fn after_grace(self) -> Self {
163 match self {
164 Self::Canceled => Self::CancelTimeout,
165 other => other,
166 }
167 }
168 }
169
170 /// Durable tool-call status within a task timeline.
171 #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
172 #[serde(rename_all = "snake_case")]
173 pub enum TaskToolStatus {
174 Running,
175 Success,
176 Failed,
177 Canceled,
178 }
179
180 /// Timeline entry for a task execution.
181 #[derive(Debug, Clone, Serialize, Deserialize)]
182 pub struct TaskTimelineEntry {
183 pub timestamp: DateTime<Utc>,
184 pub kind: String,
185 pub summary: String,
186 #[serde(skip_serializing_if = "Option::is_none")]
187 pub detail_path: Option<PathBuf>,
188 }
189
190 /// Tool call summary for a task.
191 #[derive(Debug, Clone, Serialize, Deserialize)]
192 pub struct TaskToolCallSummary {
193 pub id: String,
194 pub name: String,
195 pub status: TaskToolStatus,
196 pub started_at: DateTime<Utc>,
197 pub ended_at: Option<DateTime<Utc>>,
198 pub duration_ms: Option<u64>,
199 #[serde(skip_serializing_if = "Option::is_none")]
200 pub input_summary: Option<String>,
201 #[serde(skip_serializing_if = "Option::is_none")]
202 pub output_summary: Option<String>,
203 #[serde(skip_serializing_if = "Option::is_none")]
204 pub detail_path: Option<PathBuf>,
205 #[serde(skip_serializing_if = "Option::is_none")]
206 pub patch_ref: Option<PathBuf>,
207 }
208
209 /// Checklist item stored on durable tasks. This is the durable form behind the
210 /// model-visible checklist/todo compatibility tools.
211 #[derive(Debug, Clone, Serialize, Deserialize)]
212 pub struct TaskChecklistItem {
213 pub id: u32,
214 pub content: String,
215 pub status: String,
216 }
217
218 /// Checklist state associated with a task.
219 #[derive(Debug, Clone, Serialize, Deserialize, Default)]
220 pub struct TaskChecklistState {
221 pub items: Vec<TaskChecklistItem>,
222 pub completion_pct: u8,
223 pub in_progress_id: Option<u32>,
224 pub updated_at: Option<DateTime<Utc>>,
225 }
226
227 /// Structured verification evidence attached to a task.
228 #[derive(Debug, Clone, Serialize, Deserialize)]
229 pub struct TaskGateRecord {
230 pub id: String,
231 pub gate: String,
232 pub command: String,
233 pub cwd: PathBuf,
234 pub exit_code: Option<i32>,
235 pub status: String,
236 pub classification: String,
237 pub duration_ms: u64,
238 pub summary: String,
239 #[serde(skip_serializing_if = "Option::is_none")]
240 pub log_path: Option<PathBuf>,
241 pub recorded_at: DateTime<Utc>,
242 }
243
244 /// PR-attempt metadata and artifacts attached to a task.
245 #[derive(Debug, Clone, Serialize, Deserialize)]
246 pub struct TaskAttemptRecord {
247 pub id: String,
248 pub attempt_group_id: String,
249 pub attempt_index: u32,
250 pub attempt_count: u32,
251 #[serde(skip_serializing_if = "Option::is_none")]
252 pub base_ref: Option<String>,
253 #[serde(skip_serializing_if = "Option::is_none")]
254 pub base_sha: Option<String>,
255 #[serde(skip_serializing_if = "Option::is_none")]
256 pub head_ref: Option<String>,
257 #[serde(skip_serializing_if = "Option::is_none")]
258 pub head_sha: Option<String>,
259 pub summary: String,
260 pub changed_files: Vec<String>,
261 #[serde(skip_serializing_if = "Option::is_none")]
262 pub patch_path: Option<PathBuf>,
263 pub verification: Vec<String>,
264 pub selected: bool,
265 pub recorded_at: DateTime<Utc>,
266 }
267
268 /// Durable artifact reference produced by task-aware tools.
269 #[derive(Debug, Clone, Serialize, Deserialize)]
270 pub struct TaskArtifactRef {
271 pub label: String,
272 pub path: PathBuf,
273 pub summary: String,
274 pub created_at: DateTime<Utc>,
275 }
276
277 /// GitHub write/read evidence attached to a task timeline.
278 #[derive(Debug, Clone, Serialize, Deserialize)]
279 pub struct TaskGithubEvent {
280 pub id: String,
281 pub action: String,
282 pub target: String,
283 pub number: u64,
284 pub summary: String,
285 pub url: Option<String>,
286 pub recorded_at: DateTime<Utc>,
287 }
288
289 /// Durable task record.
290 #[derive(Debug, Clone, Serialize, Deserialize)]
291 pub struct TaskRecord {
292 #[serde(default = "default_task_schema_version")]
293 pub schema_version: u32,
294 pub id: String,
295 pub prompt: String,
296 #[serde(default, skip_serializing_if = "Option::is_none")]
297 pub name: Option<String>,
298 pub model: String,
299 #[serde(default, skip_serializing_if = "Option::is_none")]
300 pub model_provider: Option<String>,
301 #[serde(default, skip_serializing_if = "Option::is_none")]
302 pub model_provider_id: Option<String>,
303 pub workspace: PathBuf,
304 pub mode: String,
305 pub allow_shell: bool,
306 pub trust_mode: bool,
307 /// A record without the field (written before it existed) is not
308 /// auto-approved: a missing grant never widens authority.
309 #[serde(default)]
310 pub auto_approve: bool,
311 /// Permission posture the task's own thread starts on (`ask`,
312 /// `auto_review`, `full_access`). Absent on records written before the
313 /// field existed and on the in-process callers that still express authority
314 /// through `auto_approve` alone; the thread request then derives the
315 /// posture from that bit exactly as it always did.
316 #[serde(default, skip_serializing_if = "Option::is_none")]
317 pub permission_posture: Option<String>,
318 pub status: TaskStatus,
319 pub created_at: DateTime<Utc>,
320 pub started_at: Option<DateTime<Utc>>,
321 pub ended_at: Option<DateTime<Utc>>,
322 pub duration_ms: Option<u64>,
323 #[serde(skip_serializing_if = "Option::is_none")]
324 pub result_summary: Option<String>,
325 #[serde(skip_serializing_if = "Option::is_none")]
326 pub result_detail_path: Option<PathBuf>,
327 #[serde(skip_serializing_if = "Option::is_none")]
328 pub error: Option<String>,
329 #[serde(default, skip_serializing_if = "Option::is_none")]
330 pub terminal_reason: Option<String>,
331 #[serde(default, skip_serializing_if = "Option::is_none")]
332 pub thread_id: Option<String>,
333 #[serde(default, skip_serializing_if = "Option::is_none")]
334 pub turn_id: Option<String>,
335 #[serde(default, skip_serializing_if = "Option::is_none")]
336 pub owner_session_id: Option<String>,
337 /// Trusted execution provenance, distinct from model-visible ownership.
338 #[serde(default, skip_serializing_if = "Option::is_none")]
339 pub execution_scope: Option<String>,
340 #[serde(default, skip_serializing_if = "Option::is_none")]
341 pub execution_generation: Option<String>,
342 /// Durable cancellation acknowledged by the actual execution owner.
343 #[serde(default)]
344 pub cancel_requested_seq: u64,
345 #[serde(default)]
346 pub runtime_event_count: usize,
347 /// Monotonic owner-lifecycle sequence used by Work Graph reconciliation.
348 /// Output/progress events do not advance this counter; only lifecycle
349 /// transitions do, so replay after restart is stable.
350 #[serde(default)]
351 pub lifecycle_seq: u64,
352 #[serde(default)]
353 pub checklist: TaskChecklistState,
354 #[serde(default)]
355 pub gates: Vec<TaskGateRecord>,
356 #[serde(default)]
357 pub attempts: Vec<TaskAttemptRecord>,
358 #[serde(default)]
359 pub artifacts: Vec<TaskArtifactRef>,
360 #[serde(default)]
361 pub github_events: Vec<TaskGithubEvent>,
362 pub tool_calls: Vec<TaskToolCallSummary>,
363 pub timeline: Vec<TaskTimelineEntry>,
364 }
365
366 /// Lightweight task view.
367 #[derive(Debug, Clone, Serialize, Deserialize)]
368 pub struct TaskSummary {
369 pub id: String,
370 pub status: TaskStatus,
371 pub prompt_summary: String,
372 #[serde(default, skip_serializing_if = "Option::is_none")]
373 pub name: Option<String>,
374 pub model: String,
375 #[serde(default, skip_serializing_if = "Option::is_none")]
376 pub model_provider: Option<String>,
377 #[serde(default, skip_serializing_if = "Option::is_none")]
378 pub model_provider_id: Option<String>,
379 pub mode: String,
380 pub workspace: PathBuf,
381 pub created_at: DateTime<Utc>,
382 pub started_at: Option<DateTime<Utc>>,
383 pub ended_at: Option<DateTime<Utc>>,
384 pub duration_ms: Option<u64>,
385 #[serde(default)]
386 pub lifecycle_seq: u64,
387 #[serde(skip_serializing_if = "Option::is_none")]
388 pub error: Option<String>,
389 #[serde(default, skip_serializing_if = "Option::is_none")]
390 pub terminal_reason: Option<String>,
391 #[serde(default, skip_serializing_if = "Option::is_none")]
392 pub thread_id: Option<String>,
393 #[serde(default, skip_serializing_if = "Option::is_none")]
394 pub turn_id: Option<String>,
395 #[serde(default, skip_serializing_if = "Option::is_none")]
396 pub owner_session_id: Option<String>,
397 pub execution_binding_known: bool,
398 }
399
400 impl From<&TaskRecord> for TaskSummary {
401 fn from(value: &TaskRecord) -> Self {
402 Self {
403 id: value.id.clone(),
404 status: value.status,
405 prompt_summary: summarize_text(&value.prompt, TIMELINE_SUMMARY_LIMIT),
406 name: value.name.clone(),
407 model: value.model.clone(),
408 model_provider: value.model_provider.clone(),
409 model_provider_id: value.model_provider_id.clone(),
410 mode: value.mode.clone(),
411 workspace: value.workspace.clone(),
412 created_at: value.created_at,
413 started_at: value.started_at,
414 ended_at: value.ended_at,
415 duration_ms: value.duration_ms,
416 lifecycle_seq: value.lifecycle_seq,
417 error: value.error.clone(),
418 terminal_reason: value.terminal_reason.clone(),
419 thread_id: value.thread_id.clone(),
420 turn_id: value.turn_id.clone(),
421 owner_session_id: value.owner_session_id.clone(),
422 execution_binding_known: value.execution_scope.is_some()
423 && (value.status != TaskStatus::Running || value.execution_generation.is_some()),
424 }
425 }
426 }
427
428 /// Count totals by status for task dashboards.
429 #[derive(Debug, Clone, Copy, Serialize, Deserialize, Default)]
430 pub struct TaskCounts {
431 pub queued: usize,
432 pub running: usize,
433 pub completed: usize,
434 pub failed: usize,
435 pub canceled: usize,
436 }
437
438 /// Request to enqueue a new task.
439 #[derive(Debug, Clone, Serialize, Deserialize)]
440 pub struct NewTaskRequest {
441 pub prompt: String,
442 /// Caller-given run name, stored as-given. Absent names stay absent —
443 /// titles derived from the prompt are a presentation concern.
444 #[serde(default)]
445 pub name: Option<String>,
446 pub model: Option<String>,
447 #[serde(default)]
448 pub model_provider: Option<String>,
449 #[serde(default)]
450 pub model_provider_id: Option<String>,
451 pub workspace: Option<PathBuf>,
452 pub mode: Option<String>,
453 pub allow_shell: Option<bool>,
454 pub trust_mode: Option<bool>,
455 pub auto_approve: Option<bool>,
456 /// Posture for the thread this task runs on. Takes precedence over
457 /// `auto_approve`, which the runtime only reads when no posture is given.
458 pub permission_posture: Option<String>,
459 pub owner_session_id: Option<String>,
460 }
461
462 impl NewTaskRequest {
463 /// Preserve values already resolved into a staged or accepted task.
464 pub(crate) fn from_task(task: &TaskRecord) -> Self {
465 Self {
466 prompt: task.prompt.clone(),
467 name: task.name.clone(),
468 model: Some(task.model.clone()),
469 model_provider: task.model_provider.clone(),
470 model_provider_id: task.model_provider_id.clone(),
471 workspace: Some(task.workspace.clone()),
472 mode: Some(task.mode.clone()),
473 allow_shell: Some(task.allow_shell),
474 trust_mode: Some(task.trust_mode),
475 auto_approve: Some(task.auto_approve),
476 permission_posture: task.permission_posture.clone(),
477 owner_session_id: task.owner_session_id.clone(),
478 }
479 }
480
481 #[cfg(test)]
482 #[must_use]
483 pub fn from_prompt(prompt: impl Into<String>) -> Self {
484 Self {
485 prompt: prompt.into(),
486 name: None,
487 model: None,
488 model_provider: None,
489 model_provider_id: None,
490 workspace: None,
491 mode: None,
492 allow_shell: None,
493 trust_mode: None,
494 auto_approve: Some(true),
495 permission_posture: None,
496 owner_session_id: None,
497 }
498 }
499 }
500
501 /// Task manager startup options.
502 #[derive(Debug, Clone)]
503 pub struct TaskManagerConfig {
504 pub data_dir: PathBuf,
505 pub worker_count: usize,
506 pub default_workspace: PathBuf,
507 pub default_model: String,
508 pub default_mode: String,
509 pub allow_shell: bool,
510 pub trust_mode: bool,
511 pub execution_limits: TaskExecutionLimits,
512 }
513
514 /// Deadlines and persistence cadence for one durable execution.
515 #[derive(Debug, Clone, Copy)]
516 pub struct TaskExecutionLimits {
517 pub wall_time: Duration,
518 pub idle_progress: Duration,
519 pub cancel_grace: Duration,
520 pub persist_debounce: Duration,
521 }
522
523 impl Default for TaskExecutionLimits {
524 fn default() -> Self {
525 Self {
526 wall_time: Duration::from_secs(30 * 60),
527 idle_progress: Duration::from_secs(2 * 60),
528 cancel_grace: Duration::from_secs(5),
529 persist_debounce: Duration::from_millis(250),
530 }
531 }
532 }
533
534 #[cfg(test)]
535 impl TaskExecutionLimits {
536 fn short_for_tests() -> Self {
537 Self {
538 wall_time: Duration::from_millis(400),
539 idle_progress: Duration::from_millis(150),
540 cancel_grace: Duration::from_millis(50),
541 persist_debounce: Duration::from_millis(10),
542 }
543 }
544 }
545
546 /// Pure watchdog for cancel grace, wall time, and idle progress.
547 struct ExecutionGuard {
548 started_at: Instant,
549 last_progress_at: Instant,
550 interrupt_at: Option<Instant>,
551 interrupt_reason: Option<TaskTerminalReason>,
552 limits: TaskExecutionLimits,
553 }
554
555 #[derive(Debug)]
556 enum GuardAction {
557 Run { wait: Duration },
558 Interrupt { reason: TaskTerminalReason },
559 Terminalize { reason: TaskTerminalReason },
560 }
561
562 impl ExecutionGuard {
563 fn new(limits: TaskExecutionLimits, now: Instant) -> Self {
564 Self {
565 started_at: now,
566 last_progress_at: now,
567 interrupt_at: None,
568 interrupt_reason: None,
569 limits,
570 }
571 }
572
573 fn note_progress(&mut self, now: Instant) {
574 self.last_progress_at = now;
575 }
576
577 fn note_interrupt(&mut self, now: Instant, reason: TaskTerminalReason) {
578 if self.interrupt_at.is_none() {
579 self.interrupt_at = Some(now);
580 self.interrupt_reason = Some(reason);
581 }
582 }
583
584 fn evaluate(&self, now: Instant, cancel: bool, shutdown: bool) -> GuardAction {
585 if let Some(interrupt_at) = self.interrupt_at {
586 let elapsed = now.saturating_duration_since(interrupt_at);
587 if elapsed >= self.limits.cancel_grace {
588 let reason = self
589 .interrupt_reason
590 .unwrap_or(TaskTerminalReason::CancelTimeout)
591 .after_grace();
592 return GuardAction::Terminalize { reason };
593 }
594 return GuardAction::Run {
595 wait: self.limits.cancel_grace.saturating_sub(elapsed),
596 };
597 }
598
599 let wall_elapsed = now.saturating_duration_since(self.started_at);
600 let idle_elapsed = now.saturating_duration_since(self.last_progress_at);
601 // A limit whose deadline does not fit in `Instant` can never fire.
602 let wall_deadline = self.started_at.checked_add(self.limits.wall_time);
603 let idle_deadline = self.last_progress_at.checked_add(self.limits.idle_progress);
604 let pending = if shutdown {
605 Some(TaskTerminalReason::Shutdown)
606 } else if cancel {
607 Some(TaskTerminalReason::Canceled)
608 } else {
609 // Attribute the timeout to the limit that was crossed first, not
610 // to the one this tick happens to check first. When the watchdog
611 // is starved past both deadlines (a >=250 ms scheduler stall on a
612 // loaded CI runner is enough with the test budgets), the idle
613 // limit that expired earlier is still the truthful reason; a tie
614 // keeps the wall limit's precedence (issue #5898).
615 match (wall_deadline, idle_deadline) {
616 (Some(wall), Some(idle)) if now >= wall && wall <= idle => {
617 Some(TaskTerminalReason::WallTimeout)
618 }
619 (_, Some(idle)) if now >= idle => Some(TaskTerminalReason::IdleTimeout),
620 (Some(wall), _) if now >= wall => Some(TaskTerminalReason::WallTimeout),
621 _ => None,
622 }
623 };
624 if let Some(reason) = pending {
625 return GuardAction::Interrupt { reason };
626 }
627
628 let wait = self
629 .limits
630 .wall_time
631 .saturating_sub(wall_elapsed)
632 .min(self.limits.idle_progress.saturating_sub(idle_elapsed))
633 .min(EVENT_CATCHUP_POLL);
634 GuardAction::Run {
635 wait: wait.max(Duration::from_millis(1)),
636 }
637 }
638
639 fn preserve_timeout_reason(&self, result: TaskExecutionResult) -> TaskExecutionResult {
640 if result.terminal_reason != TaskTerminalReason::Canceled {
641 return result;
642 }
643 match self.interrupt_reason {
644 Some(reason @ (TaskTerminalReason::WallTimeout | TaskTerminalReason::IdleTimeout)) => {
645 TaskExecutionResult::from_reason(reason, result.result_text)
646 }
647 _ => result,
648 }
649 }
650 }
651
652 impl TaskManagerConfig {
653 #[must_use]
654 pub fn from_runtime(
655 config: &Config,
656 workspace: PathBuf,
657 default_model: Option<String>,
658 worker_count: Option<usize>,
659 ) -> Self {
660 Self {
661 data_dir: default_tasks_dir(),
662 worker_count: worker_count.unwrap_or(DEFAULT_WORKERS),
663 default_workspace: workspace,
664 default_model: default_model.unwrap_or_else(|| config.default_model()),
665 default_mode: "agent".to_string(),
666 allow_shell: config.allow_shell(),
667 trust_mode: false,
668 execution_limits: TaskExecutionLimits::default(),
669 }
670 }
671 }
672
673 #[derive(Debug, Clone)]
674 pub struct ExecutionTask {
675 id: String,
676 prompt: String,
677 model: String,
678 model_provider: Option<String>,
679 model_provider_id: Option<String>,
680 workspace: PathBuf,
681 mode_label: String,
682 allow_shell: bool,
683 trust_mode: bool,
684 auto_approve: bool,
685 permission_posture: Option<String>,
686 }
687
688 impl From<&TaskRecord> for ExecutionTask {
689 fn from(task: &TaskRecord) -> Self {
690 Self {
691 id: task.id.clone(),
692 prompt: task.prompt.clone(),
693 model: task.model.clone(),
694 model_provider: task.model_provider.clone(),
695 model_provider_id: task.model_provider_id.clone(),
696 workspace: task.workspace.clone(),
697 mode_label: task.mode.clone(),
698 allow_shell: task.allow_shell,
699 trust_mode: task.trust_mode,
700 auto_approve: task.auto_approve,
701 permission_posture: task.permission_posture.clone(),
702 }
703 }
704 }
705
706 impl ExecutionTask {
707 pub(crate) fn thread_request(&self) -> CreateThreadRequest {
708 CreateThreadRequest {
709 model: Some(self.model.clone()),
710 model_provider: self.model_provider.clone(),
711 model_provider_id: self.model_provider_id.clone(),
712 workspace: Some(self.workspace.clone()),
713 mode: Some(self.mode_label.clone()),
714 allow_shell: Some(self.allow_shell),
715 trust_mode: Some(self.trust_mode),
716 auto_approve: Some(self.auto_approve),
717 permission_posture: self.permission_posture.clone(),
718 task_id: Some(self.id.clone()),
719 ..Default::default()
720 }
721 }
722
723 /// The turn request for this task. A task that carries a pinned posture
724 /// runs its turn under that posture: the legacy `auto_approve` bit is only
725 /// sent for records that predate the pinned posture, because a per-turn
726 /// `auto_approve` without a posture re-derives the permission from the
727 /// bit alone and would override what the thread was created with.
728 pub(crate) fn turn_request(&self) -> StartTurnRequest {
729 let pinned = self.permission_posture.is_some();
730 StartTurnRequest {
731 prompt: self.prompt.clone(),
732 input_summary: Some(summarize_text(&self.prompt, TIMELINE_SUMMARY_LIMIT)),
733 model: Some(self.model.clone()),
734 mode: Some(self.mode_label.clone()),
735 permission_posture: self.permission_posture.clone(),
736 allow_shell: Some(self.allow_shell),
737 trust_mode: Some(self.trust_mode),
738 auto_approve: (!pinned).then_some(self.auto_approve),
739 ..Default::default()
740 }
741 }
742 }
743
744 /// Event stream produced by an executor while a task runs.
745 #[derive(Debug, Clone)]
746 pub enum TaskExecutionEvent {
747 ThreadLinked {
748 thread_id: String,
749 turn_id: String,
750 },
751 Status {
752 message: String,
753 },
754 MessageDelta {
755 content: String,
756 },
757 ToolStarted {
758 id: String,
759 name: String,
760 input: Value,
761 },
762 ToolProgress {
763 id: String,
764 output: String,
765 },
766 /// Emitted while a journal-tracked tool is in flight but the journal is
767 /// silent, so the worker supervisor's idle watchdog counts the window as
768 /// progress instead of idle. Supervisor-side liveness signal only: never
769 /// persisted and never shown on the task timeline.
770 ToolHeartbeat,
771 ToolCompleted {
772 id: String,
773 name: String,
774 success: bool,
775 output: String,
776 metadata: Option<Value>,
777 },
778 Error {
779 message: String,
780 },
781 RuntimeEvent {
782 seq: u64,
783 event: String,
784 summary: String,
785 },
786 }
787
788 /// Final executor result.
789 #[derive(Debug, Clone)]
790 pub struct TaskExecutionResult {
791 pub status: TaskStatus,
792 pub result_text: Option<String>,
793 pub error: Option<String>,
794 pub terminal_reason: TaskTerminalReason,
795 }
796
797 impl TaskExecutionResult {
798 fn failed(error: impl Into<String>) -> Self {
799 Self {
800 status: TaskStatus::Failed,
801 result_text: None,
802 error: Some(error.into()),
803 terminal_reason: TaskTerminalReason::Failed,
804 }
805 }
806
807 fn from_reason(reason: TaskTerminalReason, result_text: Option<String>) -> Self {
808 let error = match reason {
809 TaskTerminalReason::Completed | TaskTerminalReason::Canceled => None,
810 _ => Some(reason.receipt_message()),
811 };
812 Self {
813 status: reason.task_status(),
814 result_text,
815 error,
816 terminal_reason: reason,
817 }
818 }
819 }
820
821 /// Abstraction for task execution.
822 #[async_trait]
823 pub trait TaskExecutor: Send + Sync {
824 async fn execute(
825 &self,
826 task: ExecutionTask,
827 events: mpsc::Sender<TaskExecutionEvent>,
828 cancel: CancellationToken,
829 ) -> TaskExecutionResult;
830 }
831
832 /// Executor backed by the shared runtime and its canonical provider resolver.
833 pub struct EngineTaskExecutor {
834 runtime_threads: SharedRuntimeThreadManager,
835 limits: TaskExecutionLimits,
836 }
837
838 impl EngineTaskExecutor {
839 #[must_use]
840 pub fn new(runtime_threads: SharedRuntimeThreadManager, limits: TaskExecutionLimits) -> Self {
841 Self {
842 runtime_threads,
843 limits,
844 }
845 }
846 }
847
848 #[async_trait]
849 impl TaskExecutor for EngineTaskExecutor {
850 async fn execute(
851 &self,
852 task: ExecutionTask,
853 events: mpsc::Sender<TaskExecutionEvent>,
854 cancel: CancellationToken,
855 ) -> TaskExecutionResult {
856 if cancel.is_cancelled() {
857 return TaskExecutionResult::from_reason(TaskTerminalReason::Canceled, None);
858 }
859 let thread = match self
860 .runtime_threads
861 .create_thread(task.thread_request())
862 .await
863 {
864 Ok(thread) => thread,
865 Err(err) => {
866 return TaskExecutionResult::failed(format!(
867 "Failed to create runtime thread: {err}"
868 ));
869 }
870 };
871
872 if cancel.is_cancelled() {
873 return TaskExecutionResult::from_reason(TaskTerminalReason::Canceled, None);
874 }
875 let turn = match self
876 .runtime_threads
877 .start_turn(&thread.id, task.turn_request())
878 .await
879 {
880 Ok(turn) => turn,
881 Err(err) => {
882 return TaskExecutionResult::failed(format!("Failed to start task: {err}"));
883 }
884 };
885
886 emit_task_event(
887 &events,
888 TaskExecutionEvent::ThreadLinked {
889 thread_id: thread.id.clone(),
890 turn_id: turn.id.clone(),
891 },
892 )
893 .await;
894 emit_task_event(
895 &events,
896 TaskExecutionEvent::Status {
897 message: format!("Task {} started", task.id),
898 },
899 )
900 .await;
901
902 drive_engine_turn(
903 self.runtime_threads.as_ref(),
904 &thread.id,
905 &turn.id,
906 events,
907 cancel,
908 self.limits,
909 )
910 .await
911 }
912 }
913
914 async fn drive_engine_turn(
915 runtime_threads: &RuntimeThreadManager,
916 thread_id: &str,
917 turn_id: &str,
918 events: mpsc::Sender<TaskExecutionEvent>,
919 cancel: CancellationToken,
920 limits: TaskExecutionLimits,
921 ) -> TaskExecutionResult {
922 let mut subscription = runtime_threads.subscribe_events();
923 let mut guard = ExecutionGuard::new(limits, Instant::now());
924 let mut final_text = RuntimeTaskOutput::default();
925 let mut cursor = 0u64;
926 let mut terminal_status: Option<RuntimeTurnStatus> = None;
927 let mut terminal_error: Option<String> = None;
928 // Approval requests this turn is waiting on, each with the deadline the
929 // runtime bridge will resolve it by (#6118).
930 let mut pending_approvals: HashMap<String, Instant> = HashMap::new();
931 // Journal ids of tool items currently in flight. The journal records a
932 // tool call only at start/completion — nothing in between — so the
933 // idle-progress deadline would run unopposed during any silent build,
934 // test suite, or MCP call and kill healthy work at the idle limit. A
935 // tool that is still running IS progress; the wall-time budget remains
936 // the backstop for a hung one (a genuinely stuck tool also hits its own
937 // execution timeout long before a 30-minute wall budget in the common
938 // case). Both lifecycle edges are keyed by the journal item id:
939 // item.started also carries the engine tool-use id under "tool", but
940 // the terminal events only carry "item", so keying the insert by
941 // tool-use id would never drain the set.
942 let mut running_tools: HashSet<String> = HashSet::new();
943
944 loop {
945 let batch = match runtime_threads
946 .events_from_offset_async(thread_id, cursor, Some(EVENT_CURSOR_BATCH))
947 .await
948 {
949 Ok((batch, next_cursor)) => {
950 cursor = next_cursor;
951 batch
952 }
953 Err(err) => {
954 return TaskExecutionResult {
955 status: TaskStatus::Failed,
956 result_text: final_text.into_result(true),
957 error: Some(format!("Failed to read runtime events: {err}")),
958 terminal_reason: TaskTerminalReason::Failed,
959 };
960 }
961 };
962
963 let more_pending = batch.len() >= EVENT_CURSOR_BATCH;
964 for event in batch {
965 if event.thread_id != thread_id {
966 continue;
967 }
968 if event
969 .turn_id
970 .as_deref()
971 .is_some_and(|event_turn| event_turn != turn_id)
972 {
973 continue;
974 }
975 match event.event.as_str() {
976 "item.started" => {
977 // Journal tool items carry a `tool` payload key;
978 // non-tool items do not. A silent context compaction is
979 // therefore invisible to this set (disclosed gap: a
980 // background task mid-compaction can still hit the idle
981 // deadline — compaction has no heartbeat of its own).
982 if event.payload.get("tool").is_some()
983 && let Some(item_id) = journal_item_id(&event)
984 {
985 running_tools.insert(item_id);
986 }
987 }
988 "item.completed" | "item.failed" => {
989 if let Some(item_id) = journal_item_id(&event) {
990 running_tools.remove(&item_id);
991 }
992 }
993 _ => {}
994 }
995 match event.event.as_str() {
996 // An approval parks the turn on an external decision until
997 // the runtime bridge answers or its own window closes; note
998 // that deadline so the idle watchdog stays off it (#6118).
999 "approval.required" => {
1000 let approval_id = event
1001 .payload
1002 .get("approval_id")
1003 .and_then(Value::as_str)
1004 .map(ToString::to_string)
1005 .unwrap_or_else(|| format!("approval-{}", event.seq));
1006 let until = match runtime_threads.approval_decision_timeout() {
1007 Some(wait) => Instant::now() + wait,
1008 // `0` waits indefinitely by configuration; the wall
1009 // deadline still bounds the run.
1010 None => Instant::now() + limits.wall_time,
1011 };
1012 pending_approvals.insert(approval_id, until);
1013 }
1014 "approval.decided" => {
1015 if let Some(approval_id) =
1016 event.payload.get("approval_id").and_then(Value::as_str)
1017 {
1018 pending_approvals.remove(approval_id);
1019 }
1020 }
1021 _ => {}
1022 }
1023 if runtime_event_is_progress(&event) {
1024 guard.note_progress(Instant::now());
1025 }
1026 if let Some((status, error)) =
1027 ingest_runtime_event(&event, &mut final_text, &events).await
1028 {
1029 // The decision window closed on a pending approval: the
1030 // runtime already denied the tool, and an unattended run has
1031 // no operator to answer, so stop the turn instead of letting
1032 // it keep burning under a failure nobody sees (#6118).
1033 if event.event.as_str() == "approval.timeout" {
1034 let _ = runtime_threads.interrupt_turn(thread_id, turn_id).await;
1035 }
1036 terminal_status = Some(status);
1037 terminal_error = error;
1038 }
1039 }
1040
1041 if terminal_status.is_some() {
1042 break;
1043 }
1044
1045 // While an approval is pending the turn is deliberately waiting on an
1046 // external decision, not drifting: keep the idle deadline from firing
1047 // so the bridge's own window can resolve and record it. Entries expire
1048 // with their window, so a decision that never arrives cannot suspend
1049 // the watchdog forever (#6118).
1050 let now = Instant::now();
1051 pending_approvals.retain(|_, until| now < *until);
1052 if !pending_approvals.is_empty() {
1053 guard.note_progress(now);
1054 }
1055
1056 // Steady-state progress refresh: the journal is silent for the whole
1057 // execution window of an in-flight tool, so the idle check must be
1058 // suppressed for as long as any tool is running — not only on event
1059 // arrival. The loop below wakes at least every EVENT_CATCHUP_POLL,
1060 // so this runs throughout silent builds, test suites, and MCP calls.
1061 if !running_tools.is_empty() {
1062 guard.note_progress(now);
1063 // The worker supervisor (run_task) feeds its own guard from the
1064 // task-event channel and cannot see this journal-derived set, so
1065 // without a heartbeat its idle deadline would still fire on the
1066 // only production path that runs this loop. Publish the in-flight
1067 // window on that channel; the heartbeat refreshes the supervisor's
1068 // idle clock and is deliberately not persisted or surfaced.
1069 emit_task_event(&events, TaskExecutionEvent::ToolHeartbeat).await;
1070 }
1071
1072 match guard.evaluate(now, cancel.is_cancelled(), false) {
1073 GuardAction::Interrupt { reason } => {
1074 let _ = runtime_threads.interrupt_turn(thread_id, turn_id).await;
1075 emit_task_event(
1076 &events,
1077 TaskExecutionEvent::Status {
1078 message: reason.receipt_message(),
1079 },
1080 )
1081 .await;
1082 guard.note_interrupt(Instant::now(), reason);
1083 }
1084 GuardAction::Terminalize { reason } => {
1085 return TaskExecutionResult::from_reason(reason, final_text.into_result(true));
1086 }
1087 GuardAction::Run { wait } => {
1088 if more_pending {
1089 continue;
1090 }
1091 tokio::select! {
1092 _ = cancel.cancelled(), if !cancel.is_cancelled() => {}
1093 _ = subscription.recv() => {}
1094 _ = sleep(wait) => {}
1095 }
1096 }
1097 }
1098 }
1099
1100 let result = match terminal_status.unwrap_or(RuntimeTurnStatus::Failed) {
1101 RuntimeTurnStatus::Completed => TaskExecutionResult {
1102 status: TaskStatus::Completed,
1103 result_text: final_text.into_result(false),
1104 error: None,
1105 terminal_reason: TaskTerminalReason::Completed,
1106 },
1107 RuntimeTurnStatus::Interrupted | RuntimeTurnStatus::Canceled => TaskExecutionResult {
1108 status: TaskStatus::Canceled,
1109 result_text: final_text.into_result(true),
1110 error: None,
1111 terminal_reason: TaskTerminalReason::Canceled,
1112 },
1113 RuntimeTurnStatus::Queued | RuntimeTurnStatus::InProgress | RuntimeTurnStatus::Failed => {
1114 TaskExecutionResult {
1115 status: TaskStatus::Failed,
1116 result_text: final_text.into_result(true),
1117 error: terminal_error
1118 .or_else(|| Some(TaskTerminalReason::Failed.receipt_message())),
1119 terminal_reason: TaskTerminalReason::Failed,
1120 }
1121 }
1122 };
1123 guard.preserve_timeout_reason(result)
1124 }
1125
1126 async fn emit_task_event(events: &mpsc::Sender<TaskExecutionEvent>, event: TaskExecutionEvent) {
1127 let _ = events.send(event).await;
1128 }
1129
1130 fn optional_nonzero_text(text: String) -> Option<String> {
1131 if text.trim().is_empty() {
1132 None
1133 } else {
1134 Some(text)
1135 }
1136 }
1137
1138 /// A task result is the last message, not concatenated progress commentary.
1139 /// Only interrupted/failed execution may return a still-streaming message.
1140 #[derive(Default)]
1141 struct RuntimeTaskOutput {
1142 text: String,
1143 completed: bool,
1144 }
1145
1146 impl RuntimeTaskOutput {
1147 fn into_result(self, allow_partial: bool) -> Option<String> {
1148 (self.completed || allow_partial)
1149 .then(|| optional_nonzero_text(self.text))
1150 .flatten()
1151 }
1152 }
1153
1154 fn append_message_delta(result_text: &mut String, event: &TaskExecutionEvent) {
1155 if let TaskExecutionEvent::MessageDelta { content } = event {
1156 result_text.push_str(content);
1157 }
1158 }
1159
1160 fn runtime_event_is_progress(event: &RuntimeEventRecord) -> bool {
1161 matches!(
1162 event.event.as_str(),
1163 "item.delta" | "item.started" | "item.completed" | "item.failed" | "turn.completed"
1164 )
1165 }
1166
1167 /// The journal receipt id of the item a lifecycle event carries. Every
1168 /// lifecycle edge (`item.started` / `item.completed` / `item.failed`) repeats
1169 /// it under `payload.item.id`, which is what makes it a stable tracking key
1170 /// across both edges of one tool call.
1171 fn journal_item_id(event: &RuntimeEventRecord) -> Option<String> {
1172 event
1173 .payload
1174 .get("item")
1175 .and_then(|item| item.get("id"))
1176 .and_then(Value::as_str)
1177 .map(str::to_string)
1178 }
1179
1180 async fn ingest_runtime_event(
1181 event: &RuntimeEventRecord,
1182 final_text: &mut RuntimeTaskOutput,
1183 events: &mpsc::Sender<TaskExecutionEvent>,
1184 ) -> Option<(RuntimeTurnStatus, Option<String>)> {
1185 emit_task_event(
1186 events,
1187 TaskExecutionEvent::RuntimeEvent {
1188 seq: event.seq,
1189 event: event.event.clone(),
1190 summary: summarize_text(&event.payload.to_string(), TIMELINE_SUMMARY_LIMIT),
1191 },
1192 )
1193 .await;
1194
1195 match event.event.as_str() {
1196 "item.delta" => {
1197 let kind = event
1198 .payload
1199 .get("kind")
1200 .and_then(Value::as_str)
1201 .unwrap_or_default();
1202 if kind == "agent_message" {
1203 if let Some(content) = event.payload.get("delta").and_then(Value::as_str) {
1204 final_text.text.push_str(content);
1205 final_text.completed = false;
1206 emit_task_event(
1207 events,
1208 TaskExecutionEvent::MessageDelta {
1209 content: content.to_string(),
1210 },
1211 )
1212 .await;
1213 }
1214 } else if kind == "tool_call" {
1215 let output = event
1216 .payload
1217 .get("delta")
1218 .and_then(Value::as_str)
1219 .unwrap_or_default()
1220 .to_string();
1221 emit_task_event(
1222 events,
1223 TaskExecutionEvent::ToolProgress {
1224 id: event.item_id.clone().unwrap_or_default(),
1225 output,
1226 },
1227 )
1228 .await;
1229 }
1230 None
1231 }
1232 "item.started" => {
1233 if event.payload.pointer("/item/kind").and_then(Value::as_str) == Some("agent_message")
1234 {
1235 *final_text = RuntimeTaskOutput::default();
1236 }
1237 if let Some(tool) = event.payload.get("tool") {
1238 let id = tool
1239 .get("id")
1240 .and_then(Value::as_str)
1241 .unwrap_or_default()
1242 .to_string();
1243 let name = tool
1244 .get("name")
1245 .and_then(Value::as_str)
1246 .unwrap_or_default()
1247 .to_string();
1248 let input = tool.get("input").cloned().unwrap_or_else(|| json!({}));
1249 emit_task_event(events, TaskExecutionEvent::ToolStarted { id, name, input }).await;
1250 }
1251 None
1252 }
1253 "item.completed" | "item.failed" => {
1254 if let Some(item) = event.payload.get("item") {
1255 let kind = item.get("kind").and_then(Value::as_str).unwrap_or_default();
1256 if kind == "tool_call" || kind == "file_change" || kind == "command_execution" {
1257 let metadata = item.get("metadata");
1258 // Starts carry the call execution ID; item.id is Runtime's
1259 // separate receipt ID. Runtime preserves the call identity
1260 // in terminal metadata, including errors and redacted input.
1261 let id = metadata
1262 .and_then(|meta| {
1263 meta.get("tool_result_for")
1264 .or_else(|| meta.get("tool_use_id"))
1265 .or_else(|| meta.get("tool_call_id"))
1266 })
1267 .and_then(Value::as_str)
1268 .filter(|id| !id.is_empty())
1269 .or_else(|| item.get("id").and_then(Value::as_str))
1270 .unwrap_or_default()
1271 .to_string();
1272 let name = metadata
1273 .and_then(|meta| meta.get("tool_name"))
1274 .and_then(Value::as_str)
1275 .filter(|name| !name.is_empty())
1276 .unwrap_or_else(|| {
1277 item.get("summary")
1278 .and_then(Value::as_str)
1279 .unwrap_or("tool")
1280 .split(':')
1281 .next()
1282 .unwrap_or("tool")
1283 .trim()
1284 })
1285 .to_string();
1286 let output = item
1287 .get("detail")
1288 .and_then(Value::as_str)
1289 .unwrap_or_default()
1290 .to_string();
1291 emit_task_event(
1292 events,
1293 TaskExecutionEvent::ToolCompleted {
1294 id,
1295 name,
1296 success: event.event == "item.completed",
1297 output,
1298 metadata: metadata.cloned(),
1299 },
1300 )
1301 .await;
1302 } else if kind == "agent_message" {
1303 // The completed item is authoritative even when catch-up
1304 // did not receive its deltas. Replacing avoids duplication.
1305 final_text.text = item
1306 .get("detail")
1307 .and_then(Value::as_str)
1308 .or_else(|| item.get("summary").and_then(Value::as_str))
1309 .unwrap_or_default()
1310 .to_string();
1311 final_text.completed = event.event == "item.completed";
1312 } else if kind == "status" {
1313 let message = item
1314 .get("detail")
1315 .and_then(Value::as_str)
1316 .or_else(|| item.get("summary").and_then(Value::as_str))
1317 .unwrap_or_default()
1318 .to_string();
1319 emit_task_event(events, TaskExecutionEvent::Status { message }).await;
1320 } else if kind == "error" {
1321 let message = item
1322 .get("detail")
1323 .and_then(Value::as_str)
1324 .or_else(|| item.get("summary").and_then(Value::as_str))
1325 .unwrap_or_default()
1326 .to_string();
1327 emit_task_event(events, TaskExecutionEvent::Error { message }).await;
1328 }
1329 }
1330 None
1331 }
1332 "turn.completed" => {
1333 if let Some(turn_payload) = event.payload.get("turn") {
1334 let status = turn_payload
1335 .get("status")
1336 .and_then(Value::as_str)
1337 .unwrap_or("failed");
1338 let terminal_status = match status {
1339 "completed" => RuntimeTurnStatus::Completed,
1340 "interrupted" => RuntimeTurnStatus::Interrupted,
1341 "canceled" => RuntimeTurnStatus::Canceled,
1342 _ => RuntimeTurnStatus::Failed,
1343 };
1344 let terminal_error = turn_payload
1345 .get("error")
1346 .and_then(Value::as_str)
1347 .map(ToString::to_string);
1348 Some((terminal_status, terminal_error))
1349 } else {
1350 Some((RuntimeTurnStatus::Completed, None))
1351 }
1352 }
1353 RUNTIME_STORE_FAILURE_EVENT => {
1354 // The runtime's own store failed. The notice names the file and
1355 // the next action; `terminal` means no `turn.completed` can
1356 // follow, so the driver stops waiting instead of idling out (#5931).
1357 let message = event
1358 .payload
1359 .get("message")
1360 .and_then(Value::as_str)
1361 .unwrap_or("Session runtime store failure")
1362 .to_string();
1363 emit_task_event(
1364 events,
1365 TaskExecutionEvent::Error {
1366 message: message.clone(),
1367 },
1368 )
1369 .await;
1370 event
1371 .payload
1372 .get("terminal")
1373 .and_then(Value::as_bool)
1374 .unwrap_or(false)
1375 .then_some((RuntimeTurnStatus::Failed, Some(message)))
1376 }
1377
1378 "approval.timeout" => Some((
1379 RuntimeTurnStatus::Failed,
1380 Some(
1381 "Tool approval was not answered within the decision window; the runtime denied the tool and the run stopped."
1382 .to_string(),
1383 ),
1384 )),
1385 _ => None,
1386 }
1387 }
1388
1389 /// Thread-safe task manager.
1390 pub type SharedTaskManager = Arc<TaskManager>;
1391
1392 pub(crate) struct TaskManagerShutdownGuard(std::sync::Weak<TaskManager>);
1393
1394 impl Drop for TaskManagerShutdownGuard {
1395 fn drop(&mut self) {
1396 if let Some(manager) = self.0.upgrade() {
1397 manager.shutdown();
1398 if let Ok(runtime) = tokio::runtime::Handle::try_current() {
1399 runtime.spawn(async move {
1400 if let Err(error) = manager.shutdown_and_wait().await {
1401 tracing::error!(%error, "Task service shutdown failed");
1402 }
1403 });
1404 }
1405 }
1406 }
1407 }
1408
1409 /// A process generation stays live as long as either the manager or actual
1410 /// Runtime work retains it. The Runtime's existing scope lock excludes another
1411 /// executor for that scope; this generation lock proves a particular owner died.
1412 #[derive(Debug)]
1413 pub(crate) struct TaskExecutionLease {
1414 scope: String,
1415 generation: String,
1416 _scope_owner: Arc<RuntimeProcessOwnerLock>,
1417 _generation_owner: RuntimeProcessOwnerLock,
1418 }
1419
1420 impl TaskExecutionLease {
1421 fn new(
1422 root: &Path,
1423 scope: String,
1424 scope_owner: Arc<RuntimeProcessOwnerLock>,
1425 ) -> Result<Arc<Self>> {
1426 validate_execution_id(&scope, 64)?;
1427 let generation = Uuid::new_v4().simple().to_string();
1428 let path = execution_lease_path(root, &scope, &generation)?;
1429 let owner = RuntimeProcessOwnerLock::try_acquire_file(&path, true)?
1430 .context("Task execution generation is already owned")?;
1431 Ok(Arc::new(Self {
1432 scope,
1433 generation,
1434 _scope_owner: scope_owner,
1435 _generation_owner: owner,
1436 }))
1437 }
1438 }
1439
1440 fn validate_execution_id(value: &str, length: usize) -> Result<()> {
1441 if value.len() != length || !value.bytes().all(|byte| byte.is_ascii_hexdigit()) {
1442 bail!("Invalid task execution identity");
1443 }
1444 Ok(())
1445 }
1446
1447 fn execution_lease_path(root: &Path, scope: &str, generation: &str) -> Result<PathBuf> {
1448 validate_execution_id(scope, 64)?;
1449 validate_execution_id(generation, 32)?;
1450 Ok(root
1451 .join("execution-owners")
1452 .join(format!("{scope}.{generation}.lock")))
1453 }
1454
1455 #[cfg(test)]
1456 pub(crate) fn test_execution_scope(name: &str) -> String {
1457 use sha2::{Digest, Sha256};
1458 Sha256::digest(name.as_bytes())
1459 .iter()
1460 .map(|byte| format!("{byte:02x}"))
1461 .collect()
1462 }
1463
1464 pub struct TaskManager {
1465 cfg: TaskManagerConfig,
1466 default_workspace: Mutex<PathBuf>,
1467 executor: Arc<dyn TaskExecutor>,
1468 /// The runtime thread store this manager drives, when it owns one.
1469 runtime_threads: Option<SharedRuntimeThreadManager>,
1470 tasks_dir: PathBuf,
1471 artifacts_dir: PathBuf,
1472 queue_path: PathBuf,
1473 state: Mutex<ManagerState>,
1474 notify: Notify,
1475 cancel_token: CancellationToken,
1476 execution_lease: Arc<TaskExecutionLease>,
1477 workers: Mutex<Vec<tokio::task::JoinHandle<()>>>,
1478 shutdown_drain: Mutex<()>,
1479 /// Full store loads performed by this manager (tests only, #6573).
1480 #[cfg(test)]
1481 store_loads: std::sync::atomic::AtomicUsize,
1482 /// Queue fingerprint reads performed by worker loops (tests only, #6728).
1483 #[cfg(test)]
1484 fingerprint_reads: std::sync::atomic::AtomicUsize,
1485 }
1486
1487 /// Cheap stat-only view of the shared queue file. Everything that makes work
1488 /// claimable rewrites `queue.json` through an atomic rename: admission writes
1489 /// it before promoting the task record, and claims, cancels and startup
1490 /// recovery write it too. Running tasks flush their records into `tasks/`
1491 /// every few hundred milliseconds but leave the queue alone, so the tasks
1492 /// directory is deliberately not part of the fingerprint. Startup still rebuilds
1493 /// missing queue entries from durable records. Failed claims retry independently,
1494 /// including a queue removal whose task-status write failed. Missing or recent
1495 /// metadata keeps a fallback deadline until it can be trusted.
1496 #[derive(Debug, Clone, PartialEq, Eq)]
1497 struct StoreFingerprint(Option<(std::time::SystemTime, u64, u64)>);
1498
1499 impl StoreFingerprint {
1500 fn read(path: &Path) -> Self {
1501 Self(fs::metadata(path).ok().and_then(|meta| {
1502 #[cfg(unix)]
1503 let inode = std::os::unix::fs::MetadataExt::ino(&meta);
1504 #[cfg(not(unix))]
1505 let inode = 0;
1506 Some((meta.modified().ok()?, meta.len(), inode))
1507 }))
1508 }
1509
1510 /// True when the stamp is old enough that a later write cannot share it.
1511 fn settled(&self, now: std::time::SystemTime) -> bool {
1512 self.0.as_ref().is_none_or(|(modified, _, _)| {
1513 now.duration_since(*modified)
1514 .is_ok_and(|age| age >= READ_SNAPSHOT_SETTLE)
1515 })
1516 }
1517 }
1518
1519 /// Stat-only view of everything `load_state` reads, for read-only callers
1520 /// (#6573). Every task and queue write is an atomic rename, which bumps the
1521 /// tasks directory's or the queue file's mtime, so an unchanged snapshot
1522 /// means an unchanged store. The TUI task panel lists tasks every 2.5s; with
1523 /// this, an idle session answers from memory instead of taking the
1524 /// cross-process store lock and re-parsing every task record each time.
1525 #[derive(Debug, Clone, PartialEq, Eq)]
1526 struct ReadSnapshot {
1527 queue: StoreFingerprint,
1528 tasks_dir: StoreFingerprint,
1529 }
1530
1531 impl ReadSnapshot {
1532 fn read(tasks_dir: &Path, queue_path: &Path) -> Self {
1533 Self {
1534 queue: StoreFingerprint::read(queue_path),
1535 tasks_dir: StoreFingerprint::read(tasks_dir),
1536 }
1537 }
1538
1539 fn settled(&self, now: std::time::SystemTime) -> bool {
1540 self.queue.settled(now) && self.tasks_dir.settled(now)
1541 }
1542 }
1543
1544 /// When an idle worker should next claim from the shared store (#6573).
1545 ///
1546 /// A worker claims when it was notified in-process, when the queue
1547 /// fingerprint changed since its last attempt, or when a retry deadline passed.
1548 /// Settled empty claims have no deadline. Unreadable/recent metadata gets a
1549 /// fallback and failed claims get exponential backoff ("Task store is busy"). Workers
1550 /// keep checking the fingerprint every tick during a backoff, so a queue
1551 /// write from another process still gets an early retry.
1552 #[derive(Debug)]
1553 struct ClaimSchedule {
1554 seen: Option<StoreFingerprint>,
1555 next_claim: Option<Instant>,
1556 failure_backoff: Duration,
1557 }
1558
1559 impl ClaimSchedule {
1560 fn new(now: Instant) -> Self {
1561 Self {
1562 seen: None,
1563 next_claim: Some(now),
1564 failure_backoff: STORE_REFRESH_INTERVAL,
1565 }
1566 }
1567
1568 fn should_claim(&self, fingerprint: &StoreFingerprint, now: Instant, woken: bool) -> bool {
1569 woken
1570 || self.seen.as_ref() != Some(fingerprint)
1571 || self.next_claim.is_some_and(|deadline| now >= deadline)
1572 }
1573
1574 /// How long the worker may sleep before its next pass. Nothing claimable
1575 /// (a settled empty queue, no retry deadline) only needs to notice another
1576 /// process's queue write; pending work or a retry deadline keeps the
1577 /// short tick.
1578 fn nap(&self) -> Duration {
1579 if self.next_claim.is_none() {
1580 STORE_SETTLED_NAP
1581 } else {
1582 STORE_REFRESH_INTERVAL
1583 }
1584 }
1585
1586 fn claimed_task(&mut self, now: Instant) {
1587 self.seen = None;
1588 self.next_claim = Some(now);
1589 self.failure_backoff = STORE_REFRESH_INTERVAL;
1590 }
1591
1592 fn found_nothing(&mut self, fingerprint: StoreFingerprint, now: Instant) {
1593 // Re-read a recent stamp after its coarse mtime tick has passed before
1594 // trusting it indefinitely. Missing metadata never proves no change.
1595 let settled = fingerprint.0.as_ref().is_some_and(|(modified, _, _)| {
1596 std::time::SystemTime::now()
1597 .duration_since(*modified)
1598 .is_ok_and(|age| age >= STORE_IDLE_POLL_INTERVAL)
1599 });
1600 self.seen = Some(fingerprint);
1601 self.next_claim = (!settled).then_some(now + STORE_IDLE_POLL_INTERVAL);
1602 self.failure_backoff = STORE_REFRESH_INTERVAL;
1603 }
1604
1605 /// Returns the delay before the next claim unless the queue changes.
1606 fn claim_failed(&mut self, fingerprint: StoreFingerprint, now: Instant) -> Duration {
1607 let delay = self.failure_backoff;
1608 self.seen = Some(fingerprint);
1609 self.next_claim = Some(now + delay);
1610 self.failure_backoff = (delay * 2).min(STORE_BUSY_BACKOFF_MAX);
1611 delay
1612 }
1613 }
1614
1615 /// The "store is busy" error, naming the lock holder when its record can be
1616 /// read. Reading the record is best effort and never fails the caller (#6573).
1617 fn store_busy_error(lock_path: &Path) -> anyhow::Error {
1618 match RuntimeProcessOwnerLock::read_holder(lock_path) {
1619 Some((pid, held_for)) => anyhow!(
1620 "Task store is busy; state is unavailable (lock held by pid {pid} for {}s)",
1621 held_for.as_secs()
1622 ),
1623 None => anyhow!("Task store is busy; state is unavailable"),
1624 }
1625 }
1626
1627 struct ManagerState {
1628 tasks: HashMap<String, TaskRecord>,
1629 queue: VecDeque<String>,
1630 running_cancel: HashMap<String, CancellationToken>,
1631 /// Uncommitted typed deltas, reapplied to a fresh record before persistence.
1632 pending_events: HashMap<String, Vec<TaskExecutionEvent>>,
1633 /// Store snapshot that `tasks` and `queue` exactly reflect, set only by a
1634 /// read-only refresh and cleared by every path that may change them.
1635 read_snapshot: Option<ReadSnapshot>,
1636 }
1637
1638 #[derive(Debug, Serialize, Deserialize, Default)]
1639 struct QueueFile {
1640 queue: Vec<String>,
1641 }
1642
1643 impl TaskManager {
1644 /// Start the manager with the default DeepSeek executor.
1645 ///
1646 /// Interactive callers pass an initial session id to isolate new hosts,
1647 /// or the saved store binding to retain the same authority across resume.
1648 pub async fn start(
1649 cfg: TaskManagerConfig,
1650 api_config: Config,
1651 plugin_registry: Arc<crate::plugins::PluginRegistry>,
1652 session_id: &str,
1653 binding: Option<&crate::runtime_threads::RuntimeStoreBinding>,
1654 ) -> Result<SharedTaskManager> {
1655 // Resolve the sessions root's canonical spelling off this runtime
1656 // once, so the saved-store confinement checks here and in later
1657 // `/resume` / `/load` switches are pure comparisons (#6522).
1658 crate::runtime_threads::prepare_canonical_sessions_root().await;
1659 let runtime_threads = Arc::new(RuntimeThreadManager::open_for_session(
1660 api_config.clone(),
1661 cfg.default_workspace.clone(),
1662 RuntimeThreadManagerConfig::for_session(cfg.data_dir.clone(), session_id),
1663 plugin_registry,
1664 binding,
1665 )?);
1666 Self::start_with_runtime_manager(cfg, api_config, runtime_threads).await
1667 }
1668
1669 /// Start the manager with an injected runtime thread manager.
1670 pub async fn start_with_runtime_manager(
1671 cfg: TaskManagerConfig,
1672 _api_config: Config,
1673 runtime_threads: SharedRuntimeThreadManager,
1674 ) -> Result<SharedTaskManager> {
1675 let executor: Arc<dyn TaskExecutor> = Arc::new(EngineTaskExecutor::new(
1676 runtime_threads.clone(),
1677 cfg.execution_limits,
1678 ));
1679 let identity = runtime_threads.task_execution_identity();
1680 let manager = Self::start_with_executor_and_runtime(
1681 cfg,
1682 executor,
1683 Some(runtime_threads.clone()),
1684 identity,
1685 )
1686 .await?;
1687 runtime_threads.attach_task_manager(manager.clone());
1688 Ok(manager)
1689 }
1690
1691 /// Start the manager with a custom executor (used for tests).
1692 #[cfg(test)]
1693 pub async fn start_with_executor(
1694 cfg: TaskManagerConfig,
1695 executor: Arc<dyn TaskExecutor>,
1696 ) -> Result<SharedTaskManager> {
1697 Self::start_with_executor_in_scope(cfg, executor, "test").await
1698 }
1699
1700 #[cfg(test)]
1701 pub(crate) async fn start_with_executor_in_scope(
1702 cfg: TaskManagerConfig,
1703 executor: Arc<dyn TaskExecutor>,
1704 scope: &str,
1705 ) -> Result<SharedTaskManager> {
1706 let scope = test_execution_scope(scope);
1707 let path = cfg
1708 .data_dir
1709 .join("execution-owners")
1710 .join(format!("test-{scope}.lock"));
1711 let owner = RuntimeProcessOwnerLock::try_acquire_file(&path, true)?
1712 .context("Mock Runtime execution scope is already owned")?;
1713 Self::start_with_executor_and_runtime(cfg, executor, None, (scope, Arc::new(owner))).await
1714 }
1715
1716 async fn start_with_executor_and_runtime(
1717 cfg: TaskManagerConfig,
1718 executor: Arc<dyn TaskExecutor>,
1719 runtime_threads: Option<SharedRuntimeThreadManager>,
1720 identity: (String, Arc<RuntimeProcessOwnerLock>),
1721 ) -> Result<SharedTaskManager> {
1722 let workers = cfg.worker_count.clamp(1, MAX_WORKERS);
1723 let tasks_dir = cfg.data_dir.join("tasks");
1724 let artifacts_dir = cfg.data_dir.join("artifacts");
1725 let queue_path = cfg.data_dir.join("queue.json");
1726 tokio::fs::create_dir_all(&tasks_dir)
1727 .await
1728 .with_context(|| format!("Failed to create tasks dir {}", tasks_dir.display()))?;
1729 tokio::fs::create_dir_all(&artifacts_dir)
1730 .await
1731 .with_context(|| {
1732 format!(
1733 "Failed to create task artifacts dir {}",
1734 artifacts_dir.display()
1735 )
1736 })?;
1737
1738 let execution_lease = TaskExecutionLease::new(&cfg.data_dir, identity.0, identity.1)?;
1739 let cancel_token = CancellationToken::new();
1740 let default_workspace = cfg.default_workspace.clone();
1741 let manager = Arc::new(Self {
1742 cfg,
1743 default_workspace: Mutex::new(default_workspace),
1744 executor,
1745 runtime_threads,
1746 tasks_dir,
1747 artifacts_dir,
1748 queue_path,
1749 state: Mutex::new(ManagerState {
1750 tasks: HashMap::new(),
1751 queue: VecDeque::new(),
1752 running_cancel: HashMap::new(),
1753 pending_events: HashMap::new(),
1754 read_snapshot: None,
1755 }),
1756 notify: Notify::new(),
1757 cancel_token: cancel_token.clone(),
1758 execution_lease: execution_lease.clone(),
1759 workers: Mutex::new(Vec::new()),
1760 shutdown_drain: Mutex::new(()),
1761 #[cfg(test)]
1762 store_loads: std::sync::atomic::AtomicUsize::new(0),
1763 #[cfg(test)]
1764 fingerprint_reads: std::sync::atomic::AtomicUsize::new(0),
1765 });
1766
1767 {
1768 let mut state = manager.state.lock().await;
1769 let _transaction = manager.lock_store().await?;
1770 manager.refresh_locked(&mut state)?;
1771 manager.recover_dead_executions_locked(&mut state)?;
1772 manager.persist_queue_locked(&state.queue)?;
1773 }
1774 if let Some(runtime) = &manager.runtime_threads {
1775 runtime.retain_task_execution_lease(execution_lease)?;
1776 }
1777
1778 for _ in 0..workers {
1779 let manager_clone = Arc::clone(&manager);
1780 let worker = spawn_supervised(
1781 "task-manager-worker",
1782 std::panic::Location::caller(),
1783 async move {
1784 manager_clone.worker_loop().await;
1785 },
1786 );
1787 manager.workers.lock().await.push(worker);
1788 }
1789
1790 Ok(manager)
1791 }
1792
1793 /// Request shutdown. Ownership remains retained through actual execution.
1794 pub fn shutdown(&self) {
1795 self.cancel_token.cancel();
1796 }
1797
1798 pub(crate) fn shutdown_guard(self: &Arc<Self>) -> TaskManagerShutdownGuard {
1799 TaskManagerShutdownGuard(Arc::downgrade(self))
1800 }
1801
1802 pub(crate) async fn shutdown_and_wait(&self) -> Result<()> {
1803 self.shutdown();
1804 let _drain = self.shutdown_drain.lock().await;
1805 if let Some(runtime) = &self.runtime_threads {
1806 runtime.close_execution_admission().await;
1807 }
1808 // Waiting through the admission fence settles requests that passed
1809 // their admission check before shutdown was requested.
1810 {
1811 let _transaction = self.lock_store().await?;
1812 }
1813 let mut failure = None;
1814 {
1815 let mut workers = self.workers.lock().await;
1816 while let Some(worker) = workers.last_mut() {
1817 // Await by reference: canceling this caller leaves the join in
1818 // the manager, so a later drain still waits for actual exit.
1819 let result = worker.await;
1820 workers.pop();
1821 if let Err(error) = result {
1822 failure = Some(anyhow!("Task worker shutdown failed: {error}"));
1823 }
1824 }
1825 }
1826 if let Some(runtime) = &self.runtime_threads {
1827 runtime.shutdown_and_wait().await?;
1828 }
1829 failure.map_or(Ok(()), Err)
1830 }
1831
1832 pub(crate) fn execution_scope(&self) -> &str {
1833 &self.execution_lease.scope
1834 }
1835
1836 /// Apply `edit` to the runtime threads' authoritative config and reload
1837 /// it, so runtime-chat and queued runtime turns see a setting the UI just
1838 /// persisted instead of their startup snapshot. No runtime manager is a
1839 /// no-op.
1840 pub(crate) async fn reload_runtime_config_with(
1841 &self,
1842 edit: impl FnOnce(&mut crate::config::Config),
1843 ) -> Result<()> {
1844 let Some(runtime) = &self.runtime_threads else {
1845 return Ok(());
1846 };
1847 let mut config = runtime.read_config().clone();
1848 edit(&mut config);
1849 runtime.reload_config(config).await.map(|_| ())
1850 }
1851
1852 pub(crate) fn session_store_binding(
1853 &self,
1854 ) -> Option<crate::runtime_threads::RuntimeStoreBinding> {
1855 self.runtime_threads
1856 .as_ref()
1857 .map(|runtime| runtime.session_store_binding())
1858 }
1859
1860 pub async fn set_default_workspace(&self, workspace: PathBuf) {
1861 let mut default_workspace = self.default_workspace.lock().await;
1862 *default_workspace = workspace;
1863 }
1864
1865 pub async fn default_workspace(&self) -> PathBuf {
1866 self.default_workspace.lock().await.clone()
1867 }
1868
1869 /// Enqueue a new task.
1870 pub async fn add_task(&self, req: NewTaskRequest) -> Result<TaskRecord> {
1871 self.add_task_with_id(req, Self::new_task_id()).await
1872 }
1873
1874 /// Allocate the durable owner identity before queue insertion so callers
1875 /// can register graph spawn intent first.
1876 #[must_use]
1877 pub(crate) fn new_task_id() -> String {
1878 format!("task_{}", &Uuid::new_v4().simple().to_string()[..16])
1879 }
1880
1881 /// Read the exact durable task binding without adopting another process's
1882 /// queue. Used by the automation dispatcher while it owns the store claim.
1883 pub(crate) fn read_bound_task(&self, task_id: &str) -> Result<Option<TaskRecord>> {
1884 if task_id.is_empty()
1885 || !task_id
1886 .bytes()
1887 .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'_' | b'-'))
1888 {
1889 bail!("Invalid durable task id");
1890 }
1891 read_bound_task_file(&self.tasks_dir.join(format!("{task_id}.json")), task_id)
1892 }
1893
1894 /// Recover a persisted automation admission under its cross-process
1895 /// dispatch lock. A promoted task is accepted work, including failed or
1896 /// interrupted work; returning it must never enqueue it again.
1897 pub(crate) async fn recover_task_admission(
1898 &self,
1899 request: NewTaskRequest,
1900 task_id: String,
1901 ) -> Result<TaskRecord> {
1902 validate_preallocated_task_id(&task_id)?;
1903 if let Some(task) = self.read_bound_task(&task_id)? {
1904 validate_bound_task_request(&task, &request)?;
1905 return Ok(task);
1906 }
1907 let staged_path = self.tasks_dir.join(format!(".{task_id}.json.pending"));
1908 let admission = if let Some(staged) = read_bound_task_file(&staged_path, &task_id)? {
1909 validate_bound_task_request(&staged, &request)?;
1910 if staged.status != TaskStatus::Queued
1911 || staged.started_at.is_some()
1912 || staged.thread_id.is_some()
1913 || staged.turn_id.is_some()
1914 {
1915 bail!("Unpromoted task stage contains execution evidence; refusing to replay it");
1916 }
1917 // The original resolved settings remain durable in this stage
1918 // until promotion, including if recovery itself is interrupted.
1919 self.admit_task_record(staged, true).await
1920 } else {
1921 self.add_task_with_id(request.clone(), task_id.clone())
1922 .await
1923 };
1924 match admission {
1925 Ok(task) => Ok(task),
1926 Err(error) => {
1927 // Preserve a task promoted before a torn response; do not
1928 // replace its identity or convert accepted work into a retry.
1929 if let Some(task) = self.read_bound_task(&task_id)? {
1930 validate_bound_task_request(&task, &request)?;
1931 Ok(task)
1932 } else {
1933 Err(error)
1934 }
1935 }
1936 }
1937 }
1938
1939 /// Enqueue using a preallocated id. This is crate-visible only for the
1940 /// model tool's register-before-work transaction.
1941 pub(crate) async fn add_task_with_id(
1942 &self,
1943 req: NewTaskRequest,
1944 task_id: String,
1945 ) -> Result<TaskRecord> {
1946 let prompt = req.prompt.trim().to_string();
1947 if prompt.is_empty() {
1948 bail!("Task prompt cannot be empty");
1949 }
1950 if (req.model_provider.is_some() || req.model_provider_id.is_some())
1951 && req
1952 .model
1953 .as_deref()
1954 .is_none_or(|model| model.trim().is_empty())
1955 {
1956 bail!("A pinned task provider requires an explicit model");
1957 }
1958 // The worker runs this same projection when it opens the task's
1959 // thread. Running it here as well refuses an unknown mode or posture
1960 // at the boundary the request crossed, instead of after the task has
1961 // sat in the durable queue and a worker has claimed it.
1962 crate::runtime_policy::RuntimePolicyProjection::from_request(
1963 req.mode
1964 .as_deref()
1965 .filter(|mode| !mode.trim().is_empty())
1966 .unwrap_or(&self.cfg.default_mode),
1967 req.permission_posture.as_deref(),
1968 req.auto_approve,
1969 )?;
1970 validate_preallocated_task_id(&task_id)?;
1971
1972 let task = TaskRecord {
1973 schema_version: CURRENT_TASK_SCHEMA_VERSION,
1974 // 16 random hex chars (was 8; ~60 bits of entropy once UUIDv4's
1975 // fixed version nibble is discounted): task ids live in durable
1976 // state that accumulates across restarts, and a collision
1977 // overwrites a record while leaving a duplicate queue entry.
1978 // `resolve_task_id` matches by prefix, so short references still
1979 // work.
1980 id: task_id,
1981 prompt,
1982 name: req
1983 .name
1984 .map(|name| name.trim().to_string())
1985 .filter(|name| !name.is_empty()),
1986 model: req.model.unwrap_or_else(|| self.cfg.default_model.clone()),
1987 model_provider: req.model_provider,
1988 model_provider_id: req.model_provider_id,
1989 workspace: match req.workspace {
1990 Some(workspace) => workspace,
1991 None => self.default_workspace().await,
1992 },
1993 mode: req.mode.unwrap_or_else(|| self.cfg.default_mode.clone()),
1994 allow_shell: req.allow_shell.unwrap_or(self.cfg.allow_shell),
1995 trust_mode: req.trust_mode.unwrap_or(self.cfg.trust_mode),
1996 // Auto-approval must be opted into explicitly
1997 // (GHSA-72w5-pf8h-xfp4).
1998 auto_approve: req.auto_approve.unwrap_or(false),
1999 permission_posture: req.permission_posture,
2000 status: TaskStatus::Queued,
2001 created_at: Utc::now(),
2002 started_at: None,
2003 ended_at: None,
2004 duration_ms: None,
2005 result_summary: None,
2006 result_detail_path: None,
2007 error: None,
2008 terminal_reason: None,
2009 thread_id: None,
2010 turn_id: None,
2011 owner_session_id: req.owner_session_id,
2012 execution_scope: Some(self.execution_scope().to_string()),
2013 execution_generation: None,
2014 cancel_requested_seq: 0,
2015 runtime_event_count: 0,
2016 lifecycle_seq: 1,
2017 checklist: TaskChecklistState::default(),
2018 gates: Vec::new(),
2019 attempts: Vec::new(),
2020 artifacts: Vec::new(),
2021 github_events: Vec::new(),
2022 tool_calls: Vec::new(),
2023 timeline: vec![TaskTimelineEntry {
2024 timestamp: Utc::now(),
2025 kind: "queued".to_string(),
2026 summary: "Task queued".to_string(),
2027 detail_path: None,
2028 }],
2029 };
2030
2031 self.admit_task_record(task, false).await
2032 }
2033
2034 async fn admit_task_record(&self, task: TaskRecord, recover_stage: bool) -> Result<TaskRecord> {
2035 {
2036 let mut state = self.state.lock().await;
2037 let _transaction = self.lock_store().await?;
2038 self.refresh_locked(&mut state)?;
2039 if self.cancel_token.is_cancelled() {
2040 bail!("Task manager is shutting down; admission is closed");
2041 }
2042 if task.execution_scope.as_deref() != Some(self.execution_scope()) {
2043 bail!(
2044 "Task execution scope is unverified or belongs to another Runtime; refusing adoption"
2045 );
2046 }
2047 let task_path = self.tasks_dir.join(format!("{}.json", task.id));
2048 // The staged extension is intentionally not `.json`, so startup
2049 // replay ignores an interrupted create until the queue write has
2050 // succeeded and this file is atomically promoted.
2051 let staged_task_path = self.tasks_dir.join(format!(".{}.json.pending", task.id));
2052 if recover_stage {
2053 if let Some(accepted) = self.read_bound_task(&task.id)? {
2054 validate_bound_task_request(&accepted, &NewTaskRequest::from_task(&task))?;
2055 return Ok(accepted);
2056 }
2057 let current = read_bound_task_file(&staged_task_path, &task.id)?
2058 .context("Unaccepted task stage disappeared during recovery")?;
2059 if serde_json::to_value(&current)? != serde_json::to_value(&task)? {
2060 bail!("Unaccepted task stage changed during recovery");
2061 }
2062 }
2063 if state.tasks.contains_key(&task.id)
2064 || task_path.exists()
2065 || (!recover_stage && staged_task_path.exists())
2066 {
2067 bail!("Task id already exists: {}", task.id);
2068 }
2069 let mut next_queue = state.queue.clone();
2070 if !next_queue.contains(&task.id) {
2071 next_queue.push_back(task.id.clone());
2072 }
2073
2074 // Stage the owner record, then persist its queue membership, then
2075 // atomically promote it. A crash before promotion leaves either an
2076 // ignored staged file or a queue entry with no task (which replay
2077 // drops); a crash after promotion leaves the complete runnable
2078 // pair. In-memory scheduling is published only after all three.
2079 if !recover_stage {
2080 write_json_atomic(&staged_task_path, &task)?;
2081 }
2082 if let Err(err) = self.persist_queue_locked(&next_queue) {
2083 if !recover_stage
2084 && let Err(cleanup_err) = tokio::fs::remove_file(&staged_task_path).await
2085 {
2086 tracing::warn!(
2087 task_id = %task.id,
2088 error = %cleanup_err,
2089 "failed to remove ignored staged task after queue write failure"
2090 );
2091 }
2092 return Err(err);
2093 }
2094 if let Err(promote_err) = tokio::fs::rename(&staged_task_path, &task_path).await {
2095 let rollback_error = self.persist_queue_locked(&state.queue).err();
2096 let cleanup_error = if recover_stage {
2097 None
2098 } else {
2099 tokio::fs::remove_file(&staged_task_path).await.err()
2100 };
2101 let mut message =
2102 format!("Failed to promote staged task {}: {promote_err}", task.id);
2103 if let Some(rollback_error) = rollback_error {
2104 message.push_str(&format!("; queue rollback also failed: {rollback_error:#}"));
2105 }
2106 if let Some(cleanup_error) = cleanup_error {
2107 message.push_str(&format!(
2108 "; ignored staged-file cleanup also failed: {cleanup_error}"
2109 ));
2110 }
2111 bail!(message);
2112 }
2113 state.queue = next_queue;
2114 state.tasks.insert(task.id.clone(), task.clone());
2115 }
2116 self.notify.notify_one();
2117 Ok(task)
2118 }
2119
2120 /// List tasks, newest first.
2121 pub async fn list_tasks(&self, limit: Option<usize>) -> Result<Vec<TaskSummary>> {
2122 self.list_tasks_scoped(limit, None).await
2123 }
2124
2125 /// List tasks, newest first, optionally scoped to a workspace.
2126 pub async fn list_tasks_scoped(
2127 &self,
2128 limit: Option<usize>,
2129 workspace: Option<&Path>,
2130 ) -> Result<Vec<TaskSummary>> {
2131 self.list_tasks_visible_to(limit, workspace, None).await
2132 }
2133
2134 /// List tasks owned by a session, newest first, optionally scoped to a workspace.
2135 ///
2136 /// Ownerless legacy records fail closed and are not model-visible.
2137 pub async fn list_tasks_for_owner(
2138 &self,
2139 limit: Option<usize>,
2140 workspace: Option<&Path>,
2141 owner_session_id: &str,
2142 ) -> Result<Vec<TaskSummary>> {
2143 self.list_tasks_visible_to(limit, workspace, Some(owner_session_id))
2144 .await
2145 }
2146
2147 async fn list_tasks_visible_to(
2148 &self,
2149 limit: Option<usize>,
2150 workspace: Option<&Path>,
2151 owner_session_id: Option<&str>,
2152 ) -> Result<Vec<TaskSummary>> {
2153 let mut state = self.state.lock().await;
2154 self.refresh_for_read(&mut state).await?;
2155 let mut items = state
2156 .tasks
2157 .values()
2158 .filter(|record| {
2159 workspace.is_none_or(|workspace| record.workspace.as_path() == workspace)
2160 && owner_session_id.is_none_or(|owner_session_id| {
2161 record.owner_session_id.as_deref() == Some(owner_session_id)
2162 })
2163 })
2164 .map(TaskSummary::from)
2165 .collect::<Vec<_>>();
2166 items.sort_by_key(|i| std::cmp::Reverse(i.created_at));
2167 if let Some(limit) = limit {
2168 items.truncate(limit);
2169 }
2170 Ok(items)
2171 }
2172
2173 /// Retrieve a task by full id or prefix.
2174 pub async fn get_task(&self, id_or_prefix: &str) -> Result<TaskRecord> {
2175 self.get_task_visible_to(id_or_prefix, None).await
2176 }
2177
2178 /// Retrieve a session-owned task by full id or prefix.
2179 ///
2180 /// Ownership is applied before id resolution so foreign records cannot
2181 /// disclose their existence through exact matches or prefix ambiguity.
2182 pub async fn get_task_for_owner(
2183 &self,
2184 id_or_prefix: &str,
2185 owner_session_id: &str,
2186 ) -> Result<TaskRecord> {
2187 self.get_task_visible_to(id_or_prefix, Some(owner_session_id))
2188 .await
2189 }
2190
2191 /// Retrieve a task the interactive operator can inspect by full id or prefix.
2192 ///
2193 /// Scheduled automations have no session owner because they can outlive the
2194 /// session that configured them. They remain operator-visible only while
2195 /// bound to this manager's verified Runtime execution scope. Session-owned
2196 /// tasks keep their existing isolation, and unscoped legacy records remain
2197 /// hidden.
2198 pub(crate) async fn get_task_for_interactive_session(
2199 &self,
2200 id_or_prefix: &str,
2201 owner_session_id: &str,
2202 ) -> Result<TaskRecord> {
2203 let mut state = self.state.lock().await;
2204 self.refresh_for_read(&mut state).await?;
2205 let id = resolve_task_id_visible_to_operator(
2206 &state.tasks,
2207 id_or_prefix,
2208 owner_session_id,
2209 self.execution_scope(),
2210 )?;
2211 state
2212 .tasks
2213 .get(&id)
2214 .cloned()
2215 .ok_or_else(|| anyhow!("Task not found: {id_or_prefix}"))
2216 }
2217
2218 /// Retrieve the exact owned task stamped onto a trusted runtime thread.
2219 ///
2220 /// The runtime thread supplies a full durable id rather than model input.
2221 /// An ownerless task needs the current trusted execution scope. Legacy
2222 /// ownerless tasks without that provenance still fail closed.
2223 pub(crate) async fn get_task_for_active_runtime(&self, task_id: &str) -> Result<TaskRecord> {
2224 let mut state = self.state.lock().await;
2225 self.refresh_for_read(&mut state).await?;
2226 state
2227 .tasks
2228 .get(task_id)
2229 .filter(|task| match task.execution_scope.as_deref() {
2230 Some(scope) => scope == self.execution_scope(),
2231 None => task.owner_session_id.is_some(),
2232 })
2233 .cloned()
2234 .ok_or_else(|| anyhow!("Task not found: {task_id}"))
2235 }
2236
2237 async fn get_task_visible_to(
2238 &self,
2239 id_or_prefix: &str,
2240 owner_session_id: Option<&str>,
2241 ) -> Result<TaskRecord> {
2242 let mut state = self.state.lock().await;
2243 self.refresh_for_read(&mut state).await?;
2244 let id = resolve_task_id_visible_to(&state.tasks, id_or_prefix, owner_session_id)?;
2245 state
2246 .tasks
2247 .get(&id)
2248 .cloned()
2249 .ok_or_else(|| anyhow!("Task not found: {id_or_prefix}"))
2250 }
2251
2252 /// Cancel a queued or running task by id/prefix.
2253 pub async fn cancel_task(&self, id_or_prefix: &str) -> Result<TaskCancellation> {
2254 self.cancel_task_visible_to(id_or_prefix, None, None).await
2255 }
2256
2257 /// Cancel a queued or running task owned by the given session.
2258 ///
2259 /// Ownerless legacy and foreign records fail closed.
2260 pub async fn cancel_task_for_owner(
2261 &self,
2262 id_or_prefix: &str,
2263 owner_session_id: &str,
2264 ) -> Result<TaskCancellation> {
2265 self.cancel_task_visible_to(id_or_prefix, Some(owner_session_id), None)
2266 .await
2267 }
2268
2269 /// Cancel a task visible to the interactive operator.
2270 ///
2271 /// This is the cancellation counterpart to
2272 /// [`Self::get_task_for_interactive_session`]. It exists for human TUI
2273 /// actions, including the automation view's cancel button; model and child
2274 /// task APIs retain session-only visibility.
2275 pub(crate) async fn cancel_task_for_interactive_session(
2276 &self,
2277 id_or_prefix: &str,
2278 owner_session_id: &str,
2279 ) -> Result<TaskCancellation> {
2280 self.cancel_task_visible_to(
2281 id_or_prefix,
2282 Some(owner_session_id),
2283 Some(self.execution_scope()),
2284 )
2285 .await
2286 }
2287
2288 /// Cancel the exact owned task stamped onto a trusted runtime thread.
2289 pub(crate) async fn cancel_task_for_active_runtime(
2290 &self,
2291 task_id: &str,
2292 ) -> Result<TaskCancellation> {
2293 self.get_task_for_active_runtime(task_id).await?;
2294 self.cancel_task(task_id).await
2295 }
2296
2297 async fn cancel_task_visible_to(
2298 &self,
2299 id_or_prefix: &str,
2300 owner_session_id: Option<&str>,
2301 operator_execution_scope: Option<&str>,
2302 ) -> Result<TaskCancellation> {
2303 let mut state = self.state.lock().await;
2304 let _transaction = self.lock_store().await?;
2305 self.refresh_locked(&mut state)?;
2306 let id = match (owner_session_id, operator_execution_scope) {
2307 (Some(owner_session_id), Some(execution_scope)) => resolve_task_id_visible_to_operator(
2308 &state.tasks,
2309 id_or_prefix,
2310 owner_session_id,
2311 execution_scope,
2312 )?,
2313 _ => resolve_task_id_visible_to(&state.tasks, id_or_prefix, owner_session_id)?,
2314 };
2315 let now = Utc::now();
2316
2317 let mut cancel_running = false;
2318 let disposition = {
2319 let task = state
2320 .tasks
2321 .get_mut(&id)
2322 .ok_or_else(|| anyhow!("Task not found: {id}"))?;
2323 match task.status {
2324 TaskStatus::Queued => {
2325 task.status = TaskStatus::Canceled;
2326 task.lifecycle_seq = task.lifecycle_seq.saturating_add(1);
2327 task.ended_at = Some(now);
2328 task.duration_ms = Some(0);
2329 push_timeline_entry(
2330 task,
2331 TaskTimelineEntry {
2332 timestamp: now,
2333 kind: "canceled".to_string(),
2334 summary: "Task canceled before execution".to_string(),
2335 detail_path: None,
2336 },
2337 );
2338 state.queue.retain(|queued_id| queued_id != &id);
2339 TaskCancelDisposition::Forced
2340 }
2341 TaskStatus::Running => {
2342 cancel_running = true;
2343 task.lifecycle_seq = task.lifecycle_seq.saturating_add(1);
2344 task.cancel_requested_seq = task.lifecycle_seq;
2345 push_timeline_entry(
2346 task,
2347 TaskTimelineEntry {
2348 timestamp: now,
2349 kind: "cancel_requested".to_string(),
2350 summary: "Cancellation requested".to_string(),
2351 detail_path: None,
2352 },
2353 );
2354 TaskCancelDisposition::Requested
2355 }
2356 TaskStatus::Completed | TaskStatus::Failed | TaskStatus::Canceled => {
2357 TaskCancelDisposition::AlreadyFinished
2358 }
2359 }
2360 };
2361
2362 if cancel_running && let Some(token) = state.running_cancel.get(&id) {
2363 token.cancel();
2364 }
2365
2366 self.persist_changed_task_locked(&mut state, &id)?;
2367 self.persist_queue_locked(&state.queue)?;
2368 let task = state
2369 .tasks
2370 .get(&id)
2371 .cloned()
2372 .ok_or_else(|| anyhow!("Task not found: {id}"))?;
2373 Ok(TaskCancellation { task, disposition })
2374 }
2375
2376 /// Return aggregate status counters.
2377 pub async fn counts(&self) -> Result<TaskCounts> {
2378 let mut state = self.state.lock().await;
2379 self.refresh_for_read(&mut state).await?;
2380 let mut counts = TaskCounts::default();
2381 for task in state.tasks.values() {
2382 match task.status {
2383 TaskStatus::Queued => counts.queued += 1,
2384 TaskStatus::Running => counts.running += 1,
2385 TaskStatus::Completed => counts.completed += 1,
2386 TaskStatus::Failed => counts.failed += 1,
2387 TaskStatus::Canceled => counts.canceled += 1,
2388 }
2389 }
2390 Ok(counts)
2391 }
2392
2393 /// Root directory for durable task state.
2394 #[must_use]
2395 pub fn data_dir(&self) -> PathBuf {
2396 self.cfg.data_dir.clone()
2397 }
2398
2399 /// Live events from the runtime thread store this manager drives, when it
2400 /// owns one. The TUI taps `runtime.store_failure` here (#5931).
2401 #[must_use]
2402 pub fn subscribe_runtime_events(
2403 &self,
2404 ) -> Option<tokio::sync::broadcast::Receiver<RuntimeEventRecord>> {
2405 self.runtime_threads
2406 .as_ref()
2407 .map(|runtime| runtime.subscribe_events())
2408 }
2409
2410 /// Resolve a task artifact reference to an absolute path.
2411 #[must_use]
2412 pub fn artifact_absolute_path(&self, path: &Path) -> PathBuf {
2413 if path.is_absolute() {
2414 path.to_path_buf()
2415 } else {
2416 self.cfg.data_dir.join(path)
2417 }
2418 }
2419
2420 /// Write a durable task artifact and return the persisted path reference.
2421 pub fn write_task_artifact(
2422 &self,
2423 task_id: &str,
2424 label: &str,
2425 content: &str,
2426 ) -> Result<PathBuf> {
2427 self.write_artifact(task_id, label, content)
2428 }
2429
2430 /// Apply model-visible tool metadata to a task and persist it.
2431 pub async fn record_tool_metadata(
2432 &self,
2433 id_or_prefix: &str,
2434 metadata: &Value,
2435 ) -> Result<TaskRecord> {
2436 let mut state = self.state.lock().await;
2437 let _transaction = self.lock_store().await?;
2438 self.refresh_locked(&mut state)?;
2439 let id = resolve_task_id(&state.tasks, id_or_prefix)?;
2440 let updated = {
2441 let task = state
2442 .tasks
2443 .get_mut(&id)
2444 .ok_or_else(|| anyhow!("Task not found: {id}"))?;
2445 self.apply_task_update_metadata(task, Some(metadata))?;
2446 task.clone()
2447 };
2448 self.persist_changed_task_locked(&mut state, &id)?;
2449 Ok(updated)
2450 }
2451
2452 async fn claim_next_task(&self) -> Result<Option<(String, ExecutionTask, CancellationToken)>> {
2453 let mut state = self.state.lock().await;
2454 let _transaction = self.lock_store().await?;
2455 // The worker loop runs this repeatedly; keep the directory scan and
2456 // JSON parsing off the async runtime threads (#6573).
2457 let (tasks_dir, queue_path) = (self.tasks_dir.clone(), self.queue_path.clone());
2458 let loaded = tokio::task::spawn_blocking(move || load_state(&tasks_dir, &queue_path))
2459 .await
2460 .context("Task store load was interrupted")??;
2461 self.apply_loaded_locked(&mut state, loaded)?;
2462 if self.cancel_token.is_cancelled() {
2463 return Ok(None);
2464 }
2465 let Some(id) = state
2466 .queue
2467 .iter()
2468 .find(|id| {
2469 state.tasks.get(*id).is_some_and(|task| {
2470 task.status == TaskStatus::Queued
2471 && task.execution_scope.as_deref() == Some(self.execution_scope())
2472 })
2473 })
2474 .cloned()
2475 else {
2476 return Ok(None);
2477 };
2478 state.read_snapshot = None;
2479 state.queue.retain(|queued| queued != &id);
2480 let task = state
2481 .tasks
2482 .get_mut(&id)
2483 .context("Claimed task is missing")?;
2484 let now = Utc::now();
2485 task.status = TaskStatus::Running;
2486 task.execution_generation = Some(self.execution_lease.generation.clone());
2487 task.lifecycle_seq = task.lifecycle_seq.saturating_add(1);
2488 task.started_at = Some(now);
2489 task.ended_at = None;
2490 task.duration_ms = None;
2491 task.error = None;
2492 push_timeline_entry(
2493 task,
2494 TaskTimelineEntry {
2495 timestamp: now,
2496 kind: "running".into(),
2497 summary: "Task started".into(),
2498 detail_path: None,
2499 },
2500 );
2501 let request = ExecutionTask::from(&*task);
2502 // Removing queue membership first is recoverable from the still-Queued
2503 // record. Executor polling requires BOTH durable writes to succeed.
2504 self.persist_queue_locked(&state.queue)?;
2505 self.persist_changed_task_locked(&mut state, &id)?;
2506 let cancel = CancellationToken::new();
2507 state.running_cancel.insert(id.clone(), cancel.clone());
2508 Ok(Some((id, request, cancel)))
2509 }
2510
2511 /// Claim and run queued work.
2512 ///
2513 /// Several processes can share one data dir, and every claim takes the
2514 /// cross-process store lock and reloads the whole store. An idle worker
2515 /// therefore only claims when it was notified in-process, when the store
2516 /// fingerprint changed, or when unreadable/recent metadata needs a retry.
2517 /// A settled empty queue does no periodic full reload; failed claims still
2518 /// back off exponentially instead of retrying every tick (#6573, #6728).
2519 async fn worker_loop(self: Arc<Self>) {
2520 let mut schedule = ClaimSchedule::new(Instant::now());
2521 let mut woken = true;
2522 // True while consecutive claims are failing: the first failure of an
2523 // episode is an error, repeats are debug until a claim succeeds (#6573).
2524 let mut claim_failing = false;
2525 loop {
2526 if self.cancel_token.is_cancelled() {
2527 break;
2528 }
2529 // Read before claiming so a write racing the claim changes the
2530 // fingerprint and is picked up on the next pass.
2531 #[cfg(test)]
2532 self.fingerprint_reads
2533 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
2534 let fingerprint = StoreFingerprint::read(&self.queue_path);
2535 if schedule.should_claim(&fingerprint, Instant::now(), woken) {
2536 match self.claim_next_task().await {
2537 Ok(Some((id, request, cancel))) => {
2538 claim_failing = false;
2539 schedule.claimed_task(Instant::now());
2540 self.run_task(id, request, cancel).await;
2541 woken = true;
2542 continue;
2543 }
2544 Ok(None) => {
2545 claim_failing = false;
2546 schedule.found_nothing(fingerprint, Instant::now());
2547 }
2548 Err(error) => {
2549 let retry = schedule.claim_failed(fingerprint, Instant::now());
2550 let retry_ms = retry.as_millis() as u64;
2551 if claim_failing {
2552 tracing::debug!(
2553 %error,
2554 retry_ms,
2555 "Task claim still unavailable; executor was not polled"
2556 );
2557 } else {
2558 claim_failing = true;
2559 tracing::error!(
2560 %error,
2561 retry_ms,
2562 "Task claim unavailable; executor was not polled"
2563 );
2564 }
2565 }
2566 }
2567 }
2568 woken = false;
2569 tokio::select! {
2570 _ = self.cancel_token.cancelled() => break,
2571 _ = self.notify.notified() => woken = true,
2572 _ = sleep(schedule.nap()) => {},
2573 }
2574 }
2575 }
2576
2577 fn observe_task_cancellation(&self, task_id: &str, cancel: &CancellationToken) -> Result<()> {
2578 let task = self
2579 .read_bound_task(task_id)?
2580 .context("Running task is missing")?;
2581 self.require_execution_owner(&task)?;
2582 if task.cancel_requested_seq > 0 {
2583 cancel.cancel();
2584 }
2585 Ok(())
2586 }
2587
2588 async fn run_task(&self, task_id: String, request: ExecutionTask, cancel: CancellationToken) {
2589 let (event_tx, mut event_rx) = mpsc::channel(TASK_EVENT_CHANNEL_CAPACITY);
2590 let exec_fut = self
2591 .executor
2592 .execute(request.clone(), event_tx, cancel.clone());
2593 tokio::pin!(exec_fut);
2594
2595 let mut guard = ExecutionGuard::new(self.cfg.execution_limits, Instant::now());
2596 let mut dirty = false;
2597 let mut accumulated_result_text = String::new();
2598 let persist_debounce = self.cfg.execution_limits.persist_debounce;
2599
2600 let mut next_store_poll = Instant::now();
2601 let mut execution_started = false;
2602 let mut blocked_event: Option<TaskExecutionEvent> = None;
2603 let mut next_event_retry = Instant::now();
2604 let (mut result, manager_terminalized) = loop {
2605 if Instant::now() >= next_store_poll {
2606 if let Err(error) = self.observe_task_cancellation(&task_id, &cancel) {
2607 tracing::error!(%error, "Task ownership unavailable; requesting Runtime cancellation");
2608 cancel.cancel();
2609 }
2610 next_store_poll = Instant::now() + STORE_REFRESH_INTERVAL;
2611 }
2612 if !execution_started && (cancel.is_cancelled() || self.cancel_token.is_cancelled()) {
2613 let reason = if self.cancel_token.is_cancelled() {
2614 TaskTerminalReason::Shutdown
2615 } else {
2616 TaskTerminalReason::Canceled
2617 };
2618 break (TaskExecutionResult::from_reason(reason, None), true);
2619 }
2620 if Instant::now() >= next_event_retry {
2621 if let Some(event) = blocked_event.take()
2622 && self
2623 .process_execution_event(
2624 &task_id,
2625 event.clone(),
2626 &mut guard,
2627 &mut accumulated_result_text,
2628 &mut dirty,
2629 )
2630 .await
2631 .is_err()
2632 {
2633 blocked_event = Some(event);
2634 cancel.cancel();
2635 }
2636 next_event_retry = Instant::now() + STORE_REFRESH_INTERVAL;
2637 }
2638 let mut action = guard.evaluate(
2639 Instant::now(),
2640 cancel.is_cancelled(),
2641 self.cancel_token.is_cancelled(),
2642 );
2643 if matches!(
2644 action,
2645 GuardAction::Interrupt {
2646 reason: TaskTerminalReason::IdleTimeout
2647 }
2648 ) {
2649 // Progress already accepted by the executor must win an idle
2650 // deadline race. Drain only the events queued at this instant
2651 // so a producer cannot keep the watchdog from re-evaluating
2652 // wall time, shutdown, or explicit cancellation indefinitely.
2653 let queued = event_rx.len();
2654 for _ in 0..queued {
2655 if blocked_event.is_some() {
2656 break;
2657 }
2658 let Ok(event) = event_rx.try_recv() else {
2659 break;
2660 };
2661 if self
2662 .process_execution_event(
2663 &task_id,
2664 event.clone(),
2665 &mut guard,
2666 &mut accumulated_result_text,
2667 &mut dirty,
2668 )
2669 .await
2670 .is_err()
2671 {
2672 blocked_event = Some(event);
2673 cancel.cancel();
2674 }
2675 }
2676 action = guard.evaluate(
2677 Instant::now(),
2678 cancel.is_cancelled(),
2679 self.cancel_token.is_cancelled(),
2680 );
2681 }
2682 match action {
2683 GuardAction::Interrupt { reason } => {
2684 cancel.cancel();
2685 guard.note_interrupt(Instant::now(), reason);
2686 continue;
2687 }
2688 GuardAction::Terminalize { reason } => {
2689 break (TaskExecutionResult::from_reason(reason, None), true);
2690 }
2691 GuardAction::Run { wait } => {
2692 execution_started = true;
2693 tokio::select! {
2694 biased;
2695 exec_result = &mut exec_fut => {
2696 break (guard.preserve_timeout_reason(exec_result), false);
2697 }
2698 maybe_event = event_rx.recv(), if blocked_event.is_none() => {
2699 if let Some(event) = maybe_event
2700 && self.process_execution_event(
2701 &task_id, event.clone(), &mut guard,
2702 &mut accumulated_result_text, &mut dirty,
2703 ).await.is_err() {
2704 blocked_event = Some(event);
2705 cancel.cancel();
2706 }
2707 }
2708 _ = self.cancel_token.cancelled(), if !self.cancel_token.is_cancelled() => {
2709 cancel.cancel();
2710 }
2711 _ = sleep(persist_debounce), if dirty => {
2712 match self.flush_task(&task_id).await {
2713 Ok(()) => dirty = false,
2714 Err(err) => {
2715 tracing::error!("Failed to debounce-persist task {task_id}: {err}");
2716 cancel.cancel();
2717 }
2718 }
2719 }
2720 _ = sleep(wait.min(STORE_REFRESH_INTERVAL)) => {}
2721 }
2722 }
2723 }
2724 };
2725
2726 // Stop accepting producer events while one event cannot be retained.
2727 // The pending delta set and channel remain bounded during disk failure;
2728 // keep the execution lease until accepted events and the terminal receipt
2729 // have actually been persisted. A storage failure is not completion.
2730 event_rx.close();
2731 loop {
2732 let event = blocked_event.take().or_else(|| event_rx.try_recv().ok());
2733 let Some(event) = event else {
2734 break;
2735 };
2736 while self
2737 .process_execution_event(
2738 &task_id,
2739 event.clone(),
2740 &mut guard,
2741 &mut accumulated_result_text,
2742 &mut dirty,
2743 )
2744 .await
2745 .is_err()
2746 {
2747 sleep(STORE_REFRESH_INTERVAL).await;
2748 }
2749 }
2750 if manager_terminalized {
2751 result.result_text = optional_nonzero_text(accumulated_result_text);
2752 }
2753 loop {
2754 match self
2755 .finish_task(
2756 &task_id,
2757 result.clone(),
2758 cancel.clone(),
2759 &request.mode_label,
2760 )
2761 .await
2762 {
2763 Ok(()) => break,
2764 Err(err) => {
2765 tracing::error!(
2766 "Task {task_id} terminal receipt is pending storage recovery: {err}"
2767 );
2768 sleep(STORE_REFRESH_INTERVAL).await;
2769 }
2770 }
2771 }
2772 }
2773
2774 async fn process_execution_event(
2775 &self,
2776 task_id: &str,
2777 event: TaskExecutionEvent,
2778 guard: &mut ExecutionGuard,
2779 accumulated_result_text: &mut String,
2780 dirty: &mut bool,
2781 ) -> Result<()> {
2782 match self.apply_execution_event(task_id, event.clone()).await {
2783 Ok(outcome) => {
2784 if execution_event_is_progress(&event) {
2785 guard.note_progress(Instant::now());
2786 }
2787 append_message_delta(accumulated_result_text, &event);
2788 *dirty = !outcome.persisted;
2789 Ok(())
2790 }
2791 Err(err) => {
2792 tracing::error!("Task {task_id} event is waiting for storage recovery: {err}");
2793 Err(err)
2794 }
2795 }
2796 }
2797
2798 async fn apply_execution_event(
2799 &self,
2800 task_id: &str,
2801 event: TaskExecutionEvent,
2802 ) -> Result<EventApplyOutcome> {
2803 let urgent = execution_event_persist_urgent(&event);
2804 let mut state = self.state.lock().await;
2805 let _transaction = self.lock_store().await?;
2806 self.refresh_locked(&mut state)?;
2807 if state
2808 .pending_events
2809 .get(task_id)
2810 .is_some_and(|events| events.len() >= TASK_EVENT_CHANNEL_CAPACITY)
2811 {
2812 self.persist_changed_task_locked(&mut state, task_id)?;
2813 }
2814 let task = state
2815 .tasks
2816 .get_mut(task_id)
2817 .context("Event task is missing")?;
2818 self.require_execution_owner(task)?;
2819 self.apply_event_to_task(task, event.clone())?;
2820 let pending = state.pending_events.entry(task_id.to_string()).or_default();
2821 pending.push(event);
2822 let persist_now = urgent || pending.len() >= TASK_EVENT_CHANNEL_CAPACITY;
2823 let persisted = if persist_now {
2824 match self.persist_changed_task_locked(&mut state, task_id) {
2825 Ok(()) => true,
2826 Err(error) => {
2827 // This event is already retained in the bounded delta set.
2828 // Return acceptance so the caller must not append it again.
2829 if let Some(cancel) = state.running_cancel.get(task_id) {
2830 cancel.cancel();
2831 }
2832 tracing::error!(%error, "Task event retained pending storage recovery; cancellation requested");
2833 false
2834 }
2835 }
2836 } else {
2837 false
2838 };
2839 Ok(EventApplyOutcome { persisted })
2840 }
2841
2842 fn apply_event_to_task(&self, task: &mut TaskRecord, event: TaskExecutionEvent) -> Result<()> {
2843 let task_id = task.id.clone();
2844 match event {
2845 TaskExecutionEvent::ThreadLinked { thread_id, turn_id } => {
2846 task.thread_id = Some(thread_id.clone());
2847 task.turn_id = Some(turn_id.clone());
2848 push_timeline_entry(
2849 task,
2850 TaskTimelineEntry {
2851 timestamp: Utc::now(),
2852 kind: "runtime_link".to_string(),
2853 summary: format!("Linked runtime thread {thread_id} turn {turn_id}"),
2854 detail_path: None,
2855 },
2856 );
2857 }
2858 TaskExecutionEvent::Status { message } => {
2859 push_timeline_entry(
2860 task,
2861 TaskTimelineEntry {
2862 timestamp: Utc::now(),
2863 kind: "status".to_string(),
2864 summary: summarize_text(&message, TIMELINE_SUMMARY_LIMIT),
2865 detail_path: None,
2866 },
2867 );
2868 }
2869 TaskExecutionEvent::MessageDelta { content } => {
2870 if !content.trim().is_empty() {
2871 push_timeline_entry(
2872 task,
2873 TaskTimelineEntry {
2874 timestamp: Utc::now(),
2875 kind: "message".to_string(),
2876 summary: summarize_text(&content, TIMELINE_SUMMARY_LIMIT),
2877 detail_path: None,
2878 },
2879 );
2880 }
2881 }
2882 TaskExecutionEvent::ToolStarted { id, name, input } => {
2883 let input_summary = summarize_json(&input);
2884 task.tool_calls.push(TaskToolCallSummary {
2885 id: id.clone(),
2886 name: name.clone(),
2887 status: TaskToolStatus::Running,
2888 started_at: Utc::now(),
2889 ended_at: None,
2890 duration_ms: None,
2891 input_summary: input_summary.clone(),
2892 output_summary: None,
2893 detail_path: None,
2894 patch_ref: None,
2895 });
2896 let summary = input_summary
2897 .map(|s| format!("{name} started ({s})"))
2898 .unwrap_or_else(|| format!("{name} started"));
2899 push_timeline_entry(
2900 task,
2901 TaskTimelineEntry {
2902 timestamp: Utc::now(),
2903 kind: "tool_started".to_string(),
2904 summary,
2905 detail_path: None,
2906 },
2907 );
2908 }
2909 TaskExecutionEvent::ToolProgress { id, output } => {
2910 push_timeline_entry(
2911 task,
2912 TaskTimelineEntry {
2913 timestamp: Utc::now(),
2914 kind: "tool_progress".to_string(),
2915 summary: format!(
2916 "{id}: {}",
2917 summarize_text(&output, TIMELINE_SUMMARY_LIMIT.saturating_sub(8))
2918 ),
2919 detail_path: None,
2920 },
2921 );
2922 }
2923 // Supervisor-side liveness only: recording it would put a
2924 // timeline entry behind every poll tick of a silent build.
2925 TaskExecutionEvent::ToolHeartbeat => {}
2926 TaskExecutionEvent::ToolCompleted {
2927 id,
2928 name,
2929 success,
2930 output,
2931 metadata,
2932 } => {
2933 let now = Utc::now();
2934 let detail_path = self.artifact_if_large(&task_id, &name, &output)?;
2935 let output_summary = summarize_text(&output, TIMELINE_SUMMARY_LIMIT);
2936 let patch_ref = if name == "apply_patch" {
2937 detail_path.clone()
2938 } else {
2939 None
2940 };
2941
2942 if let Some(call) = task.tool_calls.iter_mut().find(|call| call.id == id) {
2943 call.status = if success {
2944 TaskToolStatus::Success
2945 } else {
2946 TaskToolStatus::Failed
2947 };
2948 call.ended_at = Some(now);
2949 call.duration_ms = Some(duration_ms(call.started_at, now));
2950 call.output_summary = Some(output_summary.clone());
2951 call.detail_path = detail_path.clone();
2952 call.patch_ref = patch_ref.clone();
2953
2954 if call.duration_ms.is_none()
2955 && let Some(duration) = metadata
2956 .as_ref()
2957 .and_then(|m| m.get("duration_ms"))
2958 .and_then(Value::as_u64)
2959 {
2960 call.duration_ms = Some(duration);
2961 }
2962 }
2963
2964 let status = if success { "success" } else { "failed" };
2965 push_timeline_entry(
2966 task,
2967 TaskTimelineEntry {
2968 timestamp: now,
2969 kind: "tool_completed".to_string(),
2970 summary: format!("{name} {status}: {output_summary}"),
2971 detail_path: detail_path.clone(),
2972 },
2973 );
2974 if let Some(patch_ref) = patch_ref {
2975 push_timeline_entry(
2976 task,
2977 TaskTimelineEntry {
2978 timestamp: now,
2979 kind: "patch_ref".to_string(),
2980 summary: format!("Patch artifact: {}", patch_ref.display()),
2981 detail_path: Some(patch_ref),
2982 },
2983 );
2984 }
2985
2986 self.apply_task_update_metadata(task, metadata.as_ref())?;
2987 }
2988 TaskExecutionEvent::Error { message } => {
2989 push_timeline_entry(
2990 task,
2991 TaskTimelineEntry {
2992 timestamp: Utc::now(),
2993 kind: "error".to_string(),
2994 summary: summarize_text(&message, TIMELINE_SUMMARY_LIMIT),
2995 detail_path: None,
2996 },
2997 );
2998 }
2999 TaskExecutionEvent::RuntimeEvent {
3000 seq,
3001 event,
3002 summary,
3003 } => {
3004 task.runtime_event_count = task.runtime_event_count.saturating_add(1);
3005 push_timeline_entry(
3006 task,
3007 TaskTimelineEntry {
3008 timestamp: Utc::now(),
3009 kind: "runtime_event".to_string(),
3010 summary: format!("#{seq} {event}: {summary}"),
3011 detail_path: None,
3012 },
3013 );
3014 }
3015 }
3016
3017 Ok(())
3018 }
3019
3020 async fn flush_task(&self, task_id: &str) -> Result<()> {
3021 let mut state = self.state.lock().await;
3022 let _transaction = self.lock_store().await?;
3023 self.refresh_locked(&mut state)?;
3024 let task = state
3025 .tasks
3026 .get(task_id)
3027 .context("Flushed task is missing")?;
3028 self.require_execution_owner(task)?;
3029 self.persist_changed_task_locked(&mut state, task_id)
3030 }
3031
3032 async fn finish_task(
3033 &self,
3034 task_id: &str,
3035 mut result: TaskExecutionResult,
3036 cancel: CancellationToken,
3037 mode_label: &str,
3038 ) -> Result<()> {
3039 let mut state = self.state.lock().await;
3040 let _transaction = self.lock_store().await?;
3041 self.refresh_locked(&mut state)?;
3042 state.running_cancel.remove(task_id);
3043 let task = state
3044 .tasks
3045 .get_mut(task_id)
3046 .context("Finished task is missing")?;
3047 self.require_execution_owner(task)?;
3048
3049 let now = Utc::now();
3050 if (cancel.is_cancelled() || task.cancel_requested_seq > 0)
3051 && result.status == TaskStatus::Completed
3052 {
3053 result.status = TaskStatus::Canceled;
3054 result.result_text = None;
3055 result.error = None;
3056 result.terminal_reason = TaskTerminalReason::Canceled;
3057 }
3058 if self.cancel_token.is_cancelled()
3059 && result.status != TaskStatus::Completed
3060 && matches!(
3061 result.terminal_reason,
3062 TaskTerminalReason::Canceled | TaskTerminalReason::Failed
3063 )
3064 {
3065 result.status = TaskStatus::Canceled;
3066 result.terminal_reason = TaskTerminalReason::Shutdown;
3067 result.error = Some(TaskTerminalReason::Shutdown.receipt_message());
3068 }
3069
3070 task.status = result.status;
3071 task.lifecycle_seq = task.lifecycle_seq.saturating_add(1);
3072 task.mode = mode_label.to_string();
3073 task.ended_at = Some(now);
3074 task.duration_ms = task.started_at.map(|start| duration_ms(start, now));
3075 task.error = result.error.clone();
3076 task.terminal_reason = Some(result.terminal_reason.as_str().to_string());
3077 let finished_summary = if matches!(result.status, TaskStatus::Queued | TaskStatus::Running)
3078 {
3079 format!("Task ended in unexpected state: {mode_label}")
3080 } else {
3081 match result.terminal_reason {
3082 TaskTerminalReason::Completed
3083 | TaskTerminalReason::Canceled
3084 | TaskTerminalReason::Shutdown => result.terminal_reason.receipt_message(),
3085 TaskTerminalReason::Failed
3086 | TaskTerminalReason::WallTimeout
3087 | TaskTerminalReason::IdleTimeout
3088 | TaskTerminalReason::CancelTimeout => format!(
3089 "{}: {}",
3090 result.terminal_reason.as_str(),
3091 result
3092 .error
3093 .as_deref()
3094 .map(|e| summarize_text(e, TIMELINE_SUMMARY_LIMIT))
3095 .unwrap_or_else(|| result.terminal_reason.receipt_message())
3096 ),
3097 }
3098 };
3099 push_timeline_entry(
3100 task,
3101 TaskTimelineEntry {
3102 timestamp: now,
3103 kind: "finished".to_string(),
3104 summary: finished_summary,
3105 detail_path: None,
3106 },
3107 );
3108
3109 if let Some(text) = result.result_text {
3110 let detail_path = self.artifact_if_large(task_id, "result", &text)?;
3111 task.result_summary = Some(summarize_text(&text, TIMELINE_SUMMARY_LIMIT));
3112 task.result_detail_path = detail_path.clone();
3113 if let Some(detail_path) = detail_path {
3114 push_timeline_entry(
3115 task,
3116 TaskTimelineEntry {
3117 timestamp: now,
3118 kind: "result_ref".to_string(),
3119 summary: format!("Result artifact: {}", detail_path.display()),
3120 detail_path: Some(detail_path),
3121 },
3122 );
3123 }
3124 } else if result.status == TaskStatus::Completed {
3125 task.result_summary = Some("(no textual output)".to_string());
3126 }
3127
3128 self.persist_changed_task_locked(&mut state, task_id)?;
3129 Ok(())
3130 }
3131
3132 fn artifact_if_large(
3133 &self,
3134 task_id: &str,
3135 label: &str,
3136 content: &str,
3137 ) -> Result<Option<PathBuf>> {
3138 if content.len() < ARTIFACT_THRESHOLD {
3139 return Ok(None);
3140 }
3141 self.write_artifact(task_id, label, content).map(Some)
3142 }
3143
3144 fn write_artifact(&self, task_id: &str, label: &str, content: &str) -> Result<PathBuf> {
3145 ensure_safe_storage_id("task id", task_id)?;
3146 let artifact_dir = self.artifacts_dir.join(task_id);
3147 fs::create_dir_all(&artifact_dir)
3148 .with_context(|| format!("Failed to create artifact dir {}", artifact_dir.display()))?;
3149 let stamp = Utc::now().format("%Y%m%dT%H%M%S%.3fZ");
3150 let filename = format!("{stamp}_{}.txt", sanitize_filename(label));
3151 let absolute = artifact_dir.join(filename);
3152 fs::write(&absolute, content)
3153 .with_context(|| format!("Failed to write artifact {}", absolute.display()))?;
3154 let relative = absolute
3155 .strip_prefix(&self.cfg.data_dir)
3156 .map(PathBuf::from)
3157 .unwrap_or(absolute);
3158 Ok(relative)
3159 }
3160
3161 fn apply_task_update_metadata(
3162 &self,
3163 task: &mut TaskRecord,
3164 metadata: Option<&Value>,
3165 ) -> Result<()> {
3166 let Some(updates) = metadata.and_then(|m| m.get("task_updates")) else {
3167 return Ok(());
3168 };
3169 let now = Utc::now();
3170
3171 if let Some(value) = updates.get("checklist") {
3172 let mut checklist: TaskChecklistState = serde_json::from_value(value.clone())
3173 .context("Failed to parse checklist task update")?;
3174 checklist.updated_at = checklist.updated_at.or(Some(now));
3175 task.checklist = checklist;
3176 push_timeline_entry(
3177 task,
3178 TaskTimelineEntry {
3179 timestamp: now,
3180 kind: "checklist".to_string(),
3181 summary: format!(
3182 "Checklist updated: {} item(s), {}% complete",
3183 task.checklist.items.len(),
3184 task.checklist.completion_pct
3185 ),
3186 detail_path: None,
3187 },
3188 );
3189 }
3190
3191 if let Some(value) = updates.get("gate") {
3192 let gate: TaskGateRecord = serde_json::from_value(value.clone())
3193 .context("Failed to parse gate task update")?;
3194 let summary = format!("Gate {} {}: {}", gate.gate, gate.status, gate.summary);
3195 task.gates.retain(|existing| existing.id != gate.id);
3196 task.gates.push(gate.clone());
3197 push_timeline_entry(
3198 task,
3199 TaskTimelineEntry {
3200 timestamp: now,
3201 kind: "gate".to_string(),
3202 summary: summarize_text(&summary, TIMELINE_SUMMARY_LIMIT),
3203 detail_path: gate.log_path,
3204 },
3205 );
3206 }
3207
3208 if let Some(value) = updates.get("attempt") {
3209 let attempt: TaskAttemptRecord = serde_json::from_value(value.clone())
3210 .context("Failed to parse attempt task update")?;
3211 task.attempts.retain(|existing| existing.id != attempt.id);
3212 task.attempts.push(attempt.clone());
3213 push_timeline_entry(
3214 task,
3215 TaskTimelineEntry {
3216 timestamp: now,
3217 kind: "pr_attempt".to_string(),
3218 summary: format!(
3219 "Attempt {}/{} recorded for {}",
3220 attempt.attempt_index, attempt.attempt_count, attempt.attempt_group_id
3221 ),
3222 detail_path: attempt.patch_path,
3223 },
3224 );
3225 }
3226
3227 if let Some(value) = updates.get("artifacts")
3228 && let Some(items) = value.as_array()
3229 {
3230 for item in items {
3231 let artifact: TaskArtifactRef = serde_json::from_value(item.clone())
3232 .context("Failed to parse artifact task update")?;
3233 push_timeline_entry(
3234 task,
3235 TaskTimelineEntry {
3236 timestamp: now,
3237 kind: "artifact".to_string(),
3238 summary: format!("{}: {}", artifact.label, artifact.summary),
3239 detail_path: Some(artifact.path.clone()),
3240 },
3241 );
3242 task.artifacts.push(artifact);
3243 }
3244 }
3245
3246 if let Some(value) = updates.get("github_event") {
3247 let event: TaskGithubEvent = serde_json::from_value(value.clone())
3248 .context("Failed to parse GitHub task update")?;
3249 push_timeline_entry(
3250 task,
3251 TaskTimelineEntry {
3252 timestamp: now,
3253 kind: "github".to_string(),
3254 summary: format!(
3255 "{} {}#{}: {}",
3256 event.action, event.target, event.number, event.summary
3257 ),
3258 detail_path: None,
3259 },
3260 );
3261 task.github_events.push(event);
3262 }
3263
3264 Ok(())
3265 }
3266
3267 /// Acquire the cross-process task-store lock.
3268 ///
3269 /// Polls with exponential backoff (5ms → 50ms) instead of a flat 5ms
3270 /// interval: under contention the old shape woke ~200×/s for up to its
3271 /// whole five-second deadline (#6211 R7c). The deadline and the busy
3272 /// error are unchanged. What this does not do: it does not add the
3273 /// in-process mutex the issue also suggested — in-process contenders
3274 /// just back off against the same file lock.
3275 async fn lock_store(&self) -> Result<RuntimeProcessOwnerLock> {
3276 let path = self.cfg.data_dir.join("task-store.lock");
3277 let deadline = Instant::now() + Duration::from_secs(5);
3278 let mut wait = Duration::from_millis(5);
3279 loop {
3280 if let Some(mut owner) = RuntimeProcessOwnerLock::try_acquire_file(&path, true)? {
3281 // Lets a contender that times out name this process (#6573).
3282 owner.record_holder();
3283 return Ok(owner);
3284 }
3285 if Instant::now() >= deadline {
3286 return Err(store_busy_error(&path));
3287 }
3288 sleep(wait).await;
3289 wait = (wait * 2).min(Duration::from_millis(50));
3290 }
3291 }
3292
3293 fn refresh_locked(&self, state: &mut ManagerState) -> Result<()> {
3294 // Callers go on to change `state`; a failed persist must not leave a
3295 // diverged view that a read-only caller would trust.
3296 state.read_snapshot = None;
3297 let loaded = load_state(&self.tasks_dir, &self.queue_path)?;
3298 self.apply_loaded_locked(state, loaded)
3299 }
3300
3301 /// Bring `state` up to date for a read-only caller (#6573).
3302 ///
3303 /// Skips the cross-process lock and the full reload while the store is
3304 /// provably unchanged since the last read-only refresh. Several TUIs
3305 /// sharing one data dir each list tasks every 2.5s on their event loop;
3306 /// without this, every idle session took the store lock and re-parsed
3307 /// every task record on each poll, contending with the others.
3308 async fn refresh_for_read(&self, state: &mut ManagerState) -> Result<()> {
3309 if let Some(snapshot) = &state.read_snapshot
3310 && *snapshot == ReadSnapshot::read(&self.tasks_dir, &self.queue_path)
3311 {
3312 return Ok(());
3313 }
3314 let _transaction = self.lock_store().await?;
3315 // Taken under the lock and before the load, so any later write
3316 // changes it.
3317 let snapshot = ReadSnapshot::read(&self.tasks_dir, &self.queue_path);
3318 self.refresh_locked(state)?;
3319 if snapshot.settled(std::time::SystemTime::now()) {
3320 state.read_snapshot = Some(snapshot);
3321 }
3322 Ok(())
3323 }
3324
3325 fn apply_loaded_locked(&self, state: &mut ManagerState, loaded: LoadedTaskState) -> Result<()> {
3326 #[cfg(test)]
3327 self.store_loads
3328 .fetch_add(1, std::sync::atomic::Ordering::Relaxed);
3329 state.tasks = loaded.tasks;
3330 state.queue = loaded.queue;
3331 for (id, events) in &state.pending_events {
3332 let task = state
3333 .tasks
3334 .get_mut(id)
3335 .context("Pending task disappeared")?;
3336 self.require_execution_owner(task)?;
3337 for event in events {
3338 self.apply_event_to_task(task, event.clone())?;
3339 }
3340 }
3341 Ok(())
3342 }
3343
3344 fn require_execution_owner(&self, task: &TaskRecord) -> Result<()> {
3345 if task.execution_scope.as_deref() != Some(self.execution_scope())
3346 || task.execution_generation.as_deref() != Some(&self.execution_lease.generation)
3347 || task.status != TaskStatus::Running
3348 {
3349 bail!("Task execution ownership changed; refusing a stale write");
3350 }
3351 Ok(())
3352 }
3353
3354 fn persist_changed_task_locked(&self, state: &mut ManagerState, id: &str) -> Result<()> {
3355 let task = state.tasks.get(id).context("Changed task is missing")?;
3356 self.persist_task_locked(task)?;
3357 state.pending_events.remove(id);
3358 Ok(())
3359 }
3360
3361 fn recover_dead_executions_locked(&self, state: &mut ManagerState) -> Result<()> {
3362 for task in state.tasks.values_mut() {
3363 // Unknown legacy ownership is preserved, never guessed from the
3364 // visibility owner, model spelling, or current process defaults.
3365 if task.status != TaskStatus::Running {
3366 continue;
3367 }
3368 let (Some(scope), Some(generation)) =
3369 (&task.execution_scope, &task.execution_generation)
3370 else {
3371 continue;
3372 };
3373 let path = execution_lease_path(&self.cfg.data_dir, scope, generation)?;
3374 let Some(_dead_owner) = RuntimeProcessOwnerLock::try_acquire_file(&path, false)? else {
3375 continue;
3376 };
3377 let now = Utc::now();
3378 let duration_ms = task.started_at.and_then(|started| {
3379 u64::try_from(now.signed_duration_since(started).num_milliseconds()).ok()
3380 });
3381 task.status = TaskStatus::Failed;
3382 task.lifecycle_seq = task.lifecycle_seq.saturating_add(1);
3383 task.ended_at = Some(now);
3384 task.duration_ms = duration_ms;
3385 task.terminal_reason = Some(TaskTerminalReason::Failed.as_str().to_string());
3386 task.error =
3387 Some("Interrupted by process restart; prior process is not attached".to_string());
3388 for tool in &mut task.tool_calls {
3389 if tool.status == TaskToolStatus::Running {
3390 tool.status = TaskToolStatus::Failed;
3391 tool.ended_at = Some(now);
3392 tool.duration_ms = duration_ms.or_else(|| {
3393 u64::try_from(
3394 now.signed_duration_since(tool.started_at)
3395 .num_milliseconds(),
3396 )
3397 .ok()
3398 });
3399 }
3400 }
3401 push_timeline_entry(
3402 task,
3403 TaskTimelineEntry {
3404 timestamp: now,
3405 kind: "recovered".to_string(),
3406 summary: "Interrupted by process restart; prior process is not attached"
3407 .to_string(),
3408 detail_path: None,
3409 },
3410 );
3411
3412 self.persist_task_locked(task)?;
3413 }
3414 Ok(())
3415 }
3416
3417 fn persist_queue_locked(&self, queue: &VecDeque<String>) -> Result<()> {
3418 write_json_atomic(
3419 &self.queue_path,
3420 &QueueFile {
3421 queue: queue.iter().cloned().collect(),
3422 },
3423 )
3424 }
3425
3426 fn persist_task_locked(&self, task: &TaskRecord) -> Result<()> {
3427 let path = self.tasks_dir.join(format!("{}.json", task.id));
3428 write_json_atomic(&path, task)
3429 }
3430 }
3431
3432 fn validate_preallocated_task_id(task_id: &str) -> Result<()> {
3433 if task_id.len() != 21
3434 || !task_id.starts_with("task_")
3435 || !task_id[5..].chars().all(|ch| ch.is_ascii_hexdigit())
3436 {
3437 bail!("Invalid preallocated task id: expected task_<16hex>");
3438 }
3439 Ok(())
3440 }
3441
3442 fn read_bound_task_file(path: &Path, task_id: &str) -> Result<Option<TaskRecord>> {
3443 let bytes = match fs::read(path) {
3444 Ok(bytes) => bytes,
3445 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
3446 Err(error) => return Err(error).context("read bound task"),
3447 };
3448 let task: TaskRecord = serde_json::from_slice(&bytes).context("decode bound task")?;
3449 if task.id != task_id || task.schema_version > CURRENT_TASK_SCHEMA_VERSION {
3450 bail!("Bound task identity or schema does not match its durable admission");
3451 }
3452 Ok(Some(task))
3453 }
3454
3455 pub(crate) fn validate_bound_task_request(
3456 task: &TaskRecord,
3457 request: &NewTaskRequest,
3458 ) -> Result<()> {
3459 if task.prompt != request.prompt.trim()
3460 || task.owner_session_id != request.owner_session_id
3461 || task.model_provider != request.model_provider
3462 || task.model_provider_id != request.model_provider_id
3463 || request
3464 .model
3465 .as_ref()
3466 .is_some_and(|value| value != &task.model)
3467 || request
3468 .workspace
3469 .as_ref()
3470 .is_some_and(|value| value != &task.workspace)
3471 || request
3472 .mode
3473 .as_ref()
3474 .is_some_and(|value| value != &task.mode)
3475 || request
3476 .allow_shell
3477 .is_some_and(|value| value != task.allow_shell)
3478 || request
3479 .trust_mode
3480 .is_some_and(|value| value != task.trust_mode)
3481 || request
3482 .permission_posture
3483 .as_deref()
3484 .is_some_and(|value| Some(value) != task.permission_posture.as_deref())
3485 || task.auto_approve != request.auto_approve.unwrap_or(false)
3486 {
3487 bail!("Task admission replay does not match the bound request");
3488 }
3489 Ok(())
3490 }
3491
3492 /// A read-only inventory and reconstructed queue. Execution recovery is a
3493 /// separate transaction requiring proof that a particular generation died.
3494 struct LoadedTaskState {
3495 tasks: HashMap<String, TaskRecord>,
3496 queue: VecDeque<String>,
3497 }
3498
3499 fn load_state(tasks_dir: &Path, queue_path: &Path) -> Result<LoadedTaskState> {
3500 let mut tasks = HashMap::new();
3501 if tasks_dir.exists() {
3502 for entry in fs::read_dir(tasks_dir)
3503 .with_context(|| format!("Failed to read tasks dir {}", tasks_dir.display()))?
3504 {
3505 let entry = entry?;
3506 let path = entry.path();
3507 if path.extension().is_none_or(|ext| ext != "json") {
3508 continue;
3509 }
3510 let content = fs::read_to_string(&path)
3511 .with_context(|| format!("Failed to read task file {}", path.display()))?;
3512 let task: TaskRecord = serde_json::from_str(&content)
3513 .with_context(|| format!("Failed to parse task file {}", path.display()))?;
3514 if task.schema_version > CURRENT_TASK_SCHEMA_VERSION {
3515 bail!(
3516 "Task schema v{} is newer than supported v{}",
3517 task.schema_version,
3518 CURRENT_TASK_SCHEMA_VERSION
3519 );
3520 }
3521 ensure_safe_storage_id("task id", &task.id)?;
3522 if path.file_stem().and_then(|stem| stem.to_str()) != Some(task.id.as_str()) {
3523 bail!("Task record identity differs from its path");
3524 }
3525 if let Some(scope) = &task.execution_scope {
3526 validate_execution_id(scope, 64)?;
3527 }
3528 if let Some(generation) = &task.execution_generation {
3529 validate_execution_id(generation, 32)?;
3530 if task.execution_scope.is_none() {
3531 bail!("Task generation has no execution scope");
3532 }
3533 }
3534 tasks.insert(task.id.clone(), task);
3535 }
3536 }
3537
3538 let mut queue = if queue_path.exists() {
3539 let content = fs::read_to_string(queue_path)
3540 .with_context(|| format!("Failed to read queue file {}", queue_path.display()))?;
3541 let parsed: QueueFile = serde_json::from_str(&content)
3542 .with_context(|| format!("Failed to parse queue file {}", queue_path.display()))?;
3543 VecDeque::from(parsed.queue)
3544 } else {
3545 VecDeque::new()
3546 };
3547
3548 queue.retain(|id| {
3549 tasks
3550 .get(id)
3551 .is_some_and(|task| task.status == TaskStatus::Queued)
3552 });
3553
3554 let known = queue.iter().cloned().collect::<HashSet<_>>();
3555 let mut missing = tasks
3556 .values()
3557 .filter(|task| task.status == TaskStatus::Queued && !known.contains(&task.id))
3558 .map(|task| task.id.clone())
3559 .collect::<Vec<_>>();
3560 missing.sort();
3561 for id in missing {
3562 queue.push_back(id);
3563 }
3564
3565 Ok(LoadedTaskState { tasks, queue })
3566 }
3567
3568 struct EventApplyOutcome {
3569 persisted: bool,
3570 }
3571
3572 fn execution_event_is_progress(event: &TaskExecutionEvent) -> bool {
3573 matches!(
3574 event,
3575 TaskExecutionEvent::MessageDelta { .. }
3576 | TaskExecutionEvent::ToolStarted { .. }
3577 | TaskExecutionEvent::ToolProgress { .. }
3578 | TaskExecutionEvent::ToolHeartbeat
3579 | TaskExecutionEvent::ToolCompleted { .. }
3580 )
3581 }
3582
3583 fn execution_event_persist_urgent(event: &TaskExecutionEvent) -> bool {
3584 !matches!(
3585 event,
3586 TaskExecutionEvent::MessageDelta { .. }
3587 | TaskExecutionEvent::ToolProgress { .. }
3588 | TaskExecutionEvent::RuntimeEvent { .. }
3589 // Liveness-only signal (see the variant doc): it arrives up to
3590 // ~5x/s throughout a silent build, and persisting it would
3591 // rewrite the whole task record on every tick while holding the
3592 // manager-wide state lock.
3593 | TaskExecutionEvent::ToolHeartbeat
3594 )
3595 }
3596
3597 fn timeline_kinds_coalesce(left: &str, right: &str) -> bool {
3598 matches!(
3599 (left, right),
3600 ("message", "message")
3601 | ("tool_progress", "tool_progress")
3602 | ("runtime_event", "runtime_event")
3603 )
3604 }
3605
3606 fn push_timeline_entry(task: &mut TaskRecord, entry: TaskTimelineEntry) {
3607 if let Some(last) = task.timeline.last_mut()
3608 && timeline_kinds_coalesce(last.kind.as_str(), entry.kind.as_str())
3609 {
3610 last.timestamp = entry.timestamp;
3611 last.summary = entry.summary;
3612 last.detail_path = entry.detail_path;
3613 return;
3614 }
3615 task.timeline.push(entry);
3616 trim_task_timeline(&mut task.timeline);
3617 }
3618
3619 fn trim_task_timeline(entries: &mut Vec<TaskTimelineEntry>) {
3620 if entries.len() <= TIMELINE_ENTRY_LIMIT {
3621 return;
3622 }
3623 let overflow = entries.len() - TIMELINE_ENTRY_LIMIT;
3624 let start = TIMELINE_HEAD_KEEP.min(entries.len().saturating_sub(overflow + 1));
3625 let end = start + overflow;
3626 if start >= end || end > entries.len() {
3627 entries.truncate(TIMELINE_ENTRY_LIMIT);
3628 return;
3629 }
3630 entries.drain(start..end);
3631 let omitted = TaskTimelineEntry {
3632 timestamp: Utc::now(),
3633 kind: "omitted".to_string(),
3634 summary: format!("{overflow} earlier events omitted to bound storage"),
3635 detail_path: None,
3636 };
3637 if entries
3638 .get(start)
3639 .is_none_or(|entry| entry.kind != "omitted")
3640 {
3641 entries.insert(start, omitted);
3642 }
3643 if entries.len() > TIMELINE_ENTRY_LIMIT {
3644 let extra = entries.len() - TIMELINE_ENTRY_LIMIT;
3645 let drop_at = (start + 1).min(entries.len().saturating_sub(1));
3646 let drop_end = (drop_at + extra).min(entries.len());
3647 if drop_at < drop_end {
3648 entries.drain(drop_at..drop_end);
3649 } else {
3650 entries.truncate(TIMELINE_ENTRY_LIMIT);
3651 }
3652 }
3653 }
3654
3655 fn resolve_task_id_visible_to(
3656 tasks: &HashMap<String, TaskRecord>,
3657 id_or_prefix: &str,
3658 owner_session_id: Option<&str>,
3659 ) -> Result<String> {
3660 let visible = |record: &TaskRecord| {
3661 owner_session_id.is_none_or(|owner_session_id| {
3662 record.owner_session_id.as_deref() == Some(owner_session_id)
3663 })
3664 };
3665 if tasks.get(id_or_prefix).is_some_and(visible) {
3666 return Ok(id_or_prefix.to_string());
3667 }
3668 let matches = tasks
3669 .iter()
3670 .filter(|(id, record)| id.starts_with(id_or_prefix) && visible(record))
3671 .map(|(id, _)| id)
3672 .cloned()
3673 .collect::<Vec<_>>();
3674 match matches.len() {
3675 0 => bail!("Task not found: {id_or_prefix}"),
3676 1 => Ok(matches[0].clone()),
3677 _ => bail!(
3678 "Ambiguous task prefix '{}': matches {} tasks",
3679 id_or_prefix,
3680 matches.len()
3681 ),
3682 }
3683 }
3684
3685 fn resolve_task_id_visible_to_operator(
3686 tasks: &HashMap<String, TaskRecord>,
3687 id_or_prefix: &str,
3688 owner_session_id: &str,
3689 execution_scope: &str,
3690 ) -> Result<String> {
3691 let visible = |record: &TaskRecord| {
3692 record.owner_session_id.as_deref() == Some(owner_session_id)
3693 || record.owner_session_id.is_none()
3694 && !execution_scope.is_empty()
3695 && record.execution_scope.as_deref() == Some(execution_scope)
3696 };
3697 if tasks.get(id_or_prefix).is_some_and(visible) {
3698 return Ok(id_or_prefix.to_string());
3699 }
3700 let matches = tasks
3701 .iter()
3702 .filter(|(id, record)| id.starts_with(id_or_prefix) && visible(record))
3703 .map(|(id, _)| id)
3704 .cloned()
3705 .collect::<Vec<_>>();
3706 match matches.len() {
3707 0 => bail!("Task not found: {id_or_prefix}"),
3708 1 => Ok(matches[0].clone()),
3709 _ => bail!(
3710 "Ambiguous task prefix '{}': matches {} tasks",
3711 id_or_prefix,
3712 matches.len()
3713 ),
3714 }
3715 }
3716
3717 fn resolve_task_id(tasks: &HashMap<String, TaskRecord>, id_or_prefix: &str) -> Result<String> {
3718 resolve_task_id_visible_to(tasks, id_or_prefix, None)
3719 }
3720
3721 fn summarize_json(value: &Value) -> Option<String> {
3722 let text = serde_json::to_string(value).ok()?;
3723 Some(summarize_text(&text, TIMELINE_SUMMARY_LIMIT))
3724 }
3725
3726 fn summarize_text(text: &str, limit: usize) -> String {
3727 let take = limit.saturating_sub(3);
3728 let mut count = 0;
3729 let mut out = String::new();
3730 for ch in text.chars() {
3731 if count >= take {
3732 out.push_str("...");
3733 return out;
3734 }
3735 if ch.is_control() && ch != '\n' && ch != '\t' {
3736 continue;
3737 }
3738 out.push(ch);
3739 count += 1;
3740 }
3741 out
3742 }
3743
3744 fn ensure_safe_storage_id(kind: &str, value: &str) -> Result<()> {
3745 let mut components = Path::new(value).components();
3746 let Some(component) = components.next() else {
3747 bail!("{kind} must not be empty");
3748 };
3749 if components.next().is_some() || !matches!(component, std::path::Component::Normal(_)) {
3750 bail!("{kind} must be a single path component");
3751 }
3752 Ok(())
3753 }
3754
3755 fn sanitize_filename(input: &str) -> String {
3756 let mut out = String::new();
3757 for ch in input.chars() {
3758 if ch.is_ascii_alphanumeric() || ch == '_' || ch == '-' {
3759 out.push(ch);
3760 } else {
3761 out.push('_');
3762 }
3763 }
3764 if out.is_empty() {
3765 "artifact".to_string()
3766 } else {
3767 out
3768 }
3769 }
3770
3771 fn duration_ms(start: DateTime<Utc>, end: DateTime<Utc>) -> u64 {
3772 let millis = (end - start).num_milliseconds();
3773 if millis.is_negative() {
3774 0
3775 } else {
3776 u64::try_from(millis).unwrap_or(u64::MAX)
3777 }
3778 }
3779
3780 fn write_json_atomic<T: Serialize>(path: &Path, value: &T) -> Result<()> {
3781 if let Some(parent) = path.parent() {
3782 fs::create_dir_all(parent)
3783 .with_context(|| format!("Failed to create directory {}", parent.display()))?;
3784 }
3785 let payload = serde_json::to_string_pretty(value)?;
3786 crate::utils::write_atomic(path, payload.as_bytes())
3787 .with_context(|| format!("Failed to write {}", path.display()))
3788 }
3789
3790 /// Default task manager data location (`~/.codewhale/tasks`, or legacy
3791 /// `~/.deepseek/tasks` when only the legacy directory exists).
3792 #[must_use]
3793 pub fn default_tasks_dir() -> PathBuf {
3794 for var in ["CODEWHALE_TASKS_DIR", "DEEPSEEK_TASKS_DIR"] {
3795 if let Ok(path) = std::env::var(var)
3796 && !path.trim().is_empty()
3797 {
3798 return PathBuf::from(path);
3799 }
3800 }
3801 if let Some(home) = codewhale_paths::codewhale_home_override().ok().flatten() {
3802 return home.join("tasks");
3803 }
3804 codewhale_paths::user_home()
3805 .map(|home| default_tasks_dir_for_home(&home))
3806 .unwrap_or_else(|| PathBuf::from(".codewhale").join("tasks"))
3807 }
3808
3809 fn default_tasks_dir_for_home(home: &Path) -> PathBuf {
3810 let primary = home.join(".codewhale").join("tasks");
3811 if primary.is_dir() {
3812 return primary;
3813 }
3814 let legacy = home.join(".deepseek").join("tasks");
3815 if legacy.is_dir() {
3816 return legacy;
3817 }
3818 primary
3819 }
3820
3821 /// Wait for a task to reach a terminal status (tests and API helpers).
3822 #[cfg(test)]
3823 pub async fn wait_for_terminal_state(
3824 manager: &TaskManager,
3825 task_id: &str,
3826 timeout: StdDuration,
3827 ) -> Result<TaskRecord> {
3828 let deadline = std::time::Instant::now() + timeout;
3829 loop {
3830 let task = manager.get_task(task_id).await?;
3831 if task.status.is_terminal() {
3832 return Ok(task);
3833 }
3834 if std::time::Instant::now() >= deadline {
3835 bail!("Timed out waiting for task {task_id}");
3836 }
3837 sleep(StdDuration::from_millis(50)).await;
3838 }
3839 }
3840
3841 #[cfg(test)]
3842 mod tests {
3843 use super::*;
3844 use crate::test_support::{EnvVarGuard, lock_test_env};
3845 use std::fs;
3846 use std::sync::atomic::{AtomicUsize, Ordering};
3847 use tokio::time::Duration;
3848
3849 struct MockExecutor;
3850
3851 /// Poll until the task is claimed as `Running`, or fail at `timeout`.
3852 ///
3853 /// A worker claims the task and installs its cancel token under one state
3854 /// lock, so observing `Running` means a cancel or shutdown now reaches a
3855 /// live executor rather than a still-queued record.
3856 async fn wait_for_running(
3857 manager: &TaskManager,
3858 task_id: &str,
3859 timeout: Duration,
3860 ) -> Result<TaskRecord> {
3861 let deadline = std::time::Instant::now() + timeout;
3862 loop {
3863 let task = manager.get_task(task_id).await?;
3864 if task.status == TaskStatus::Running {
3865 return Ok(task);
3866 }
3867 if task.status.is_terminal() || std::time::Instant::now() >= deadline {
3868 bail!("task {task_id} never started running: {task:?}");
3869 }
3870 sleep(Duration::from_millis(5)).await;
3871 }
3872 }
3873
3874 fn provider_default_model_cases() -> Vec<(&'static str, Config, &'static str)> {
3875 let deepseek = Config {
3876 provider: Some("deepseek".to_string()),
3877 default_text_model: Some("deepseek-v4-flash".to_string()),
3878 ..Config::default()
3879 };
3880
3881 let zai = Config {
3882 provider: Some("zai".to_string()),
3883 // Exercise provider-aware rejection of a stale DeepSeek root default.
3884 default_text_model: Some(crate::config::DEFAULT_TEXT_MODEL.to_string()),
3885 ..Config::default()
3886 };
3887
3888 let mut custom_providers = crate::config::ProvidersConfig::default();
3889 custom_providers.custom.insert(
3890 "acme".to_string(),
3891 crate::config::ProviderConfig {
3892 base_url: Some("http://127.0.0.1:1/v1".to_string()),
3893 model: Some("acme-coder".to_string()),
3894 kind: Some("openai-compatible".to_string()),
3895 ..crate::config::ProviderConfig::default()
3896 },
3897 );
3898 let custom = Config {
3899 provider: Some("acme".to_string()),
3900 providers: Some(custom_providers),
3901 ..Config::default()
3902 };
3903
3904 vec![
3905 ("deepseek", deepseek, "deepseek-v4-flash"),
3906 ("zai", zai, crate::config::DEFAULT_ZAI_MODEL),
3907 ("custom", custom, "acme-coder"),
3908 ]
3909 }
3910
3911 #[test]
3912 fn task_manager_config_uses_the_active_provider_default() {
3913 for (label, config, expected) in provider_default_model_cases() {
3914 let task_config =
3915 TaskManagerConfig::from_runtime(&config, PathBuf::from("."), None, Some(1));
3916 assert_eq!(
3917 task_config.default_model, expected,
3918 "{label} durable task default"
3919 );
3920 }
3921 }
3922
3923 #[async_trait]
3924 impl TaskExecutor for MockExecutor {
3925 async fn execute(
3926 &self,
3927 task: ExecutionTask,
3928 events: mpsc::Sender<TaskExecutionEvent>,
3929 cancel: CancellationToken,
3930 ) -> TaskExecutionResult {
3931 let _ = events
3932 .send(TaskExecutionEvent::Status {
3933 message: format!("running {}", task.id),
3934 })
3935 .await;
3936 let _ = events
3937 .send(TaskExecutionEvent::ThreadLinked {
3938 thread_id: "thr_test".to_string(),
3939 turn_id: "turn_test".to_string(),
3940 })
3941 .await;
3942 let _ = events
3943 .send(TaskExecutionEvent::ToolStarted {
3944 id: "tool_1".to_string(),
3945 name: "read_file".to_string(),
3946 input: serde_json::json!({ "path": "README.md" }),
3947 })
3948 .await;
3949 sleep(Duration::from_millis(50)).await;
3950 if cancel.is_cancelled() {
3951 return TaskExecutionResult {
3952 status: TaskStatus::Canceled,
3953 result_text: None,
3954 error: None,
3955 terminal_reason: TaskTerminalReason::Canceled,
3956 };
3957 }
3958 let _ = events
3959 .send(TaskExecutionEvent::ToolCompleted {
3960 id: "tool_1".to_string(),
3961 name: "read_file".to_string(),
3962 success: true,
3963 output: "read ok".to_string(),
3964 metadata: Some(serde_json::json!({
3965 "duration_ms": 10,
3966 "task_updates": {
3967 "checklist": {
3968 "items": [
3969 { "id": 1, "content": "read fixture", "status": "in_progress" }
3970 ],
3971 "completion_pct": 0,
3972 "in_progress_id": 1,
3973 "updated_at": null
3974 }
3975 }
3976 })),
3977 })
3978 .await;
3979 TaskExecutionResult {
3980 status: TaskStatus::Completed,
3981 result_text: Some("done".to_string()),
3982 error: None,
3983 terminal_reason: TaskTerminalReason::Completed,
3984 }
3985 }
3986 }
3987
3988 fn test_config(root: PathBuf) -> TaskManagerConfig {
3989 TaskManagerConfig {
3990 data_dir: root,
3991 worker_count: 1,
3992 default_workspace: PathBuf::from("."),
3993 default_model: "deepseek-v4-flash".to_string(),
3994 default_mode: "agent".to_string(),
3995 allow_shell: false,
3996 trust_mode: false,
3997 execution_limits: TaskExecutionLimits::default(),
3998 }
3999 }
4000
4001 fn short_test_config(root: PathBuf) -> TaskManagerConfig {
4002 TaskManagerConfig {
4003 execution_limits: TaskExecutionLimits::short_for_tests(),
4004 ..test_config(root)
4005 }
4006 }
4007
4008 fn wall_timeout_test_config(root: PathBuf) -> TaskManagerConfig {
4009 let mut config = short_test_config(root);
4010 config.execution_limits.idle_progress = config
4011 .execution_limits
4012 .wall_time
4013 .saturating_add(config.execution_limits.cancel_grace);
4014 config
4015 }
4016
4017 /// Wait until every idle worker of `managers` has stopped polling.
4018 ///
4019 /// A worker that saw a fresh queue stamp (another manager starting writes
4020 /// one) re-reads it once, `STORE_IDLE_POLL_INTERVAL` later and on a
4021 /// `STORE_REFRESH_INTERVAL` tick boundary, to confirm the stamp has settled.
4022 /// A fixed `interval + 300ms` sleep leaves under one tick of margin, so a
4023 /// slow runner (Windows CI) could land that confirming read inside the
4024 /// caller's "no reloads" window. A window longer than the interval in which
4025 /// no worker reloaded proves every worker has settled; a worker that polls
4026 /// forever never produces one, so the caller's assertion still catches it.
4027 async fn settle_idle_store_polling(managers: &[&TaskManager]) {
4028 for _ in 0..5 {
4029 for manager in managers {
4030 manager.store_loads.store(0, Ordering::Relaxed);
4031 }
4032 sleep(STORE_IDLE_POLL_INTERVAL + Duration::from_millis(300)).await;
4033 let loads: usize = managers
4034 .iter()
4035 .map(|manager| manager.store_loads.load(Ordering::Relaxed))
4036 .sum();
4037 if loads == 0 {
4038 return;
4039 }
4040 }
4041 }
4042
4043 /// #6573: idle managers sharing one data dir must not reload the store
4044 /// (and take its cross-process lock) every 200ms per worker.
4045 #[tokio::test]
4046 async fn idle_managers_sharing_a_store_do_not_poll_it_continuously() -> Result<()> {
4047 let root = tempfile::tempdir()?;
4048 let config = || TaskManagerConfig {
4049 worker_count: 2,
4050 ..test_config(root.path().to_path_buf())
4051 };
4052 let first =
4053 TaskManager::start_with_executor_in_scope(config(), Arc::new(MockExecutor), "first")
4054 .await?;
4055 let second =
4056 TaskManager::start_with_executor_in_scope(config(), Arc::new(MockExecutor), "second")
4057 .await?;
4058 settle_idle_store_polling(&[&first, &second]).await;
4059 for manager in [&first, &second] {
4060 manager.store_loads.store(0, Ordering::Relaxed);
4061 }
4062
4063 let window = Duration::from_secs(2);
4064 sleep(window).await;
4065 let loads =
4066 first.store_loads.load(Ordering::Relaxed) + second.store_loads.load(Ordering::Relaxed);
4067 // Once the queue stamp settles, unchanged stores never reload just
4068 // because time passed. The former 2s fallback reloads in this window.
4069 assert!(
4070 loads == 0,
4071 "idle workers reloaded the shared store {loads} times in {window:?}"
4072 );
4073
4074 // An in-process submission still wakes a worker immediately.
4075 let task = first
4076 .add_task(NewTaskRequest::from_prompt("wake an idle worker"))
4077 .await?;
4078 let finished = wait_for_terminal_state(&first, &task.id, Duration::from_secs(1)).await?;
4079 assert_eq!(finished.status, TaskStatus::Completed);
4080
4081 first.shutdown_and_wait().await?;
4082 second.shutdown_and_wait().await?;
4083 Ok(())
4084 }
4085
4086 /// #6573: the TUI task panel lists tasks every 2.5s. An unchanged store
4087 /// must be answered from memory, without the cross-process lock or a
4088 /// reload, while a write from another process is still seen.
4089 #[tokio::test]
4090 async fn repeated_listing_of_an_unchanged_store_does_not_reload_it() -> Result<()> {
4091 let root = tempfile::tempdir()?;
4092 let lister = TaskManager::start_with_executor_in_scope(
4093 test_config(root.path().to_path_buf()),
4094 Arc::new(MockExecutor),
4095 "lister",
4096 )
4097 .await?;
4098 let writer = TaskManager::start_with_executor_in_scope(
4099 test_config(root.path().to_path_buf()),
4100 Arc::new(MockExecutor),
4101 "writer",
4102 )
4103 .await?;
4104 let first = writer
4105 .add_task(NewTaskRequest::from_prompt("first"))
4106 .await?;
4107 wait_for_terminal_state(&writer, &first.id, Duration::from_secs(2)).await?;
4108 // Let the store's mtimes settle so a snapshot can be trusted.
4109 sleep(READ_SNAPSHOT_SETTLE + Duration::from_millis(300)).await;
4110 assert_eq!(lister.list_tasks(None).await?.len(), 1);
4111
4112 lister.store_loads.store(0, Ordering::Relaxed);
4113 for _ in 0..20 {
4114 assert_eq!(lister.list_tasks(None).await?.len(), 1);
4115 assert_eq!(lister.counts().await?.completed, 1);
4116 }
4117 let loads = lister.store_loads.load(Ordering::Relaxed);
4118 // Only the idle worker's fallback poll may reload in this window.
4119 assert!(
4120 loads <= 1,
4121 "40 read-only calls on an unchanged store reloaded it {loads} times"
4122 );
4123
4124 let second = writer
4125 .add_task(NewTaskRequest::from_prompt("second"))
4126 .await?;
4127 assert!(
4128 lister
4129 .list_tasks(None)
4130 .await?
4131 .iter()
4132 .any(|task| task.id == second.id),
4133 "a write from another process must invalidate the read snapshot"
4134 );
4135
4136 lister.shutdown_and_wait().await?;
4137 writer.shutdown_and_wait().await?;
4138 Ok(())
4139 }
4140
4141 /// #6573: a task running in one process rewrites its record on every
4142 /// persisted event. Those writes must not make idle workers in another
4143 /// process reload the shared store on every tick.
4144 #[tokio::test]
4145 async fn running_task_flushes_do_not_wake_idle_workers_elsewhere() -> Result<()> {
4146 struct StatusStreamExecutor;
4147
4148 #[async_trait]
4149 impl TaskExecutor for StatusStreamExecutor {
4150 async fn execute(
4151 &self,
4152 _task: ExecutionTask,
4153 events: mpsc::Sender<TaskExecutionEvent>,
4154 cancel: CancellationToken,
4155 ) -> TaskExecutionResult {
4156 for chunk in 0.. {
4157 if cancel.is_cancelled() {
4158 break;
4159 }
4160 let _ = events
4161 .send(TaskExecutionEvent::Status {
4162 message: format!("step {chunk}"),
4163 })
4164 .await;
4165 // Status events persist the task record immediately.
4166 sleep(Duration::from_millis(50)).await;
4167 }
4168 TaskExecutionResult::from_reason(TaskTerminalReason::Canceled, None)
4169 }
4170 }
4171
4172 let root = tempfile::tempdir()?;
4173 let busy = TaskManager::start_with_executor_in_scope(
4174 test_config(root.path().to_path_buf()),
4175 Arc::new(StatusStreamExecutor),
4176 "busy",
4177 )
4178 .await?;
4179 let task = busy
4180 .add_task(NewTaskRequest::from_prompt("stream status"))
4181 .await?;
4182 wait_for_running(&busy, &task.id, Duration::from_secs(2)).await?;
4183
4184 let idle = TaskManager::start_with_executor_in_scope(
4185 TaskManagerConfig {
4186 worker_count: 2,
4187 ..test_config(root.path().to_path_buf())
4188 },
4189 Arc::new(MockExecutor),
4190 "idle",
4191 )
4192 .await?;
4193 settle_idle_store_polling(&[&idle]).await;
4194 idle.store_loads.store(0, Ordering::Relaxed);
4195 let record = root.path().join("tasks").join(format!("{}.json", task.id));
4196 let flushed_before = fs::metadata(&record)?.modified()?;
4197
4198 let window = Duration::from_secs(2);
4199 sleep(window).await;
4200 let loads = idle.store_loads.load(Ordering::Relaxed);
4201 assert_ne!(
4202 fs::metadata(&record)?.modified()?,
4203 flushed_before,
4204 "the running task should have flushed its record during the window"
4205 );
4206 // Task-event writes never change queue eligibility, so settled idle
4207 // workers do not need to re-read them.
4208 assert!(
4209 loads == 0,
4210 "idle workers reloaded the shared store {loads} times in {window:?} while a task ran elsewhere"
4211 );
4212
4213 busy.cancel_task(&task.id).await?;
4214 busy.shutdown_and_wait().await?;
4215 idle.shutdown_and_wait().await?;
4216 Ok(())
4217 }
4218
4219 #[test]
4220 fn claim_schedule_backs_off_after_failures_and_retries_on_queue_change() {
4221 let fingerprint = |len| StoreFingerprint(Some((std::time::SystemTime::UNIX_EPOCH, len, 1)));
4222 let start = Instant::now();
4223 let mut schedule = ClaimSchedule::new(start);
4224 assert!(schedule.should_claim(&fingerprint(1), start, false));
4225
4226 // Consecutive failures double the retry delay up to the ceiling.
4227 let mut now = start;
4228 let mut delays = Vec::new();
4229 for _ in 0..8 {
4230 let delay = schedule.claim_failed(fingerprint(1), now);
4231 delays.push(delay.as_millis());
4232 assert!(!schedule.should_claim(&fingerprint(1), now + delay / 2, false));
4233 now += delay;
4234 assert!(schedule.should_claim(&fingerprint(1), now, false));
4235 }
4236 assert_eq!(delays, [200, 400, 800, 1600, 3200, 6400, 8000, 8000]);
4237
4238 // During a backoff, a queue change or an in-process wakeup retries
4239 // at once.
4240 let delay = schedule.claim_failed(fingerprint(1), now);
4241 assert_eq!(delay, STORE_BUSY_BACKOFF_MAX);
4242 assert!(schedule.should_claim(&fingerprint(2), now, false));
4243 assert!(schedule.should_claim(&fingerprint(1), now, true));
4244
4245 // An empty claim on a settled queue resets backoff without scheduling
4246 // another scan, however long the worker stays idle.
4247 schedule.found_nothing(fingerprint(2), now);
4248 assert!(!schedule.should_claim(&fingerprint(2), now, false));
4249 assert!(!schedule.should_claim(&fingerprint(2), now + Duration::from_secs(86_400), false));
4250 assert_eq!(
4251 schedule.claim_failed(fingerprint(2), now),
4252 STORE_REFRESH_INTERVAL
4253 );
4254 }
4255
4256 #[test]
4257 fn claim_schedule_retries_unknown_or_recent_queue_metadata() {
4258 let now = Instant::now();
4259 for fingerprint in [
4260 StoreFingerprint(None),
4261 StoreFingerprint(Some((std::time::SystemTime::now(), 1, 1))),
4262 ] {
4263 let mut schedule = ClaimSchedule::new(now);
4264 schedule.found_nothing(fingerprint.clone(), now);
4265 assert!(!schedule.should_claim(&fingerprint, now, false));
4266 assert!(schedule.should_claim(&fingerprint, now + STORE_IDLE_POLL_INTERVAL, false));
4267 }
4268 }
4269
4270 #[tokio::test]
4271 async fn settled_idle_workers_notice_external_queue_writes_without_notify() -> Result<()> {
4272 let root = tempfile::tempdir()?;
4273 let executions = Arc::new(AtomicUsize::new(0));
4274 let manager = TaskManager::start_with_executor(
4275 TaskManagerConfig {
4276 worker_count: 2,
4277 ..test_config(root.path().to_path_buf())
4278 },
4279 Arc::new(AdmissionCountingExecutor(executions.clone())),
4280 )
4281 .await?;
4282 sleep(STORE_IDLE_POLL_INTERVAL + Duration::from_millis(300)).await;
4283
4284 let mut task = sample_task_record();
4285 task.status = TaskStatus::Queued;
4286 task.started_at = None;
4287 let mut second = task.clone();
4288 second.id = "task_0123456789abcdee".into();
4289 {
4290 // Simulate another process's durable admission, under the real
4291 // cross-process transaction lock. No local state or Notify changes.
4292 let _transaction = manager.lock_store().await?;
4293 manager.persist_task_locked(&task)?;
4294 manager.persist_task_locked(&second)?;
4295 manager.persist_queue_locked(&VecDeque::from([task.id.clone(), second.id.clone()]))?;
4296 }
4297 for id in [&task.id, &second.id] {
4298 // Settled idle workers recheck the queue every STORE_SETTLED_NAP.
4299 let done =
4300 wait_for_terminal_state(&manager, id, STORE_SETTLED_NAP + Duration::from_secs(2))
4301 .await?;
4302 assert_eq!(done.status, TaskStatus::Completed);
4303 }
4304 assert_eq!(executions.load(Ordering::SeqCst), 2);
4305 manager.shutdown_and_wait().await?;
4306 Ok(())
4307 }
4308
4309 /// #6728: with nothing claimable, each worker rechecks the queue about once
4310 /// per STORE_SETTLED_NAP rather than every STORE_REFRESH_INTERVAL. The
4311 /// bound is an upper limit, so a slow machine only lowers the count.
4312 #[tokio::test]
4313 async fn settled_idle_workers_check_the_queue_at_the_settled_nap() -> Result<()> {
4314 let root = tempfile::tempdir()?;
4315 let workers = 2;
4316 let manager = TaskManager::start_with_executor(
4317 TaskManagerConfig {
4318 worker_count: workers,
4319 ..test_config(root.path().to_path_buf())
4320 },
4321 Arc::new(MockExecutor),
4322 )
4323 .await?;
4324 sleep(STORE_IDLE_POLL_INTERVAL + Duration::from_millis(300)).await;
4325 manager.fingerprint_reads.store(0, Ordering::Relaxed);
4326 manager.store_loads.store(0, Ordering::Relaxed);
4327
4328 let window = Duration::from_secs(3);
4329 sleep(window).await;
4330 let reads = manager.fingerprint_reads.load(Ordering::Relaxed);
4331 let loads = manager.store_loads.load(Ordering::Relaxed);
4332 let ticks = (window.as_millis() / STORE_SETTLED_NAP.as_millis()) as usize + 1;
4333 assert!(
4334 reads <= workers * ticks,
4335 "{workers} idle workers read the queue fingerprint {reads} times in {window:?}; \
4336 expected at most {}",
4337 workers * ticks
4338 );
4339 assert_eq!(loads, 0, "idle workers reloaded a settled store");
4340 manager.shutdown_and_wait().await?;
4341 Ok(())
4342 }
4343
4344 #[test]
4345 fn claim_schedule_naps_long_only_when_nothing_is_claimable() {
4346 let now = Instant::now();
4347 let settled = StoreFingerprint(Some((std::time::SystemTime::UNIX_EPOCH, 1, 1)));
4348 let mut schedule = ClaimSchedule::new(now);
4349 // A fresh schedule has a claim pending.
4350 assert_eq!(schedule.nap(), STORE_REFRESH_INTERVAL);
4351 schedule.found_nothing(settled.clone(), now);
4352 assert_eq!(schedule.nap(), STORE_SETTLED_NAP);
4353 // A failed claim keeps the short tick so a queue change is seen early.
4354 schedule.claim_failed(settled.clone(), now);
4355 assert_eq!(schedule.nap(), STORE_REFRESH_INTERVAL);
4356 // Recent or missing metadata retries on a deadline, so it stays short.
4357 schedule.found_nothing(StoreFingerprint(None), now);
4358 assert_eq!(schedule.nap(), STORE_REFRESH_INTERVAL);
4359 schedule.claimed_task(now);
4360 assert_eq!(schedule.nap(), STORE_REFRESH_INTERVAL);
4361 }
4362
4363 /// #6573: a busy store names the holder; an unreadable record falls back
4364 /// to the bare message without failing.
4365 #[test]
4366 fn busy_store_error_names_the_lock_holder_when_known() -> Result<()> {
4367 let root = tempfile::tempdir()?;
4368 let path = root.path().join("task-store.lock");
4369 let mut holder =
4370 RuntimeProcessOwnerLock::try_acquire_file(&path, true)?.context("first acquire")?;
4371 holder.record_holder();
4372 assert!(RuntimeProcessOwnerLock::try_acquire_file(&path, true)?.is_none());
4373
4374 let message = store_busy_error(&path).to_string();
4375 assert!(
4376 message.contains(&format!("held by pid {}", std::process::id())),
4377 "{message}"
4378 );
4379
4380 // Releasing clears the record, so a later error cannot name a stale pid.
4381 drop(holder);
4382 assert_eq!(
4383 store_busy_error(&path).to_string(),
4384 "Task store is busy; state is unavailable"
4385 );
4386 // Garbage and a missing file also fall back.
4387 fs::write(&path, "pid=oops since_ms=\n")?;
4388 assert!(RuntimeProcessOwnerLock::read_holder(&path).is_none());
4389 fs::remove_file(&path)?;
4390 assert!(RuntimeProcessOwnerLock::read_holder(&path).is_none());
4391 Ok(())
4392 }
4393
4394 #[tokio::test]
4395 async fn running_task_observes_durable_cancel_without_a_queue_change() -> Result<()> {
4396 let root = tempfile::tempdir()?;
4397 let manager = TaskManager::start_with_executor(
4398 test_config(root.path().to_path_buf()),
4399 Arc::new(CooperativeIdleCancelExecutor),
4400 )
4401 .await?;
4402 let task = manager
4403 .add_task(NewTaskRequest::from_prompt("durable cancel"))
4404 .await?;
4405 wait_for_running(&manager, &task.id, Duration::from_secs(2)).await?;
4406 let queue_before = StoreFingerprint::read(&manager.queue_path);
4407 {
4408 let _transaction = manager.lock_store().await?;
4409 let mut task = manager
4410 .read_bound_task(&task.id)?
4411 .context("running fixture")?;
4412 task.lifecycle_seq += 1;
4413 task.cancel_requested_seq = task.lifecycle_seq;
4414 manager.persist_task_locked(&task)?;
4415 }
4416 let done = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(2)).await?;
4417 assert_eq!(done.status, TaskStatus::Canceled);
4418 assert_eq!(StoreFingerprint::read(&manager.queue_path), queue_before);
4419 manager.shutdown_and_wait().await?;
4420 Ok(())
4421 }
4422
4423 #[tokio::test]
4424 async fn persists_and_recovers_task_records() -> Result<()> {
4425 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
4426 let manager =
4427 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
4428 .await?;
4429
4430 let task = manager
4431 .add_task(NewTaskRequest {
4432 owner_session_id: Some("session-persist".to_string()),
4433 ..NewTaskRequest::from_prompt("test persistence")
4434 })
4435 .await?;
4436 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
4437 assert_eq!(finished.status, TaskStatus::Completed);
4438 assert_eq!(finished.thread_id.as_deref(), Some("thr_test"));
4439 assert_eq!(finished.turn_id.as_deref(), Some("turn_test"));
4440 assert_eq!(finished.checklist.items.len(), 1);
4441 assert_eq!(finished.checklist.in_progress_id, Some(1));
4442 assert!(
4443 finished.lifecycle_seq >= 3,
4444 "queued, running, and terminal owner transitions must advance the sequence"
4445 );
4446
4447 manager.shutdown_and_wait().await?;
4448 drop(manager);
4449
4450 let recovered =
4451 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
4452 .await?;
4453 let loaded = recovered.get_task(&task.id).await?;
4454 assert_eq!(loaded.status, TaskStatus::Completed);
4455 assert_eq!(
4456 loaded.owner_session_id.as_deref(),
4457 Some("session-persist"),
4458 "session ownership should survive persistence and restart"
4459 );
4460 assert!(!loaded.timeline.is_empty());
4461 assert_eq!(loaded.checklist.items[0].content, "read fixture");
4462 Ok(())
4463 }
4464
4465 struct AdmissionCountingExecutor(Arc<AtomicUsize>);
4466
4467 #[async_trait]
4468 impl TaskExecutor for AdmissionCountingExecutor {
4469 async fn execute(
4470 &self,
4471 _task: ExecutionTask,
4472 _events: mpsc::Sender<TaskExecutionEvent>,
4473 _cancel: CancellationToken,
4474 ) -> TaskExecutionResult {
4475 self.0.fetch_add(1, Ordering::SeqCst);
4476 TaskExecutionResult {
4477 status: TaskStatus::Completed,
4478 result_text: Some("admission fixture completed".into()),
4479 error: None,
4480 terminal_reason: TaskTerminalReason::Completed,
4481 }
4482 }
4483 }
4484
4485 #[tokio::test]
4486 async fn interrupted_task_stage_preserves_resolved_request_and_executes_once() -> Result<()> {
4487 for queue_was_written in [false, true] {
4488 let root = tempfile::tempdir()?;
4489 let tasks_dir = root.path().join("tasks");
4490 fs::create_dir_all(&tasks_dir)?;
4491 let mut staged = sample_task_record();
4492 staged.status = TaskStatus::Queued;
4493 staged.started_at = None;
4494 staged.model = "staged-model".into();
4495 staged.workspace = root.path().join("staged-workspace");
4496 fs::create_dir(&staged.workspace)?;
4497 let mut request = NewTaskRequest::from_task(&staged);
4498 request.model = None;
4499 request.workspace = None;
4500 request.mode = None;
4501 request.allow_shell = None;
4502 request.trust_mode = None;
4503 let staged_path = tasks_dir.join(format!(".{}.json.pending", staged.id));
4504 write_json_atomic(&staged_path, &staged)?;
4505 if queue_was_written {
4506 write_json_atomic(
4507 &root.path().join("queue.json"),
4508 &QueueFile {
4509 queue: vec![staged.id.clone()],
4510 },
4511 )?;
4512 }
4513 let executions = Arc::new(AtomicUsize::new(0));
4514 let mut config = test_config(root.path().to_path_buf());
4515 config.default_model = "new-default-model".into();
4516 config.default_workspace = root.path().join("new-workspace");
4517 config.default_mode = "plan".into();
4518 config.allow_shell = true;
4519 config.trust_mode = true;
4520 let manager = TaskManager::start_with_executor(
4521 config,
4522 Arc::new(AdmissionCountingExecutor(executions.clone())),
4523 )
4524 .await?;
4525 // Interrupt recovery while its admission is waiting for the queue
4526 // lock. The resolved intent must remain durable for another retry.
4527 let stage_before = fs::read(&staged_path)?;
4528 let queue_guard = manager.state.lock().await;
4529 let mut recovery =
4530 Box::pin(manager.recover_task_admission(request.clone(), staged.id.clone()));
4531 assert!(
4532 tokio::time::timeout(Duration::from_millis(25), &mut recovery)
4533 .await
4534 .is_err()
4535 );
4536 assert_eq!(
4537 fs::read(&staged_path)?,
4538 stage_before,
4539 "interrupted recovery must preserve its resolved staged intent"
4540 );
4541 assert!(manager.read_bound_task(&staged.id)?.is_none());
4542 assert_eq!(executions.load(Ordering::SeqCst), 0);
4543 drop(recovery);
4544 drop(queue_guard);
4545 let admitted = manager
4546 .recover_task_admission(request.clone(), staged.id.clone())
4547 .await?;
4548 assert_eq!(admitted.id, staged.id);
4549 assert_eq!(admitted.model, staged.model);
4550 assert_eq!(admitted.workspace, staged.workspace);
4551 assert_eq!(admitted.mode, staged.mode);
4552 assert_eq!(admitted.allow_shell, staged.allow_shell);
4553 assert_eq!(admitted.trust_mode, staged.trust_mode);
4554 assert!(
4555 !staged_path.exists(),
4556 "unaccepted stage recovered through TaskManager"
4557 );
4558 let completed =
4559 wait_for_terminal_state(&manager, &admitted.id, Duration::from_secs(5)).await?;
4560 assert_eq!(completed.status, TaskStatus::Completed);
4561 assert!(
4562 completed
4563 .result_summary
4564 .as_deref()
4565 .unwrap_or_default()
4566 .contains("admission fixture completed")
4567 );
4568 assert_eq!(executions.load(Ordering::SeqCst), 1);
4569 let replay = manager
4570 .recover_task_admission(request.clone(), staged.id.clone())
4571 .await?;
4572 assert_eq!(replay.status, TaskStatus::Completed);
4573 assert_eq!(replay.id, staged.id);
4574 let canonical = tasks_dir.join(format!("{}.json", staged.id));
4575 let before = fs::read(&canonical)?;
4576 let mut mismatched = request;
4577 mismatched.prompt = "a different operation".into();
4578 let error = manager
4579 .recover_task_admission(mismatched, staged.id)
4580 .await
4581 .expect_err("mismatched replay must be rejected");
4582 assert!(error.to_string().contains("does not match"));
4583 assert_eq!(
4584 fs::read(canonical)?,
4585 before,
4586 "replay cannot rewrite accepted work"
4587 );
4588 assert_eq!(executions.load(Ordering::SeqCst), 1);
4589 manager.shutdown();
4590 }
4591 Ok(())
4592 }
4593
4594 #[tokio::test]
4595 async fn accepted_task_interrupted_by_restart_is_reconciled_without_execution() -> Result<()> {
4596 let root = tempfile::tempdir()?;
4597 let tasks_dir = root.path().join("tasks");
4598 fs::create_dir_all(&tasks_dir)?;
4599 let mut accepted = sample_task_record();
4600 accepted.execution_generation = Some(Uuid::new_v4().simple().to_string());
4601 let lease_path = execution_lease_path(
4602 root.path(),
4603 accepted.execution_scope.as_deref().unwrap(),
4604 accepted.execution_generation.as_deref().unwrap(),
4605 )?;
4606 drop(
4607 RuntimeProcessOwnerLock::try_acquire_file(&lease_path, true)?
4608 .context("fixture generation")?,
4609 );
4610 let request = NewTaskRequest::from_task(&accepted);
4611 write_json_atomic(&tasks_dir.join(format!("{}.json", accepted.id)), &accepted)?;
4612 write_json_atomic(
4613 &root.path().join("queue.json"),
4614 &QueueFile {
4615 queue: vec![accepted.id.clone()],
4616 },
4617 )?;
4618 let executions = Arc::new(AtomicUsize::new(0));
4619 let manager = TaskManager::start_with_executor(
4620 test_config(root.path().to_path_buf()),
4621 Arc::new(AdmissionCountingExecutor(executions.clone())),
4622 )
4623 .await?;
4624 let recovered = manager
4625 .recover_task_admission(request, accepted.id.clone())
4626 .await?;
4627 assert_eq!(recovered.id, accepted.id);
4628 assert_eq!(recovered.status, TaskStatus::Failed);
4629 assert!(
4630 recovered
4631 .error
4632 .as_deref()
4633 .unwrap_or_default()
4634 .contains("Interrupted by process restart")
4635 );
4636 assert_eq!(manager.list_tasks(None).await?.len(), 1);
4637 assert_eq!(
4638 executions.load(Ordering::SeqCst),
4639 0,
4640 "accepted work cannot be replayed after restart"
4641 );
4642 manager.shutdown();
4643 Ok(())
4644 }
4645
4646 #[tokio::test]
4647 async fn preallocated_task_ids_are_validated_and_collision_safe() -> Result<()> {
4648 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
4649 let manager =
4650 TaskManager::start_with_executor(test_config(root), Arc::new(MockExecutor)).await?;
4651 let request = NewTaskRequest::from_prompt("preallocated owner identity");
4652
4653 let invalid = manager
4654 .add_task_with_id(request.clone(), "task_short".to_string())
4655 .await
4656 .expect_err("invalid preallocated id");
4657 assert!(invalid.to_string().contains("task_<16hex>"), "{invalid:#}");
4658
4659 let id = "task_0123456789abcdef".to_string();
4660 let created = manager
4661 .add_task_with_id(request.clone(), id.clone())
4662 .await?;
4663 assert_eq!(created.id, id);
4664 assert_eq!(
4665 created.schema_version, CURRENT_TASK_SCHEMA_VERSION,
4666 "execution provenance requires readers that preserve the binding"
4667 );
4668 assert_eq!(created.lifecycle_seq, 1);
4669 let collision = manager
4670 .add_task_with_id(request, id)
4671 .await
4672 .expect_err("task id collision");
4673 assert!(
4674 collision.to_string().contains("already exists"),
4675 "{collision:#}"
4676 );
4677 assert_eq!(manager.list_tasks(None).await?.len(), 1);
4678 Ok(())
4679 }
4680
4681 #[tokio::test]
4682 async fn failed_queue_write_leaves_no_replayable_task_record() -> Result<()> {
4683 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
4684 let manager =
4685 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
4686 .await?;
4687 std::fs::remove_file(root.join("queue.json"))?;
4688 std::fs::create_dir(root.join("queue.json"))?;
4689
4690 let id = "task_fedcba9876543210".to_string();
4691 let error = manager
4692 .add_task_with_id(
4693 NewTaskRequest::from_prompt("must not resurrect"),
4694 id.clone(),
4695 )
4696 .await
4697 .expect_err("queue path directory must reject the atomic queue write");
4698 assert!(error.to_string().contains("queue.json"), "{error:#}");
4699 assert!(
4700 manager.list_tasks(None).await.is_err(),
4701 "unavailable storage cannot be reported as empty"
4702 );
4703 assert!(!root.join("tasks").join(format!("{id}.json")).exists());
4704 assert!(
4705 !root
4706 .join("tasks")
4707 .join(format!(".{id}.json.pending"))
4708 .exists(),
4709 "a failed queue write may leave no replayable or staged task record"
4710 );
4711 Ok(())
4712 }
4713
4714 #[tokio::test]
4715 async fn list_tasks_scopes_results_to_workspace_before_limit() -> Result<()> {
4716 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
4717 let manager =
4718 TaskManager::start_with_executor(test_config(root), Arc::new(MockExecutor)).await?;
4719
4720 manager
4721 .add_task(NewTaskRequest {
4722 prompt: "task in workspace a".to_string(),
4723 workspace: Some(PathBuf::from("/tmp/workspace-a")),
4724 ..NewTaskRequest::from_prompt("task in workspace a")
4725 })
4726 .await?;
4727 manager
4728 .add_task(NewTaskRequest {
4729 prompt: "task in workspace b".to_string(),
4730 workspace: Some(PathBuf::from("/tmp/workspace-b")),
4731 ..NewTaskRequest::from_prompt("task in workspace b")
4732 })
4733 .await?;
4734
4735 let scoped = manager
4736 .list_tasks_scoped(Some(1), Some(Path::new("/tmp/workspace-a")))
4737 .await?;
4738 assert_eq!(scoped.len(), 1);
4739 assert_eq!(scoped[0].workspace, PathBuf::from("/tmp/workspace-a"));
4740 Ok(())
4741 }
4742
4743 #[tokio::test]
4744 async fn task_controls_are_session_owned_and_legacy_records_fail_closed() -> Result<()> {
4745 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
4746 let manager =
4747 TaskManager::start_with_executor(test_config(root), Arc::new(MockExecutor)).await?;
4748
4749 let mut session_a = sample_task_record();
4750 session_a.id = "task_dead000000000001".to_string();
4751 session_a.owner_session_id = Some("session-a".to_string());
4752 session_a.status = TaskStatus::Completed;
4753 session_a.created_at = Utc::now() - chrono::Duration::seconds(3);
4754
4755 let mut session_b = sample_task_record();
4756 session_b.id = "task_dead000000000002".to_string();
4757 session_b.owner_session_id = Some("session-b".to_string());
4758 session_b.status = TaskStatus::Completed;
4759 session_b.created_at = Utc::now() - chrono::Duration::seconds(2);
4760
4761 let mut session_b_newest = sample_task_record();
4762 session_b_newest.id = "task_beef000000000002".to_string();
4763 session_b_newest.owner_session_id = Some("session-b".to_string());
4764 session_b_newest.status = TaskStatus::Completed;
4765 session_b_newest.created_at = Utc::now();
4766
4767 let mut legacy = sample_task_record();
4768 legacy.id = "task_dead000000000003".to_string();
4769 legacy.owner_session_id = None;
4770 legacy.execution_scope = None;
4771 legacy.status = TaskStatus::Completed;
4772 legacy.created_at = Utc::now() - chrono::Duration::seconds(1);
4773
4774 {
4775 let mut state = manager.state.lock().await;
4776 for record in [
4777 session_a.clone(),
4778 session_b.clone(),
4779 session_b_newest.clone(),
4780 legacy.clone(),
4781 ] {
4782 manager.persist_task_locked(&record)?;
4783 state.tasks.insert(record.id.clone(), record);
4784 }
4785 }
4786
4787 let session_b_list = manager
4788 .list_tasks_for_owner(Some(1), None, "session-b")
4789 .await?;
4790 assert_eq!(session_b_list.len(), 1);
4791 assert_eq!(session_b_list[0].id, session_b_newest.id);
4792
4793 let session_b_prefix = manager.get_task_for_owner("task_dead", "session-b").await?;
4794 assert_eq!(session_b_prefix.id, session_b.id);
4795 assert!(
4796 manager
4797 .get_task_for_owner(&session_a.id, "session-b")
4798 .await
4799 .unwrap_err()
4800 .to_string()
4801 .contains("Task not found")
4802 );
4803 assert!(
4804 manager
4805 .get_task_for_owner(&legacy.id, "session-b")
4806 .await
4807 .unwrap_err()
4808 .to_string()
4809 .contains("Task not found")
4810 );
4811 assert!(
4812 manager
4813 .get_task_for_active_runtime(&legacy.id)
4814 .await
4815 .unwrap_err()
4816 .to_string()
4817 .contains("Task not found"),
4818 "legacy ownerless active tasks must fail closed"
4819 );
4820 assert_eq!(
4821 manager.get_task_for_active_runtime(&session_a.id).await?.id,
4822 session_a.id
4823 );
4824
4825 assert!(
4826 manager
4827 .cancel_task_for_owner(&session_a.id, "session-b")
4828 .await
4829 .unwrap_err()
4830 .to_string()
4831 .contains("Task not found")
4832 );
4833 assert_eq!(
4834 manager.get_task(&session_a.id).await?.status,
4835 TaskStatus::Completed
4836 );
4837 assert!(
4838 manager
4839 .cancel_task_for_owner(&legacy.id, "session-b")
4840 .await
4841 .unwrap_err()
4842 .to_string()
4843 .contains("Task not found")
4844 );
4845
4846 let own = manager
4847 .cancel_task_for_owner(&session_b.id, "session-b")
4848 .await?;
4849 assert_eq!(own.disposition, TaskCancelDisposition::AlreadyFinished);
4850 let active_own = manager
4851 .cancel_task_for_active_runtime(&session_a.id)
4852 .await?;
4853 assert_eq!(
4854 active_own.disposition,
4855 TaskCancelDisposition::AlreadyFinished
4856 );
4857 assert_eq!(
4858 manager
4859 .get_task_for_owner(&session_a.id, "session-a")
4860 .await?
4861 .id,
4862 session_a.id,
4863 "switching A to B and back must restore A's controls"
4864 );
4865 Ok(())
4866 }
4867
4868 #[tokio::test]
4869 async fn interactive_task_controls_include_only_owned_or_same_scope_records() -> Result<()> {
4870 let root = tempfile::tempdir()?;
4871 let manager = TaskManager::start_with_executor(
4872 test_config(root.path().to_path_buf()),
4873 Arc::new(MockExecutor),
4874 )
4875 .await?;
4876 // Exercise persisted controls without a worker racing to execute the
4877 // queued fixture records. The manager retains its verified scope.
4878 manager.shutdown_and_wait().await?;
4879 let mut scheduled = sample_task_record();
4880 scheduled.id = "task_dead000000000001".to_string();
4881 scheduled.owner_session_id = None;
4882 scheduled.execution_scope = Some(manager.execution_scope().to_string());
4883 scheduled.status = TaskStatus::Queued;
4884
4885 let mut owned = scheduled.clone();
4886 owned.id = "task_beef000000000001".to_string();
4887 owned.owner_session_id = Some("session-a".to_string());
4888 owned.status = TaskStatus::Completed;
4889
4890 let mut foreign_scope = scheduled.clone();
4891 foreign_scope.id = "task_dead000000000002".to_string();
4892 foreign_scope.execution_scope = Some(test_execution_scope("other"));
4893 let mut other_session = scheduled.clone();
4894 other_session.id = "task_dead000000000003".to_string();
4895 other_session.owner_session_id = Some("session-b".to_string());
4896 let mut legacy = scheduled.clone();
4897 legacy.id = "task_dead000000000004".to_string();
4898 legacy.execution_scope = None;
4899 let hidden = [foreign_scope, other_session, legacy];
4900 {
4901 let mut state = manager.state.lock().await;
4902 for record in std::iter::once(&scheduled)
4903 .chain(std::iter::once(&owned))
4904 .chain(hidden.iter())
4905 {
4906 manager.persist_task_locked(record)?;
4907 state.tasks.insert(record.id.clone(), record.clone());
4908 }
4909 }
4910
4911 for record in [&scheduled, &owned] {
4912 assert_eq!(
4913 manager
4914 .get_task_for_interactive_session(&record.id, "session-a")
4915 .await?
4916 .id,
4917 record.id
4918 );
4919 }
4920 // Hidden records sharing this prefix must not make it ambiguous.
4921 assert_eq!(
4922 manager
4923 .get_task_for_interactive_session("task_dead", "session-a")
4924 .await?
4925 .id,
4926 scheduled.id
4927 );
4928 let ambiguous = manager
4929 .get_task_for_interactive_session("task_", "session-a")
4930 .await
4931 .unwrap_err();
4932 assert!(ambiguous.to_string().contains("matches 2 tasks"));
4933 for record in &hidden {
4934 assert!(
4935 manager
4936 .get_task_for_interactive_session(&record.id, "session-a")
4937 .await
4938 .unwrap_err()
4939 .to_string()
4940 .contains("Task not found")
4941 );
4942 assert!(
4943 manager
4944 .cancel_task_for_interactive_session(&record.id, "session-a")
4945 .await
4946 .unwrap_err()
4947 .to_string()
4948 .contains("Task not found")
4949 );
4950 assert_eq!(
4951 manager.get_task(&record.id).await?.status,
4952 TaskStatus::Queued
4953 );
4954 }
4955 // Model and child-session APIs do not inherit the human-only access.
4956 assert!(
4957 manager
4958 .get_task_for_owner(&scheduled.id, "session-a")
4959 .await
4960 .is_err()
4961 );
4962 assert!(
4963 manager
4964 .cancel_task_for_owner(&scheduled.id, "session-a")
4965 .await
4966 .is_err()
4967 );
4968 let canceled = manager
4969 .cancel_task_for_interactive_session("task_dead", "session-a")
4970 .await?;
4971 assert_eq!(canceled.task.id, scheduled.id);
4972 assert_eq!(canceled.task.status, TaskStatus::Canceled);
4973 assert_eq!(
4974 manager.get_task(&scheduled.id).await?.status,
4975 TaskStatus::Canceled
4976 );
4977 manager.shutdown_and_wait().await?;
4978 Ok(())
4979 }
4980
4981 #[test]
4982 fn interactive_task_controls_reject_an_empty_manager_scope() {
4983 let mut record = sample_task_record();
4984 record.owner_session_id = None;
4985 record.execution_scope = Some(String::new());
4986 let id = record.id.clone();
4987 let tasks = HashMap::from([(id.clone(), record)]);
4988 assert!(resolve_task_id_visible_to_operator(&tasks, &id, "session-a", "").is_err());
4989 }
4990
4991 #[tokio::test]
4992 async fn boot_does_not_rewrite_non_recovered_task_files() -> Result<()> {
4993 // #3757 boot-persist narrowing: TaskManager::start must persist only
4994 // the reconciled queue and the running->failed recoveries — a
4995 // completed task's file must be byte-identical across a restart.
4996 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
4997 let manager =
4998 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
4999 .await?;
5000 let task = manager
5001 .add_task(NewTaskRequest::from_prompt("finish then persist"))
5002 .await?;
5003 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
5004 assert_eq!(finished.status, TaskStatus::Completed);
5005 manager.shutdown_and_wait().await?;
5006 drop(manager);
5007
5008 let task_file = root.join("tasks").join(format!("{}.json", task.id));
5009 let before = fs::read(&task_file)?;
5010
5011 let recovered =
5012 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
5013 .await?;
5014 // Give start() a beat to run its (narrowed) boot persist.
5015 sleep(Duration::from_millis(50)).await;
5016 drop(recovered);
5017
5018 let after = fs::read(&task_file)?;
5019 assert_eq!(
5020 before, after,
5021 "a completed task file must not be rewritten on boot"
5022 );
5023 Ok(())
5024 }
5025
5026 #[test]
5027 fn legacy_running_tasks_are_preserved_without_assuming_owner_death() -> Result<()> {
5028 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5029 let tasks_dir = root.join("tasks");
5030 fs::create_dir_all(&tasks_dir)?;
5031 let queue_path = root.join("queue.json");
5032 let task_id = "task_stale_running".to_string();
5033 let started_at = Utc::now() - chrono::Duration::seconds(30);
5034 let task = TaskRecord {
5035 schema_version: CURRENT_TASK_SCHEMA_VERSION,
5036 id: task_id.clone(),
5037 prompt: "long-running shell work".to_string(),
5038 name: None,
5039 model: "deepseek-v4-flash".to_string(),
5040 model_provider: None,
5041 model_provider_id: None,
5042 workspace: PathBuf::from("."),
5043 mode: "agent".to_string(),
5044 allow_shell: true,
5045 trust_mode: false,
5046 auto_approve: false,
5047 permission_posture: None,
5048 status: TaskStatus::Running,
5049 created_at: started_at,
5050 started_at: Some(started_at),
5051 ended_at: None,
5052 duration_ms: None,
5053 result_summary: None,
5054 result_detail_path: None,
5055 error: None,
5056 terminal_reason: None,
5057 thread_id: Some("thr_stale".to_string()),
5058 turn_id: Some("turn_stale".to_string()),
5059 owner_session_id: Some("session-old".to_string()),
5060 execution_scope: None,
5061 execution_generation: None,
5062 cancel_requested_seq: 0,
5063 runtime_event_count: 0,
5064 lifecycle_seq: 2,
5065 checklist: TaskChecklistState::default(),
5066 gates: Vec::new(),
5067 attempts: Vec::new(),
5068 artifacts: Vec::new(),
5069 github_events: Vec::new(),
5070 tool_calls: vec![TaskToolCallSummary {
5071 id: "tool_shell".to_string(),
5072 name: "task_shell_start".to_string(),
5073 status: TaskToolStatus::Running,
5074 started_at,
5075 ended_at: None,
5076 duration_ms: None,
5077 input_summary: Some("shell: sleep 999".to_string()),
5078 output_summary: None,
5079 detail_path: None,
5080 patch_ref: None,
5081 }],
5082 timeline: vec![TaskTimelineEntry {
5083 timestamp: started_at,
5084 kind: "running".to_string(),
5085 summary: "Task started".to_string(),
5086 detail_path: None,
5087 }],
5088 };
5089 fs::write(
5090 tasks_dir.join(format!("{task_id}.json")),
5091 serde_json::to_string_pretty(&task)?,
5092 )?;
5093 fs::write(
5094 &queue_path,
5095 serde_json::to_string_pretty(&QueueFile {
5096 queue: vec![task_id.clone()],
5097 })?,
5098 )?;
5099
5100 let loaded = load_state(&tasks_dir, &queue_path)?;
5101 let queue = loaded.queue;
5102 let recovered = loaded.tasks.get(&task_id).expect("task loaded");
5103
5104 assert!(queue.is_empty(), "stale running task must not be requeued");
5105 assert_eq!(recovered.status, TaskStatus::Running);
5106 assert!(recovered.ended_at.is_none());
5107 assert!(recovered.error.is_none());
5108 assert_eq!(recovered.tool_calls[0].status, TaskToolStatus::Running);
5109 assert!(!TaskSummary::from(recovered).execution_binding_known);
5110 assert_eq!(
5111 fs::read(tasks_dir.join(format!("{task_id}.json")))?,
5112 serde_json::to_string_pretty(&task)?.as_bytes()
5113 );
5114 Ok(())
5115 }
5116
5117 #[tokio::test]
5118 async fn default_workspace_updates_for_future_tasks() -> Result<()> {
5119 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5120 let new_workspace =
5121 std::env::temp_dir().join(format!("deepseek-workspace-{}", Uuid::new_v4()));
5122 let manager =
5123 TaskManager::start_with_executor(test_config(root), Arc::new(MockExecutor)).await?;
5124
5125 manager.set_default_workspace(new_workspace.clone()).await;
5126 let task = manager
5127 .add_task(NewTaskRequest::from_prompt("test workspace default"))
5128 .await?;
5129
5130 assert_eq!(manager.default_workspace().await, new_workspace);
5131 assert_eq!(task.workspace, new_workspace);
5132 Ok(())
5133 }
5134
5135 #[tokio::test]
5136 async fn record_tool_metadata_updates_explicit_task() -> Result<()> {
5137 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5138 let manager =
5139 TaskManager::start_with_executor(test_config(root), Arc::new(MockExecutor)).await?;
5140
5141 let task = manager
5142 .add_task(NewTaskRequest::from_prompt("test metadata"))
5143 .await?;
5144 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
5145 let updated = manager
5146 .record_tool_metadata(
5147 &finished.id,
5148 &serde_json::json!({
5149 "task_updates": {
5150 "gate": {
5151 "id": "gate_test",
5152 "gate": "test",
5153 "command": "cargo test -p codewhale-tui --lib",
5154 "cwd": ".",
5155 "exit_code": 0,
5156 "status": "passed",
5157 "classification": "passed",
5158 "duration_ms": 1,
5159 "summary": "ok",
5160 "log_path": null,
5161 "recorded_at": Utc::now()
5162 }
5163 }
5164 }),
5165 )
5166 .await?;
5167
5168 assert_eq!(updated.gates.len(), 1);
5169 assert_eq!(updated.gates[0].classification, "passed");
5170 Ok(())
5171 }
5172
5173 #[tokio::test]
5174 async fn write_task_artifact_rejects_traversal_task_id() -> Result<()> {
5175 let temp = tempfile::tempdir()?;
5176 let root = temp.path().join("tasks-root");
5177 let escaped = temp.path().join("escape");
5178 let manager =
5179 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
5180 .await?;
5181
5182 let err = manager
5183 .write_task_artifact("../escape", "result", "artifact body")
5184 .expect_err("traversal task ids must be rejected");
5185
5186 assert!(err.to_string().contains("single path component"));
5187 assert!(!escaped.exists(), "artifact write escaped the task root");
5188 Ok(())
5189 }
5190
5191 #[tokio::test]
5192 async fn cancel_running_task_marks_canceled() -> Result<()> {
5193 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5194 let manager = TaskManager::start_with_executor(
5195 test_config(root),
5196 Arc::new(CooperativeIdleCancelExecutor),
5197 )
5198 .await?;
5199
5200 let task = manager
5201 .add_task(NewTaskRequest::from_prompt("test cancellation"))
5202 .await?;
5203
5204 wait_for_running(&manager, &task.id, Duration::from_secs(5)).await?;
5205 let cancellation = manager.cancel_task(&task.id).await?;
5206 assert_eq!(cancellation.disposition, TaskCancelDisposition::Requested);
5207 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
5208 assert_eq!(finished.status, TaskStatus::Canceled);
5209 Ok(())
5210 }
5211
5212 #[tokio::test]
5213 async fn cancel_finished_task_returns_atomic_already_finished_outcome() -> Result<()> {
5214 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5215 let manager =
5216 TaskManager::start_with_executor(test_config(root), Arc::new(MockExecutor)).await?;
5217 let task = manager
5218 .add_task(NewTaskRequest::from_prompt("finish before cancellation"))
5219 .await?;
5220 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
5221 assert_eq!(finished.status, TaskStatus::Completed);
5222
5223 let cancellation = manager.cancel_task(&task.id).await?;
5224
5225 assert_eq!(
5226 cancellation.disposition,
5227 TaskCancelDisposition::AlreadyFinished
5228 );
5229 assert_eq!(cancellation.task.status, TaskStatus::Completed);
5230 Ok(())
5231 }
5232
5233 // GHSA-72w5-pf8h-xfp4 — regression: omitted optional fields must not
5234 // silently elevate the spawned task's privileges.
5235 #[tokio::test]
5236 async fn add_task_without_optional_fields_does_not_grant_shell_or_auto_approve() -> Result<()> {
5237 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5238 let manager =
5239 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
5240 .await?;
5241
5242 let req = NewTaskRequest {
5243 prompt: "fix TODOs and write a README".to_string(),
5244 name: None,
5245 model: None,
5246 model_provider: None,
5247 model_provider_id: None,
5248 workspace: None,
5249 mode: None,
5250 allow_shell: None,
5251 trust_mode: None,
5252 auto_approve: None,
5253 permission_posture: None,
5254 owner_session_id: None,
5255 };
5256 let task = manager.add_task(req).await?;
5257
5258 assert!(
5259 !task.allow_shell,
5260 "model-omitted allow_shell must default to false (no silent shell grant)"
5261 );
5262 assert!(
5263 !task.auto_approve,
5264 "model-omitted auto_approve must default to false (no silent auto-approval)"
5265 );
5266 assert!(
5267 !task.trust_mode,
5268 "model-omitted trust_mode must default to false"
5269 );
5270 Ok(())
5271 }
5272
5273 /// A task's own thread starts on the posture the request pinned, and the
5274 /// posture is what the thread's policy is derived from — the legacy
5275 /// `auto_approve` bit is only read when no posture is given.
5276 #[tokio::test]
5277 async fn add_task_pins_the_posture_its_thread_starts_on() -> Result<()> {
5278 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5279 let manager =
5280 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
5281 .await?;
5282
5283 let task = manager
5284 .add_task(NewTaskRequest {
5285 permission_posture: Some("auto_review".to_string()),
5286 ..NewTaskRequest::from_prompt("pin the posture")
5287 })
5288 .await?;
5289
5290 assert_eq!(task.permission_posture.as_deref(), Some("auto_review"));
5291 let request = ExecutionTask::from(&task).thread_request();
5292 assert_eq!(request.permission_posture.as_deref(), Some("auto_review"));
5293 // `from_prompt` asks for auto-approval; the pinned posture outranks it,
5294 // so the thread must not silently run wider than what was requested.
5295 assert_eq!(request.auto_approve, Some(true));
5296 Ok(())
5297 }
5298
5299 /// The worker's thread projection refuses these postures, so admission
5300 /// refuses them too: the request came through the Runtime API, and that
5301 /// is where the refusal belongs, not in a worker after the task was
5302 /// durably queued.
5303 #[tokio::test]
5304 async fn add_task_refuses_a_posture_the_thread_would_reject() -> Result<()> {
5305 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5306 let manager =
5307 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
5308 .await?;
5309
5310 for posture in ["sideways", "never"] {
5311 let error = manager
5312 .add_task(NewTaskRequest {
5313 permission_posture: Some(posture.to_string()),
5314 ..NewTaskRequest::from_prompt("refuse me")
5315 })
5316 .await
5317 .expect_err("a posture the thread cannot honour is refused at admission");
5318 assert!(error.to_string().contains("permission posture"), "{error}");
5319 }
5320 assert!(manager.list_tasks(None).await?.is_empty());
5321 Ok(())
5322 }
5323
5324 /// The Runtime's `POST /v1/tasks` body may omit the posture entirely: it is
5325 /// optional on the wire, and absent means "derive it from the legacy bits",
5326 /// which is what every client that predates the field sends.
5327 #[test]
5328 fn new_task_request_accepts_a_body_without_a_posture() {
5329 let request: NewTaskRequest =
5330 serde_json::from_str(r#"{"prompt":"ship it","mode":"agent"}"#).expect("wire body");
5331 assert!(request.permission_posture.is_none());
5332 assert_eq!(request.mode.as_deref(), Some("agent"));
5333 }
5334
5335 #[tokio::test]
5336 async fn rejects_newer_task_schema_on_recovery() -> Result<()> {
5337 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5338 let manager =
5339 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(MockExecutor))
5340 .await?;
5341
5342 let task = manager
5343 .add_task(NewTaskRequest::from_prompt("test schema gate"))
5344 .await?;
5345 let _ = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
5346 manager.shutdown_and_wait().await?;
5347 drop(manager);
5348
5349 let task_path = root.join("tasks").join(format!("{}.json", task.id));
5350 let mut value: serde_json::Value = serde_json::from_str(&fs::read_to_string(&task_path)?)?;
5351 value["schema_version"] = serde_json::json!(999);
5352 fs::write(&task_path, serde_json::to_string_pretty(&value)?)?;
5353
5354 match TaskManager::start_with_executor(test_config(root), Arc::new(MockExecutor)).await {
5355 Ok(_) => panic!("manager should reject newer task schema"),
5356 Err(err) => assert!(err.to_string().contains("newer than supported")),
5357 }
5358 Ok(())
5359 }
5360
5361 #[test]
5362 fn default_tasks_dir_falls_back_to_legacy_deepseek_tasks() {
5363 let temp_home = tempfile::tempdir().unwrap();
5364 let home = temp_home.path();
5365 let legacy_tasks = home.join(".deepseek").join("tasks");
5366 std::fs::create_dir_all(&legacy_tasks).unwrap();
5367
5368 assert_eq!(default_tasks_dir_for_home(home), legacy_tasks);
5369 }
5370
5371 #[test]
5372 fn default_tasks_dir_prefers_existing_codewhale_tasks() {
5373 let temp_home = tempfile::tempdir().unwrap();
5374 let home = temp_home.path();
5375 let primary_tasks = home.join(".codewhale").join("tasks");
5376 let legacy_tasks = home.join(".deepseek").join("tasks");
5377 std::fs::create_dir_all(&primary_tasks).unwrap();
5378 std::fs::create_dir_all(&legacy_tasks).unwrap();
5379
5380 assert_eq!(default_tasks_dir_for_home(home), primary_tasks);
5381 }
5382
5383 #[test]
5384 fn default_tasks_dir_falls_back_to_legacy_when_primary_is_file() {
5385 let temp_home = tempfile::tempdir().unwrap();
5386 let home = temp_home.path();
5387 let primary_tasks = home.join(".codewhale").join("tasks");
5388 let legacy_tasks = home.join(".deepseek").join("tasks");
5389 std::fs::create_dir_all(primary_tasks.parent().unwrap()).unwrap();
5390 std::fs::write(&primary_tasks, "not a directory").unwrap();
5391 std::fs::create_dir_all(&legacy_tasks).unwrap();
5392
5393 assert_eq!(default_tasks_dir_for_home(home), legacy_tasks);
5394 }
5395
5396 #[test]
5397 fn default_tasks_dir_ignores_legacy_file_for_new_installs() {
5398 let temp_home = tempfile::tempdir().unwrap();
5399 let home = temp_home.path();
5400 let primary_tasks = home.join(".codewhale").join("tasks");
5401 let legacy_tasks = home.join(".deepseek").join("tasks");
5402 std::fs::create_dir_all(legacy_tasks.parent().unwrap()).unwrap();
5403 std::fs::write(&legacy_tasks, "not a directory").unwrap();
5404
5405 assert_eq!(default_tasks_dir_for_home(home), primary_tasks);
5406 }
5407
5408 #[test]
5409 fn default_tasks_dir_uses_codewhale_tasks_for_new_installs() {
5410 let temp_home = tempfile::tempdir().unwrap();
5411 let home = temp_home.path();
5412
5413 assert_eq!(
5414 default_tasks_dir_for_home(home),
5415 home.join(".codewhale").join("tasks")
5416 );
5417 }
5418
5419 #[test]
5420 fn task_and_runtime_roots_honor_explicit_codewhale_home() {
5421 let _lock = lock_test_env();
5422 let temp_root = tempfile::tempdir().unwrap();
5423 let ambient_home = temp_root.path().join("ambient-home");
5424 let explicit_home = temp_root.path().join("explicit-home");
5425 std::fs::create_dir_all(ambient_home.join(".deepseek").join("tasks")).unwrap();
5426 let _home = EnvVarGuard::set("HOME", &ambient_home);
5427 let _userprofile = EnvVarGuard::set("USERPROFILE", &ambient_home);
5428 let _codewhale_home = EnvVarGuard::set("CODEWHALE_HOME", &explicit_home);
5429 let _tasks_override = EnvVarGuard::remove("CODEWHALE_TASKS_DIR");
5430 let _legacy_tasks_override = EnvVarGuard::remove("DEEPSEEK_TASKS_DIR");
5431 let _runtime_override = EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
5432 let _legacy_runtime_override = EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
5433
5434 let task_root = default_tasks_dir();
5435 let task_manager =
5436 TaskManagerConfig::from_runtime(&Config::default(), PathBuf::from("."), None, None);
5437 let runtime = RuntimeThreadManagerConfig::from_task_data_dir(task_manager.data_dir.clone());
5438
5439 assert_eq!(task_root, explicit_home.join("tasks"));
5440 assert_eq!(task_manager.data_dir, task_root);
5441 assert_eq!(runtime.task_data_dir, task_root);
5442 assert_eq!(
5443 runtime.data_dir,
5444 explicit_home.join("tasks").join("runtime")
5445 );
5446 }
5447
5448 #[test]
5449 fn whitespace_codewhale_home_keeps_ambient_legacy_task_and_runtime_fallbacks() {
5450 let _lock = lock_test_env();
5451 let temp_root = tempfile::tempdir().unwrap();
5452 let ambient_home = temp_root.path().join("ambient-home");
5453 let legacy_tasks = ambient_home.join(".deepseek").join("tasks");
5454 std::fs::create_dir_all(&legacy_tasks).unwrap();
5455 let _home = EnvVarGuard::set("HOME", &ambient_home);
5456 let _userprofile = EnvVarGuard::set("USERPROFILE", &ambient_home);
5457 let _codewhale_home = EnvVarGuard::set("CODEWHALE_HOME", " \t ");
5458 let _tasks_override = EnvVarGuard::remove("CODEWHALE_TASKS_DIR");
5459 let _legacy_tasks_override = EnvVarGuard::remove("DEEPSEEK_TASKS_DIR");
5460 let _runtime_override = EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
5461 let _legacy_runtime_override = EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
5462
5463 let task_root = default_tasks_dir();
5464 let task_manager =
5465 TaskManagerConfig::from_runtime(&Config::default(), PathBuf::from("."), None, None);
5466 let runtime = RuntimeThreadManagerConfig::from_task_data_dir(task_manager.data_dir.clone());
5467
5468 assert_eq!(task_root, legacy_tasks);
5469 assert_eq!(task_manager.data_dir, task_root);
5470 assert_eq!(runtime.task_data_dir, task_root);
5471 assert_eq!(runtime.data_dir, task_root.join("runtime"));
5472 }
5473
5474 #[cfg(unix)]
5475 #[test]
5476 fn non_unicode_codewhale_home_is_preserved_by_task_and_runtime_roots() {
5477 use std::os::unix::ffi::OsStringExt;
5478
5479 let _lock = lock_test_env();
5480 let temp_root = tempfile::tempdir().unwrap();
5481 let explicit_home = temp_root.path().join(std::ffi::OsString::from_vec(
5482 b"codewhale-\xff-home".to_vec(),
5483 ));
5484 let _codewhale_home = EnvVarGuard::set("CODEWHALE_HOME", &explicit_home);
5485 let _tasks_override = EnvVarGuard::remove("CODEWHALE_TASKS_DIR");
5486 let _legacy_tasks_override = EnvVarGuard::remove("DEEPSEEK_TASKS_DIR");
5487 let _runtime_override = EnvVarGuard::remove("CODEWHALE_RUNTIME_DIR");
5488 let _legacy_runtime_override = EnvVarGuard::remove("DEEPSEEK_RUNTIME_DIR");
5489
5490 let task_root = default_tasks_dir();
5491 let task_manager =
5492 TaskManagerConfig::from_runtime(&Config::default(), PathBuf::from("."), None, None);
5493 let runtime = RuntimeThreadManagerConfig::from_task_data_dir(task_manager.data_dir.clone());
5494
5495 assert_eq!(task_root, explicit_home.join("tasks"));
5496 assert_eq!(task_manager.data_dir, task_root);
5497 assert_eq!(runtime.task_data_dir, task_root);
5498 assert_eq!(
5499 runtime.data_dir,
5500 explicit_home.join("tasks").join("runtime")
5501 );
5502 }
5503
5504 struct DeafHangExecutor;
5505
5506 #[async_trait]
5507 impl TaskExecutor for DeafHangExecutor {
5508 async fn execute(
5509 &self,
5510 _task: ExecutionTask,
5511 _events: mpsc::Sender<TaskExecutionEvent>,
5512 _cancel: CancellationToken,
5513 ) -> TaskExecutionResult {
5514 std::future::pending().await
5515 }
5516 }
5517
5518 struct PartialThenHangExecutor;
5519
5520 #[async_trait]
5521 impl TaskExecutor for PartialThenHangExecutor {
5522 async fn execute(
5523 &self,
5524 _task: ExecutionTask,
5525 events: mpsc::Sender<TaskExecutionEvent>,
5526 _cancel: CancellationToken,
5527 ) -> TaskExecutionResult {
5528 let _ = events
5529 .send(TaskExecutionEvent::MessageDelta {
5530 content: "partial ".to_string(),
5531 })
5532 .await;
5533 let _ = events
5534 .send(TaskExecutionEvent::MessageDelta {
5535 content: "result".to_string(),
5536 })
5537 .await;
5538 std::future::pending().await
5539 }
5540 }
5541
5542 struct PollCountingHangExecutor {
5543 polls: Arc<AtomicUsize>,
5544 }
5545
5546 #[async_trait]
5547 impl TaskExecutor for PollCountingHangExecutor {
5548 async fn execute(
5549 &self,
5550 _task: ExecutionTask,
5551 _events: mpsc::Sender<TaskExecutionEvent>,
5552 _cancel: CancellationToken,
5553 ) -> TaskExecutionResult {
5554 std::future::poll_fn(|_| {
5555 self.polls.fetch_add(1, Ordering::Relaxed);
5556 std::task::Poll::Pending
5557 })
5558 .await
5559 }
5560 }
5561
5562 struct HeartbeatExecutor;
5563
5564 #[async_trait]
5565 impl TaskExecutor for HeartbeatExecutor {
5566 async fn execute(
5567 &self,
5568 _task: ExecutionTask,
5569 events: mpsc::Sender<TaskExecutionEvent>,
5570 _cancel: CancellationToken,
5571 ) -> TaskExecutionResult {
5572 loop {
5573 let _ = events
5574 .send(TaskExecutionEvent::Status {
5575 message: "heartbeat".to_string(),
5576 })
5577 .await;
5578 sleep(Duration::from_millis(10)).await;
5579 }
5580 }
5581 }
5582
5583 struct ProgressHeartbeatExecutor;
5584
5585 #[async_trait]
5586 impl TaskExecutor for ProgressHeartbeatExecutor {
5587 async fn execute(
5588 &self,
5589 _task: ExecutionTask,
5590 events: mpsc::Sender<TaskExecutionEvent>,
5591 _cancel: CancellationToken,
5592 ) -> TaskExecutionResult {
5593 loop {
5594 let _ = events
5595 .send(TaskExecutionEvent::MessageDelta {
5596 content: "working".to_string(),
5597 })
5598 .await;
5599 sleep(Duration::from_millis(10)).await;
5600 }
5601 }
5602 }
5603
5604 struct CooperativeIdleCancelExecutor;
5605
5606 #[async_trait]
5607 impl TaskExecutor for CooperativeIdleCancelExecutor {
5608 async fn execute(
5609 &self,
5610 _task: ExecutionTask,
5611 _events: mpsc::Sender<TaskExecutionEvent>,
5612 cancel: CancellationToken,
5613 ) -> TaskExecutionResult {
5614 cancel.cancelled().await;
5615 TaskExecutionResult::from_reason(TaskTerminalReason::Canceled, None)
5616 }
5617 }
5618
5619 struct CooperativeProgressCancelExecutor;
5620
5621 #[async_trait]
5622 impl TaskExecutor for CooperativeProgressCancelExecutor {
5623 async fn execute(
5624 &self,
5625 _task: ExecutionTask,
5626 events: mpsc::Sender<TaskExecutionEvent>,
5627 cancel: CancellationToken,
5628 ) -> TaskExecutionResult {
5629 loop {
5630 tokio::select! {
5631 _ = cancel.cancelled() => {
5632 return TaskExecutionResult::from_reason(
5633 TaskTerminalReason::Canceled,
5634 None,
5635 );
5636 }
5637 _ = sleep(Duration::from_millis(10)) => {
5638 let _ = events
5639 .send(TaskExecutionEvent::MessageDelta {
5640 content: "working".to_string(),
5641 })
5642 .await;
5643 }
5644 }
5645 }
5646 }
5647 }
5648
5649 struct PromptRouterExecutor;
5650
5651 #[async_trait]
5652 impl TaskExecutor for PromptRouterExecutor {
5653 async fn execute(
5654 &self,
5655 task: ExecutionTask,
5656 events: mpsc::Sender<TaskExecutionEvent>,
5657 _cancel: CancellationToken,
5658 ) -> TaskExecutionResult {
5659 if task.prompt.starts_with("hang ") {
5660 std::future::pending().await
5661 } else {
5662 // The follow-up task must complete without a single await
5663 // point: `run_task` polls the executor future before its
5664 // guard can observe the (test-shortened) idle/wall budgets,
5665 // so an await-free future always finishes first and an
5666 // interrupt can never be recorded against it. The previous
5667 // MockExecutor delegation (`send(...).await` x4 plus a 50 ms
5668 // sleep) left windows where CI scheduler/storage stalls of
5669 // >=150 ms tripped the idle watchdog mid-flight; the executor
5670 // then observed the cancellation and returned `Canceled`,
5671 // which `preserve_timeout_reason` rewrote into the timeout
5672 // reason -> `Failed` (issue #5898). `try_send` keeps the
5673 // released worker's event pipeline exercised without
5674 // suspending this future.
5675 let _ = events.try_send(TaskExecutionEvent::Status {
5676 message: format!("running after forced release {}", task.id),
5677 });
5678 TaskExecutionResult {
5679 status: TaskStatus::Completed,
5680 result_text: Some("done after hang".to_string()),
5681 error: None,
5682 terminal_reason: TaskTerminalReason::Completed,
5683 }
5684 }
5685 }
5686 }
5687
5688 struct FloodExecutor;
5689
5690 #[async_trait]
5691 impl TaskExecutor for FloodExecutor {
5692 async fn execute(
5693 &self,
5694 _task: ExecutionTask,
5695 events: mpsc::Sender<TaskExecutionEvent>,
5696 _cancel: CancellationToken,
5697 ) -> TaskExecutionResult {
5698 for i in 0..400 {
5699 // Mirror the runtime path: each raw event is followed by its
5700 // derived message delta. Alternating the two non-urgent stream
5701 // kinds prevents timeline coalescing without turning this
5702 // storage-bound test into hundreds of synchronous fsyncs.
5703 let _ = events
5704 .send(TaskExecutionEvent::RuntimeEvent {
5705 seq: i,
5706 event: "item.delta".to_string(),
5707 summary: format!("tick {i}"),
5708 })
5709 .await;
5710 let _ = events
5711 .send(TaskExecutionEvent::MessageDelta {
5712 content: format!("chunk {i}"),
5713 })
5714 .await;
5715 }
5716 TaskExecutionResult {
5717 status: TaskStatus::Completed,
5718 result_text: Some("flooded".to_string()),
5719 error: None,
5720 terminal_reason: TaskTerminalReason::Completed,
5721 }
5722 }
5723 }
5724
5725 struct CompleteAfterCancelExecutor;
5726
5727 #[async_trait]
5728 impl TaskExecutor for CompleteAfterCancelExecutor {
5729 async fn execute(
5730 &self,
5731 _task: ExecutionTask,
5732 _events: mpsc::Sender<TaskExecutionEvent>,
5733 cancel: CancellationToken,
5734 ) -> TaskExecutionResult {
5735 cancel.cancelled().await;
5736 TaskExecutionResult {
5737 status: TaskStatus::Completed,
5738 result_text: Some("late complete".to_string()),
5739 error: None,
5740 terminal_reason: TaskTerminalReason::Completed,
5741 }
5742 }
5743 }
5744
5745 fn sample_task_record() -> TaskRecord {
5746 TaskRecord {
5747 schema_version: CURRENT_TASK_SCHEMA_VERSION,
5748 id: "task_0123456789abcdef".to_string(),
5749 prompt: "bound timeline".to_string(),
5750 name: None,
5751 model: "deepseek-v4-flash".to_string(),
5752 model_provider: None,
5753 model_provider_id: None,
5754 workspace: PathBuf::from("."),
5755 mode: "agent".to_string(),
5756 allow_shell: false,
5757 trust_mode: false,
5758 auto_approve: false,
5759 permission_posture: None,
5760 status: TaskStatus::Running,
5761 created_at: Utc::now(),
5762 started_at: Some(Utc::now()),
5763 ended_at: None,
5764 duration_ms: None,
5765 result_summary: None,
5766 result_detail_path: None,
5767 error: None,
5768 terminal_reason: None,
5769 thread_id: None,
5770 turn_id: None,
5771 owner_session_id: None,
5772 execution_scope: Some(test_execution_scope("test")),
5773 execution_generation: None,
5774 cancel_requested_seq: 0,
5775 runtime_event_count: 0,
5776 lifecycle_seq: 2,
5777 checklist: TaskChecklistState::default(),
5778 gates: Vec::new(),
5779 attempts: Vec::new(),
5780 artifacts: Vec::new(),
5781 github_events: Vec::new(),
5782 tool_calls: Vec::new(),
5783 timeline: Vec::new(),
5784 }
5785 }
5786
5787 /// Records persisted before `auto_approve` existed must load as not
5788 /// auto-approved, and the thread request they build must say so.
5789 #[test]
5790 fn task_record_missing_auto_approve_loads_fail_closed() {
5791 let mut record = sample_task_record();
5792 record.auto_approve = true;
5793 let mut value = serde_json::to_value(&record).expect("serialize");
5794 assert_eq!(value["auto_approve"], serde_json::json!(true));
5795 value
5796 .as_object_mut()
5797 .expect("object")
5798 .remove("auto_approve");
5799 let loaded: TaskRecord = serde_json::from_value(value).expect("legacy record loads");
5800 assert!(!loaded.auto_approve);
5801 assert_eq!(
5802 ExecutionTask::from(&loaded).thread_request().auto_approve,
5803 Some(false)
5804 );
5805 assert_eq!(
5806 ExecutionTask::from(&loaded).turn_request().auto_approve,
5807 Some(false)
5808 );
5809 }
5810
5811 #[test]
5812 fn execution_guard_idle_does_not_reset_without_progress() {
5813 let start = Instant::now();
5814 let limits = TaskExecutionLimits::short_for_tests();
5815 let guard = ExecutionGuard::new(limits, start);
5816 match guard.evaluate(start + limits.idle_progress, false, false) {
5817 GuardAction::Interrupt { reason } => {
5818 assert_eq!(reason, TaskTerminalReason::IdleTimeout);
5819 }
5820 other => panic!("expected idle interrupt, got {other:?}"),
5821 }
5822 }
5823
5824 #[test]
5825 fn execution_guard_reports_the_limit_that_expired_first_when_both_elapsed() {
5826 let start = Instant::now();
5827 let limits = TaskExecutionLimits::short_for_tests();
5828 let guard = ExecutionGuard::new(limits, start);
5829 // A starved watchdog that first ticks after both budgets ran out must
5830 // still report the idle limit, which expired first.
5831 match guard.evaluate(
5832 start + limits.wall_time + limits.idle_progress,
5833 false,
5834 false,
5835 ) {
5836 GuardAction::Interrupt { reason } => {
5837 assert_eq!(reason, TaskTerminalReason::IdleTimeout);
5838 }
5839 other => panic!("expected idle interrupt, got {other:?}"),
5840 }
5841
5842 // Late progress pushes the idle deadline past the wall deadline, so
5843 // the same starved tick reports the wall limit instead.
5844 let mut guard = ExecutionGuard::new(limits, start);
5845 guard.note_progress(start + limits.wall_time - Duration::from_millis(1));
5846 match guard.evaluate(
5847 start + limits.wall_time + limits.idle_progress,
5848 false,
5849 false,
5850 ) {
5851 GuardAction::Interrupt { reason } => {
5852 assert_eq!(reason, TaskTerminalReason::WallTimeout);
5853 }
5854 other => panic!("expected wall interrupt, got {other:?}"),
5855 }
5856 }
5857
5858 #[test]
5859 fn execution_guard_progress_refreshes_idle_until_wall_timeout() {
5860 let start = Instant::now();
5861 let limits = TaskExecutionLimits::short_for_tests();
5862 let mut guard = ExecutionGuard::new(limits, start);
5863 let progressed = start + (limits.idle_progress / 2);
5864 guard.note_progress(progressed);
5865 match guard.evaluate(progressed + (limits.idle_progress / 2), false, false) {
5866 GuardAction::Run { .. } => {}
5867 other => panic!("progress should keep idle from firing, got {other:?}"),
5868 }
5869 // Progress keeps arriving, so the idle deadline never expires before
5870 // the wall deadline does.
5871 guard.note_progress(start + limits.wall_time - (limits.idle_progress / 2));
5872 match guard.evaluate(start + limits.wall_time, false, false) {
5873 GuardAction::Interrupt { reason } => {
5874 assert_eq!(reason, TaskTerminalReason::WallTimeout);
5875 }
5876 other => panic!("expected wall interrupt, got {other:?}"),
5877 }
5878 }
5879
5880 #[test]
5881 fn execution_guard_cancel_grace_terminalizes_stuck_work() {
5882 let start = Instant::now();
5883 let limits = TaskExecutionLimits::short_for_tests();
5884 let mut guard = ExecutionGuard::new(limits, start);
5885 match guard.evaluate(start, true, false) {
5886 GuardAction::Interrupt { reason } => {
5887 assert_eq!(reason, TaskTerminalReason::Canceled);
5888 guard.note_interrupt(start, reason);
5889 }
5890 other => panic!("expected cancel interrupt, got {other:?}"),
5891 }
5892 match guard.evaluate(start + limits.cancel_grace, true, false) {
5893 GuardAction::Terminalize { reason } => {
5894 assert_eq!(reason, TaskTerminalReason::CancelTimeout);
5895 }
5896 other => panic!("expected cancel timeout, got {other:?}"),
5897 }
5898 }
5899
5900 #[test]
5901 fn consecutive_message_deltas_coalesce_on_the_timeline() {
5902 let mut task = sample_task_record();
5903 for i in 0..50 {
5904 push_timeline_entry(
5905 &mut task,
5906 TaskTimelineEntry {
5907 timestamp: Utc::now(),
5908 kind: "message".to_string(),
5909 summary: format!("chunk {i}"),
5910 detail_path: None,
5911 },
5912 );
5913 }
5914 assert_eq!(
5915 task.timeline
5916 .iter()
5917 .filter(|entry| entry.kind == "message")
5918 .count(),
5919 1
5920 );
5921 assert_eq!(
5922 task.timeline.last().map(|e| e.summary.as_str()),
5923 Some("chunk 49")
5924 );
5925 }
5926
5927 #[test]
5928 fn timeline_trim_bounds_growth_and_keeps_a_head() {
5929 let mut task = sample_task_record();
5930 for i in 0..400 {
5931 push_timeline_entry(
5932 &mut task,
5933 TaskTimelineEntry {
5934 timestamp: Utc::now(),
5935 kind: "status".to_string(),
5936 summary: format!("tick {i}"),
5937 detail_path: None,
5938 },
5939 );
5940 }
5941 assert!(task.timeline.len() <= TIMELINE_ENTRY_LIMIT);
5942 assert_eq!(task.timeline[0].summary, "tick 0");
5943 assert!(
5944 task.timeline.iter().any(|entry| entry.kind == "omitted"),
5945 "bounded timeline should record omitted history: {:?}",
5946 task.timeline
5947 .iter()
5948 .map(|e| e.kind.as_str())
5949 .collect::<Vec<_>>()
5950 );
5951 }
5952
5953 #[tokio::test]
5954 async fn never_terminalizing_execution_fails_with_idle_timeout() -> Result<()> {
5955 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5956 let manager =
5957 TaskManager::start_with_executor(short_test_config(root), Arc::new(DeafHangExecutor))
5958 .await?;
5959 let task = manager
5960 .add_task(NewTaskRequest::from_prompt("never finish"))
5961 .await?;
5962 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
5963 assert_eq!(finished.status, TaskStatus::Failed);
5964 assert_eq!(finished.terminal_reason.as_deref(), Some("idle_timeout"));
5965 assert!(
5966 finished
5967 .error
5968 .as_deref()
5969 .is_some_and(|err| err.contains("idle")),
5970 "idle timeout must be visible on the receipt: {finished:?}"
5971 );
5972 Ok(())
5973 }
5974
5975 #[tokio::test]
5976 async fn forced_timeout_keeps_all_partial_message_output() -> Result<()> {
5977 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5978 let manager = TaskManager::start_with_executor(
5979 short_test_config(root),
5980 Arc::new(PartialThenHangExecutor),
5981 )
5982 .await?;
5983 let task = manager
5984 .add_task(NewTaskRequest::from_prompt("retain partial result"))
5985 .await?;
5986 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
5987
5988 assert_eq!(finished.status, TaskStatus::Failed);
5989 assert_eq!(finished.terminal_reason.as_deref(), Some("idle_timeout"));
5990 assert_eq!(finished.result_summary.as_deref(), Some("partial result"));
5991 Ok(())
5992 }
5993
5994 #[tokio::test]
5995 async fn cooperative_cancel_after_idle_timeout_keeps_timeout_reason() -> Result<()> {
5996 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
5997 let manager = TaskManager::start_with_executor(
5998 short_test_config(root),
5999 Arc::new(CooperativeIdleCancelExecutor),
6000 )
6001 .await?;
6002 let task = manager
6003 .add_task(NewTaskRequest::from_prompt("cooperative idle timeout"))
6004 .await?;
6005 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6006 assert_eq!(finished.status, TaskStatus::Failed);
6007 assert_eq!(finished.terminal_reason.as_deref(), Some("idle_timeout"));
6008 assert!(
6009 finished
6010 .error
6011 .as_deref()
6012 .is_some_and(|error| error.contains("idle")),
6013 "cooperative cancellation must retain the timeout receipt: {finished:?}"
6014 );
6015 Ok(())
6016 }
6017
6018 #[tokio::test]
6019 async fn heartbeat_status_does_not_refresh_idle_timeout() -> Result<()> {
6020 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
6021 let manager =
6022 TaskManager::start_with_executor(short_test_config(root), Arc::new(HeartbeatExecutor))
6023 .await?;
6024 let task = manager
6025 .add_task(NewTaskRequest::from_prompt("heartbeat only"))
6026 .await?;
6027 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6028 assert_eq!(finished.status, TaskStatus::Failed);
6029 assert_eq!(finished.terminal_reason.as_deref(), Some("idle_timeout"));
6030 Ok(())
6031 }
6032
6033 #[tokio::test]
6034 async fn active_progress_keeps_idle_alive_until_wall_timeout() -> Result<()> {
6035 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
6036 let manager = TaskManager::start_with_executor(
6037 short_test_config(root),
6038 Arc::new(ProgressHeartbeatExecutor),
6039 )
6040 .await?;
6041 let task = manager
6042 .add_task(NewTaskRequest::from_prompt("genuine progress"))
6043 .await?;
6044 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6045 assert_eq!(finished.status, TaskStatus::Failed);
6046 assert_eq!(finished.terminal_reason.as_deref(), Some("wall_timeout"));
6047 Ok(())
6048 }
6049
6050 #[tokio::test]
6051 async fn cooperative_cancel_after_wall_timeout_keeps_timeout_reason() -> Result<()> {
6052 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
6053 let manager = TaskManager::start_with_executor(
6054 wall_timeout_test_config(root),
6055 Arc::new(CooperativeProgressCancelExecutor),
6056 )
6057 .await?;
6058 let task = manager
6059 .add_task(NewTaskRequest::from_prompt("cooperative wall timeout"))
6060 .await?;
6061 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6062 assert_eq!(finished.status, TaskStatus::Failed);
6063 assert_eq!(finished.terminal_reason.as_deref(), Some("wall_timeout"));
6064 assert!(
6065 finished
6066 .error
6067 .as_deref()
6068 .is_some_and(|error| error.contains("wall-time")),
6069 "cooperative cancellation must retain the timeout receipt: {finished:?}"
6070 );
6071 Ok(())
6072 }
6073
6074 #[tokio::test]
6075 async fn shutdown_terminalizes_a_stuck_running_task() -> Result<()> {
6076 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
6077 let manager =
6078 TaskManager::start_with_executor(short_test_config(root), Arc::new(DeafHangExecutor))
6079 .await?;
6080 let task = manager
6081 .add_task(NewTaskRequest::from_prompt("stuck during shutdown"))
6082 .await?;
6083 wait_for_running(&manager, &task.id, Duration::from_secs(5)).await?;
6084 manager.shutdown();
6085 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6086 assert_eq!(finished.status, TaskStatus::Canceled);
6087 assert_eq!(finished.terminal_reason.as_deref(), Some("shutdown"));
6088 Ok(())
6089 }
6090
6091 #[tokio::test]
6092 async fn shutdown_cancel_signal_does_not_spin_during_grace() -> Result<()> {
6093 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
6094 let polls = Arc::new(AtomicUsize::new(0));
6095 let manager = TaskManager::start_with_executor(
6096 short_test_config(root),
6097 Arc::new(PollCountingHangExecutor {
6098 polls: Arc::clone(&polls),
6099 }),
6100 )
6101 .await?;
6102 let task = manager
6103 .add_task(NewTaskRequest::from_prompt("stuck during shutdown"))
6104 .await?;
6105
6106 wait_for_running(&manager, &task.id, Duration::from_secs(5)).await?;
6107
6108 manager.shutdown();
6109 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6110 assert_eq!(finished.terminal_reason.as_deref(), Some("shutdown"));
6111 assert!(
6112 polls.load(Ordering::Relaxed) <= 10,
6113 "already-canceled shutdown signal repeatedly repolled the executor during grace: {} polls",
6114 polls.load(Ordering::Relaxed)
6115 );
6116 Ok(())
6117 }
6118
6119 #[tokio::test]
6120 async fn forced_idle_timeout_releases_the_worker_for_later_tasks() -> Result<()> {
6121 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
6122 let manager = TaskManager::start_with_executor(
6123 short_test_config(root),
6124 Arc::new(PromptRouterExecutor),
6125 )
6126 .await?;
6127 let stuck = manager
6128 .add_task(NewTaskRequest::from_prompt("hang until idle timeout"))
6129 .await?;
6130 let finished =
6131 wait_for_terminal_state(&manager, &stuck.id, Duration::from_secs(10)).await?;
6132 assert_eq!(
6133 finished.terminal_reason.as_deref(),
6134 Some("idle_timeout"),
6135 "stuck task terminal record: {finished:?}"
6136 );
6137
6138 let next = manager
6139 .add_task(NewTaskRequest::from_prompt("run after hang"))
6140 .await?;
6141 let completed =
6142 wait_for_terminal_state(&manager, &next.id, Duration::from_secs(10)).await?;
6143 assert_eq!(
6144 completed.status,
6145 TaskStatus::Completed,
6146 "follow-up task terminal record: {completed:?}"
6147 );
6148 assert_eq!(
6149 completed.terminal_reason.as_deref(),
6150 Some("completed"),
6151 "follow-up task terminal record: {completed:?}"
6152 );
6153 Ok(())
6154 }
6155
6156 #[tokio::test]
6157 async fn cancel_then_completed_result_is_recorded_as_canceled() -> Result<()> {
6158 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
6159 let manager = TaskManager::start_with_executor(
6160 test_config(root),
6161 Arc::new(CompleteAfterCancelExecutor),
6162 )
6163 .await?;
6164 let task = manager
6165 .add_task(NewTaskRequest::from_prompt("race complete after cancel"))
6166 .await?;
6167 wait_for_running(&manager, &task.id, Duration::from_secs(5)).await?;
6168 let cancellation = manager.cancel_task(&task.id).await?;
6169 assert_eq!(cancellation.disposition, TaskCancelDisposition::Requested);
6170 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6171 assert_eq!(finished.status, TaskStatus::Canceled);
6172 assert_eq!(finished.terminal_reason.as_deref(), Some("canceled"));
6173 Ok(())
6174 }
6175
6176 #[tokio::test]
6177 async fn long_stream_timeline_is_bounded() -> Result<()> {
6178 let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4()));
6179 let manager =
6180 TaskManager::start_with_executor(test_config(root.clone()), Arc::new(FloodExecutor))
6181 .await?;
6182 let task = manager
6183 .add_task(NewTaskRequest::from_prompt("flood the timeline"))
6184 .await?;
6185 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6186 assert_eq!(finished.status, TaskStatus::Completed);
6187 assert_eq!(finished.runtime_event_count, 400);
6188 assert!(
6189 finished.timeline.len() <= TIMELINE_ENTRY_LIMIT,
6190 "timeline grew to {}",
6191 finished.timeline.len()
6192 );
6193 assert!(
6194 finished
6195 .timeline
6196 .iter()
6197 .any(|entry| entry.kind == "omitted"),
6198 "long streams should drop older timeline entries: {:?}",
6199 finished
6200 .timeline
6201 .iter()
6202 .map(|e| e.kind.as_str())
6203 .collect::<Vec<_>>()
6204 );
6205
6206 let persisted_path = root.join("tasks").join(format!("{}.json", task.id));
6207 let persisted: TaskRecord = serde_json::from_slice(&fs::read(&persisted_path)?)?;
6208 assert_eq!(persisted.status, TaskStatus::Completed);
6209 assert_eq!(persisted.runtime_event_count, 400);
6210 assert!(persisted.timeline.len() <= TIMELINE_ENTRY_LIMIT);
6211 assert!(
6212 persisted
6213 .timeline
6214 .iter()
6215 .any(|entry| entry.kind == "omitted")
6216 );
6217 Ok(())
6218 }
6219
6220 async fn test_runtime_manager() -> Result<RuntimeThreadManager> {
6221 let root = tempfile::tempdir()?.keep();
6222 RuntimeThreadManager::open(
6223 Config::default(),
6224 PathBuf::from("."),
6225 RuntimeThreadManagerConfig::from_task_data_dir(root),
6226 )
6227 }
6228
6229 async fn drain_task_events(mut rx: mpsc::Receiver<TaskExecutionEvent>) {
6230 while rx.recv().await.is_some() {}
6231 }
6232
6233 struct RuntimeProjectionExecutor(Vec<(&'static str, Value)>);
6234
6235 #[async_trait]
6236 impl TaskExecutor for RuntimeProjectionExecutor {
6237 async fn execute(
6238 &self,
6239 _task: ExecutionTask,
6240 events: mpsc::Sender<TaskExecutionEvent>,
6241 cancel: CancellationToken,
6242 ) -> TaskExecutionResult {
6243 let runtime = test_runtime_manager().await.expect("fixture runtime");
6244 let thread = runtime
6245 .create_thread(CreateThreadRequest::default())
6246 .await
6247 .expect("fixture thread");
6248 for (event, payload) in &self.0 {
6249 runtime
6250 .emit_event_for_test(
6251 &thread.id,
6252 Some("turn_projection"),
6253 event,
6254 payload.clone(),
6255 )
6256 .await
6257 .expect("persist fixture runtime event");
6258 }
6259 drive_engine_turn(
6260 &runtime,
6261 &thread.id,
6262 "turn_projection",
6263 events,
6264 cancel,
6265 TaskExecutionLimits::default(),
6266 )
6267 .await
6268 }
6269 }
6270
6271 async fn project_runtime_task(events: Vec<(&'static str, Value)>) -> Result<TaskRecord> {
6272 let root = tempfile::tempdir()?;
6273 let manager = TaskManager::start_with_executor(
6274 test_config(root.path().to_path_buf()),
6275 Arc::new(RuntimeProjectionExecutor(events)),
6276 )
6277 .await?;
6278 let task = manager
6279 .add_task(NewTaskRequest::from_prompt("runtime projection fixture"))
6280 .await?;
6281 wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6282 manager.shutdown_and_wait().await?;
6283 // Assert the durable task projection, not just an adapter event.
6284 let path = root.path().join("tasks").join(format!("{}.json", task.id));
6285 Ok(serde_json::from_slice(&fs::read(path)?)?)
6286 }
6287
6288 #[tokio::test]
6289 async fn runtime_task_projection_preserves_provider_ids_and_terminal_tool_statuses()
6290 -> Result<()> {
6291 let mut events = Vec::new();
6292 for (item, provider) in [
6293 ("item_a", "call_a"),
6294 ("item_b", "call_b"),
6295 ("item_c", "call_c"),
6296 ] {
6297 events.push((
6298 "item.started",
6299 json!({
6300 "item": { "id": item, "kind": "tool_call" },
6301 "tool": { "id": provider, "name": "read", "input": {} }
6302 }),
6303 ));
6304 }
6305 // Parallel same-name calls finish out of order. Success preserves
6306 // tool_result_for; errors retain tool_use_id; redaction uses tool_call_id.
6307 for (item, provider, identity_key, terminal) in [
6308 ("item_c", "call_c", "tool_use_id", "item.failed"),
6309 ("item_b", "call_b", "tool_call_id", "item.completed"),
6310 ("item_a", "call_a", "tool_result_for", "item.completed"),
6311 ] {
6312 events.push((
6313 terminal,
6314 json!({ "item": {
6315 "id": item, "kind": "tool_call", "summary": "redacted receipt",
6316 "detail": format!("result for {provider}"),
6317 "metadata": { identity_key: provider, "tool_name": "read" }
6318 }}),
6319 ));
6320 }
6321 // Old event shapes with no metadata retain their existing identity.
6322 events.extend([
6323 ("item.started", json!({ "tool": { "id": "legacy", "name": "read", "input": {} } })),
6324 ("item.completed", json!({ "item": { "id": "legacy", "kind": "tool_call", "summary": "read: ok", "detail": "ok" } })),
6325 ("turn.completed", json!({ "turn": { "status": "completed" } })),
6326 ]);
6327 let task = project_runtime_task(events).await?;
6328 assert_eq!(task.status, TaskStatus::Completed);
6329 assert_eq!(task.tool_calls.len(), 4);
6330 for (call, expected_id, expected_status) in [
6331 (&task.tool_calls[0], "call_a", TaskToolStatus::Success),
6332 (&task.tool_calls[1], "call_b", TaskToolStatus::Success),
6333 (&task.tool_calls[2], "call_c", TaskToolStatus::Failed),
6334 (&task.tool_calls[3], "legacy", TaskToolStatus::Success),
6335 ] {
6336 assert_eq!(call.id, expected_id);
6337 assert_eq!(call.status, expected_status);
6338 assert!(call.ended_at.is_some());
6339 assert!(call.duration_ms.is_some());
6340 }
6341 assert_eq!(
6342 task.tool_calls[0].output_summary.as_deref(),
6343 Some("result for call_a")
6344 );
6345 assert_eq!(
6346 task.tool_calls[2].output_summary.as_deref(),
6347 Some("result for call_c")
6348 );
6349 Ok(())
6350 }
6351
6352 #[tokio::test]
6353 async fn runtime_task_projection_uses_last_completed_message_without_delta_duplication()
6354 -> Result<()> {
6355 let commentary = "Checking fixture state. ".repeat(20);
6356 let task = project_runtime_task(vec![
6357 ("item.started", json!({ "item": { "id": "commentary", "kind": "agent_message" } })),
6358 ("item.delta", json!({ "kind": "agent_message", "delta": commentary })),
6359 ("item.completed", json!({ "item": { "id": "commentary", "kind": "agent_message", "detail": commentary } })),
6360 ("item.started", json!({ "item": { "id": "final", "kind": "agent_message" } })),
6361 ("item.delta", json!({ "kind": "agent_message", "delta": "NOTHING_" })),
6362 ("item.completed", json!({ "item": { "id": "final", "kind": "agent_message", "detail": "NOTHING_TO_REPORT" } })),
6363 ("turn.completed", json!({ "turn": { "status": "completed" } })),
6364 ]).await?;
6365 assert_eq!(task.result_summary.as_deref(), Some("NOTHING_TO_REPORT"));
6366 assert!(task.result_detail_path.is_none());
6367 assert!(
6368 task.timeline.iter().any(|entry| entry.kind == "message"
6369 && entry.summary.starts_with("Checking fixture state."))
6370 );
6371 Ok(())
6372 }
6373
6374 #[tokio::test]
6375 async fn runtime_task_projection_preserves_partial_output_only_for_unfinished_results()
6376 -> Result<()> {
6377 for (status, expected) in [
6378 ("interrupted", "partial final"),
6379 ("failed", "partial final"),
6380 ("completed", "(no textual output)"),
6381 ] {
6382 let task = project_runtime_task(vec![
6383 (
6384 "item.completed",
6385 json!({ "item": { "kind": "agent_message", "detail": "earlier commentary" } }),
6386 ),
6387 (
6388 "item.started",
6389 json!({ "item": { "kind": "agent_message" } }),
6390 ),
6391 (
6392 "item.delta",
6393 json!({ "kind": "agent_message", "delta": "partial final" }),
6394 ),
6395 ("turn.completed", json!({ "turn": { "status": status } })),
6396 ])
6397 .await?;
6398 assert_eq!(task.result_summary.as_deref(), Some(expected), "{status}");
6399 }
6400 Ok(())
6401 }
6402
6403 #[tokio::test]
6404 async fn engine_turn_without_terminal_event_idle_times_out() -> Result<()> {
6405 let runtime = test_runtime_manager().await?;
6406 let thread = runtime
6407 .create_thread(CreateThreadRequest::default())
6408 .await?;
6409 let (tx, rx) = mpsc::channel(64);
6410 tokio::spawn(drain_task_events(rx));
6411 let result = drive_engine_turn(
6412 &runtime,
6413 &thread.id,
6414 "turn_missing",
6415 tx,
6416 CancellationToken::new(),
6417 TaskExecutionLimits::short_for_tests(),
6418 )
6419 .await;
6420 assert_eq!(result.status, TaskStatus::Failed);
6421 assert_eq!(result.terminal_reason, TaskTerminalReason::IdleTimeout);
6422 Ok(())
6423 }
6424
6425 #[tokio::test]
6426 async fn engine_turn_keeps_idle_timeout_when_runtime_interrupts_during_grace() -> Result<()> {
6427 let runtime = Arc::new(test_runtime_manager().await?);
6428 let thread = runtime
6429 .create_thread(CreateThreadRequest::default())
6430 .await?;
6431 let (tx, mut rx) = mpsc::channel(64);
6432 let runtime_for_drive = Arc::clone(&runtime);
6433 let thread_id = thread.id.clone();
6434 let drive = tokio::spawn(async move {
6435 drive_engine_turn(
6436 runtime_for_drive.as_ref(),
6437 &thread_id,
6438 "turn_timeout",
6439 tx,
6440 CancellationToken::new(),
6441 TaskExecutionLimits {
6442 wall_time: Duration::from_secs(2),
6443 idle_progress: Duration::from_millis(80),
6444 cancel_grace: Duration::from_millis(500),
6445 persist_debounce: Duration::from_millis(10),
6446 },
6447 )
6448 .await
6449 });
6450
6451 tokio::time::timeout(Duration::from_secs(2), async {
6452 loop {
6453 match rx.recv().await {
6454 Some(TaskExecutionEvent::Status { message })
6455 if message.contains("idle deadline") =>
6456 {
6457 break;
6458 }
6459 Some(_) => {}
6460 None => panic!("engine task event stream closed before timeout interrupt"),
6461 }
6462 }
6463 })
6464 .await
6465 .context("engine did not request the idle-timeout interrupt")?;
6466
6467 runtime
6468 .emit_event_for_test(
6469 &thread.id,
6470 Some("turn_timeout"),
6471 "turn.completed",
6472 json!({ "turn": { "status": "interrupted" } }),
6473 )
6474 .await?;
6475 let result = drive.await?;
6476 assert_eq!(result.status, TaskStatus::Failed);
6477 assert_eq!(result.terminal_reason, TaskTerminalReason::IdleTimeout);
6478 Ok(())
6479 }
6480
6481 #[tokio::test]
6482 async fn engine_turn_uses_cursor_catchup_for_completed_event() -> Result<()> {
6483 let runtime = test_runtime_manager().await?;
6484 let thread = runtime
6485 .create_thread(CreateThreadRequest::default())
6486 .await?;
6487 runtime
6488 .emit_event_for_test(
6489 &thread.id,
6490 Some("turn_done"),
6491 "turn.completed",
6492 json!({ "turn": { "status": "completed" } }),
6493 )
6494 .await?;
6495 let (tx, rx) = mpsc::channel(64);
6496 tokio::spawn(drain_task_events(rx));
6497 let result = drive_engine_turn(
6498 &runtime,
6499 &thread.id,
6500 "turn_done",
6501 tx,
6502 CancellationToken::new(),
6503 TaskExecutionLimits::short_for_tests(),
6504 )
6505 .await;
6506 assert_eq!(result.status, TaskStatus::Completed);
6507 assert_eq!(result.terminal_reason, TaskTerminalReason::Completed);
6508 Ok(())
6509 }
6510
6511 #[tokio::test]
6512 async fn pending_approval_suspends_idle_and_timeout_denial_settles_failed() -> Result<()> {
6513 // #6118: a run that needs a tool approval must not die as a silent
6514 // idle-timeout cancel; the pending approval suspends the idle
6515 // watchdog, and the bridge's own deadline denial then settles the
6516 // run Failed with the reason recorded.
6517 let runtime = Arc::new(test_runtime_manager().await?);
6518 let thread = runtime
6519 .create_thread(CreateThreadRequest::default())
6520 .await?;
6521 let thread_id = thread.id.clone();
6522 runtime
6523 .emit_event_for_test(
6524 &thread.id,
6525 Some("turn_approval"),
6526 "approval.required",
6527 json!({
6528 "approval_id": "approval_fixture_1",
6529 "tool_call_id": "call_fixture_1",
6530 "tool_name": "shell",
6531 }),
6532 )
6533 .await?;
6534 let (tx, mut rx) = mpsc::channel(64);
6535 let runtime_for_drive = Arc::clone(&runtime);
6536 let drive = tokio::spawn(async move {
6537 drive_engine_turn(
6538 runtime_for_drive.as_ref(),
6539 &thread_id,
6540 "turn_approval",
6541 tx,
6542 CancellationToken::new(),
6543 TaskExecutionLimits {
6544 wall_time: Duration::from_secs(5),
6545 idle_progress: Duration::from_millis(120),
6546 cancel_grace: Duration::from_millis(200),
6547 persist_debounce: Duration::from_millis(10),
6548 },
6549 )
6550 .await
6551 });
6552
6553 // Well past the idle window, the pending approval must keep the run
6554 // alive; the old behavior killed it here with no receipt.
6555 tokio::time::sleep(Duration::from_millis(400)).await;
6556 assert!(
6557 !drive.is_finished(),
6558 "a pending approval must suspend the idle watchdog (#6118)"
6559 );
6560 while let Ok(event) = rx.try_recv() {
6561 if let TaskExecutionEvent::Status { message } = event {
6562 assert!(
6563 !message.contains("idle deadline"),
6564 "no idle interrupt may fire while an approval is pending: {message}"
6565 );
6566 }
6567 }
6568
6569 // The decision window closes: the run settles Failed with the reason.
6570 runtime
6571 .emit_event_for_test(
6572 &thread.id,
6573 Some("turn_approval"),
6574 "approval.timeout",
6575 json!({ "approval_id": "approval_fixture_1", "tool_call_id": "call_fixture_1" }),
6576 )
6577 .await?;
6578 let result = tokio::time::timeout(Duration::from_secs(2), drive)
6579 .await
6580 .context("the decision-window denial must settle the run promptly")??;
6581 assert_eq!(result.status, TaskStatus::Failed);
6582 assert_eq!(result.terminal_reason, TaskTerminalReason::Failed);
6583 assert!(
6584 result
6585 .error
6586 .as_deref()
6587 .is_some_and(|error| error.contains("Tool approval was not answered")),
6588 "the run must record why it stopped, got {:?}",
6589 result.error
6590 );
6591 Ok(())
6592 }
6593
6594 #[tokio::test]
6595 async fn resolved_approval_restores_the_idle_watchdog() -> Result<()> {
6596 // #6118 counter-check: the suspension ends with the decision, so a
6597 // run that then stops making progress is idle-killed exactly as
6598 // before.
6599 let runtime = test_runtime_manager().await?;
6600 let thread = runtime
6601 .create_thread(CreateThreadRequest::default())
6602 .await?;
6603 runtime
6604 .emit_event_for_test(
6605 &thread.id,
6606 Some("turn_approval_resolved"),
6607 "approval.required",
6608 json!({
6609 "approval_id": "approval_fixture_2",
6610 "tool_call_id": "call_fixture_2",
6611 "tool_name": "shell",
6612 }),
6613 )
6614 .await?;
6615 runtime
6616 .emit_event_for_test(
6617 &thread.id,
6618 Some("turn_approval_resolved"),
6619 "approval.decided",
6620 json!({
6621 "approval_id": "approval_fixture_2",
6622 "tool_call_id": "call_fixture_2",
6623 "decision": "allow",
6624 }),
6625 )
6626 .await?;
6627 let (tx, rx) = mpsc::channel(64);
6628 tokio::spawn(drain_task_events(rx));
6629 let result = drive_engine_turn(
6630 &runtime,
6631 &thread.id,
6632 "turn_approval_resolved",
6633 tx,
6634 CancellationToken::new(),
6635 TaskExecutionLimits::short_for_tests(),
6636 )
6637 .await;
6638 assert_eq!(result.status, TaskStatus::Failed);
6639 assert_eq!(result.terminal_reason, TaskTerminalReason::IdleTimeout);
6640 Ok(())
6641 }
6642
6643 #[tokio::test]
6644 async fn silent_running_tool_suppresses_idle_until_it_completes() -> Result<()> {
6645 let runtime = Arc::new(test_runtime_manager().await?);
6646 let thread = runtime
6647 .create_thread(CreateThreadRequest::default())
6648 .await?;
6649 let (tx, mut rx) = mpsc::channel(64);
6650 let runtime_for_drive = Arc::clone(&runtime);
6651 let thread_id = thread.id.clone();
6652 let drive = tokio::spawn(async move {
6653 drive_engine_turn(
6654 runtime_for_drive.as_ref(),
6655 &thread_id,
6656 "turn_silent_tool",
6657 tx,
6658 CancellationToken::new(),
6659 TaskExecutionLimits {
6660 wall_time: Duration::from_secs(5),
6661 idle_progress: Duration::from_millis(80),
6662 cancel_grace: Duration::from_millis(500),
6663 persist_debounce: Duration::from_millis(10),
6664 },
6665 )
6666 .await
6667 });
6668
6669 // Keyed by the journal item id on the started edge; the payload also
6670 // carries the engine tool-use id under "tool" (which must NOT become
6671 // the tracking key: the terminal events only repeat "item").
6672 runtime
6673 .emit_event_for_test(
6674 &thread.id,
6675 Some("turn_silent_tool"),
6676 "item.started",
6677 json!({
6678 "item": { "id": "item_tool_silent", "status": "in_progress" },
6679 "tool": { "id": "call_build", "name": "exec_command", "input": {} }
6680 }),
6681 )
6682 .await?;
6683
6684 // A silent build produces no journal traffic for its whole window;
6685 // far past the idle deadline the watchdog must still not fire.
6686 tokio::time::sleep(Duration::from_millis(400)).await;
6687 let mut seen = Vec::new();
6688 while let Ok(event) = rx.try_recv() {
6689 seen.push(event);
6690 }
6691 assert!(
6692 !seen.iter().any(|event| matches!(
6693 event,
6694 TaskExecutionEvent::Status { message }
6695 if message.contains("idle deadline")
6696 )),
6697 "idle watchdog fired while a tool was still running"
6698 );
6699
6700 runtime
6701 .emit_event_for_test(
6702 &thread.id,
6703 Some("turn_silent_tool"),
6704 "item.completed",
6705 json!({ "item": { "id": "item_tool_silent", "status": "completed" } }),
6706 )
6707 .await?;
6708
6709 // Completion drains the set; silence afterwards must let the idle
6710 // deadline fire again (a drain bug would surface as WallTimeout).
6711 tokio::time::timeout(Duration::from_secs(2), async {
6712 loop {
6713 match rx.recv().await {
6714 Some(TaskExecutionEvent::Status { message })
6715 if message.contains("idle deadline") =>
6716 {
6717 break;
6718 }
6719 Some(_) => {}
6720 None => panic!("task event stream closed before idle deadline"),
6721 }
6722 }
6723 })
6724 .await
6725 .context("idle deadline did not fire after the tool completed")?;
6726 runtime
6727 .emit_event_for_test(
6728 &thread.id,
6729 Some("turn_silent_tool"),
6730 "turn.completed",
6731 json!({ "turn": { "status": "interrupted" } }),
6732 )
6733 .await?;
6734 let result = drive.await?;
6735 assert_eq!(result.status, TaskStatus::Failed);
6736 assert_eq!(result.terminal_reason, TaskTerminalReason::IdleTimeout);
6737 Ok(())
6738 }
6739
6740 #[tokio::test]
6741 async fn silent_running_tool_heartbeats_the_worker_watchdog() -> Result<()> {
6742 // The supervisor's guard (run_task) only sees task events, so the
6743 // steady-state suppression must also publish the in-flight window on
6744 // the event channel — otherwise the production path still idles out
6745 // even though this loop's own guard is satisfied.
6746 let runtime = Arc::new(test_runtime_manager().await?);
6747 let thread = runtime
6748 .create_thread(CreateThreadRequest::default())
6749 .await?;
6750 let (tx, mut rx) = mpsc::channel(64);
6751 let runtime_for_drive = Arc::clone(&runtime);
6752 let thread_id = thread.id.clone();
6753 let drive = tokio::spawn(async move {
6754 drive_engine_turn(
6755 runtime_for_drive.as_ref(),
6756 &thread_id,
6757 "turn_heartbeat",
6758 tx,
6759 CancellationToken::new(),
6760 TaskExecutionLimits {
6761 wall_time: Duration::from_secs(5),
6762 idle_progress: Duration::from_millis(80),
6763 cancel_grace: Duration::from_millis(500),
6764 persist_debounce: Duration::from_millis(10),
6765 },
6766 )
6767 .await
6768 });
6769
6770 runtime
6771 .emit_event_for_test(
6772 &thread.id,
6773 Some("turn_heartbeat"),
6774 "item.started",
6775 json!({
6776 "item": { "id": "item_tool_hb", "status": "in_progress" },
6777 "tool": { "id": "call_build", "name": "exec_command", "input": {} }
6778 }),
6779 )
6780 .await?;
6781
6782 tokio::time::timeout(Duration::from_secs(2), async {
6783 loop {
6784 match rx.recv().await {
6785 Some(TaskExecutionEvent::ToolHeartbeat) => break,
6786 Some(_) => {}
6787 None => panic!("task event stream closed before any heartbeat"),
6788 }
6789 }
6790 })
6791 .await
6792 .context("no ToolHeartbeat reached the worker channel while a tool was in flight")?;
6793
6794 runtime
6795 .emit_event_for_test(
6796 &thread.id,
6797 Some("turn_heartbeat"),
6798 "item.completed",
6799 json!({ "item": { "id": "item_tool_hb", "status": "completed" } }),
6800 )
6801 .await?;
6802 runtime
6803 .emit_event_for_test(
6804 &thread.id,
6805 Some("turn_heartbeat"),
6806 "turn.completed",
6807 json!({ "turn": { "status": "completed" } }),
6808 )
6809 .await?;
6810 let result = drive.await?;
6811 assert_eq!(result.status, TaskStatus::Completed);
6812 Ok(())
6813 }
6814
6815 #[tokio::test]
6816 async fn hung_running_tool_still_hits_the_wall_budget() -> Result<()> {
6817 let runtime = Arc::new(test_runtime_manager().await?);
6818 let thread = runtime
6819 .create_thread(CreateThreadRequest::default())
6820 .await?;
6821 let (tx, mut rx) = mpsc::channel(64);
6822 let runtime_for_drive = Arc::clone(&runtime);
6823 let thread_id = thread.id.clone();
6824 let drive = tokio::spawn(async move {
6825 drive_engine_turn(
6826 runtime_for_drive.as_ref(),
6827 &thread_id,
6828 "turn_hung_tool",
6829 tx,
6830 CancellationToken::new(),
6831 TaskExecutionLimits {
6832 wall_time: Duration::from_millis(300),
6833 idle_progress: Duration::from_millis(80),
6834 cancel_grace: Duration::from_millis(500),
6835 persist_debounce: Duration::from_millis(10),
6836 },
6837 )
6838 .await
6839 });
6840
6841 runtime
6842 .emit_event_for_test(
6843 &thread.id,
6844 Some("turn_hung_tool"),
6845 "item.started",
6846 json!({
6847 "item": { "id": "item_tool_hung", "status": "in_progress" },
6848 "tool": { "id": "call_hang", "name": "exec_command", "input": {} }
6849 }),
6850 )
6851 .await?;
6852
6853 // The running-tool suppression only defers the idle watchdog; the
6854 // wall-time budget remains the backstop for a tool that never
6855 // completes.
6856 tokio::time::timeout(Duration::from_secs(2), async {
6857 loop {
6858 match rx.recv().await {
6859 Some(TaskExecutionEvent::Status { message })
6860 if message.contains("wall-time") =>
6861 {
6862 break;
6863 }
6864 Some(_) => {}
6865 None => panic!("task event stream closed before wall deadline"),
6866 }
6867 }
6868 })
6869 .await
6870 .context("wall budget did not fire during a hung tool")?;
6871 runtime
6872 .emit_event_for_test(
6873 &thread.id,
6874 Some("turn_hung_tool"),
6875 "turn.completed",
6876 json!({ "turn": { "status": "interrupted" } }),
6877 )
6878 .await?;
6879 let result = drive.await?;
6880 assert_eq!(result.status, TaskStatus::Failed);
6881 assert_eq!(result.terminal_reason, TaskTerminalReason::WallTimeout);
6882 Ok(())
6883 }
6884
6885 /// Mirrors the fixed runtime contract: a ToolStarted edge, then a silent
6886 /// window longer than the idle deadline fed only by ToolHeartbeat, then
6887 /// completion. Pins the supervisor side of the heartbeat chain — the
6888 /// worker idle watchdog must count heartbeats as progress.
6889 ///
6890 /// Margins: heartbeats every 30 ms across a 1.2 s window against a 500 ms
6891 /// idle deadline. The window still outlasts the deadline (so the test
6892 /// fails if heartbeats stop counting), while a loaded runner would need a
6893 /// half-second stall to starve it; the 150 ms default deadline did starve
6894 /// on hosted Windows.
6895 struct ToolHeartbeatExecutor;
6896
6897 const HEARTBEAT_TEST_IDLE_PROGRESS: Duration = Duration::from_millis(500);
6898 const HEARTBEAT_TEST_TICKS: u32 = 40;
6899
6900 #[async_trait]
6901 impl TaskExecutor for ToolHeartbeatExecutor {
6902 async fn execute(
6903 &self,
6904 _task: ExecutionTask,
6905 events: mpsc::Sender<TaskExecutionEvent>,
6906 cancel: CancellationToken,
6907 ) -> TaskExecutionResult {
6908 let _ = events
6909 .send(TaskExecutionEvent::ToolStarted {
6910 id: "tool_hb".to_string(),
6911 name: "exec_command".to_string(),
6912 input: serde_json::json!({}),
6913 })
6914 .await;
6915 for _ in 0..HEARTBEAT_TEST_TICKS {
6916 sleep(Duration::from_millis(30)).await;
6917 if cancel.is_cancelled() {
6918 return TaskExecutionResult {
6919 status: TaskStatus::Canceled,
6920 result_text: None,
6921 error: None,
6922 terminal_reason: TaskTerminalReason::Canceled,
6923 };
6924 }
6925 let _ = events.send(TaskExecutionEvent::ToolHeartbeat).await;
6926 }
6927 let _ = events
6928 .send(TaskExecutionEvent::ToolCompleted {
6929 id: "tool_hb".to_string(),
6930 name: "exec_command".to_string(),
6931 success: true,
6932 output: "build ok".to_string(),
6933 metadata: None,
6934 })
6935 .await;
6936 TaskExecutionResult {
6937 status: TaskStatus::Completed,
6938 result_text: Some("done".to_string()),
6939 error: None,
6940 terminal_reason: TaskTerminalReason::Completed,
6941 }
6942 }
6943 }
6944
6945 #[tokio::test]
6946 async fn worker_supervisor_honors_tool_heartbeats_during_silent_tools() -> Result<()> {
6947 let root = tempfile::tempdir()?;
6948 let mut config = short_test_config(root.path().to_path_buf());
6949 config.execution_limits.wall_time = Duration::from_secs(5);
6950 config.execution_limits.idle_progress = HEARTBEAT_TEST_IDLE_PROGRESS;
6951 let manager =
6952 TaskManager::start_with_executor(config, Arc::new(ToolHeartbeatExecutor)).await?;
6953
6954 let task = manager
6955 .add_task(NewTaskRequest::from_prompt("silent build with heartbeats"))
6956 .await?;
6957 let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?;
6958 assert_eq!(
6959 finished.status,
6960 TaskStatus::Completed,
6961 "heartbeats from a silent in-flight tool must keep the worker idle watchdog fed"
6962 );
6963 assert!(
6964 !finished
6965 .timeline
6966 .iter()
6967 .any(|entry| entry.summary.contains("heartbeat")),
6968 "heartbeats are supervisor-side liveness only and must not surface on the timeline"
6969 );
6970 Ok(())
6971 }
6972
6973 #[tokio::test]
6974 async fn tool_heartbeat_is_liveness_only_and_never_persists() -> Result<()> {
6975 // The heartbeat arrives up to ~5x/s during a silent build; if it
6976 // counted as persist-urgent it would rewrite the whole task record
6977 // on every tick while holding the manager-wide state lock. The
6978 // exclusion list must keep treating it as transient state.
6979 assert!(!execution_event_persist_urgent(
6980 &TaskExecutionEvent::ToolHeartbeat
6981 ));
6982 assert!(execution_event_persist_urgent(
6983 &TaskExecutionEvent::ToolStarted {
6984 id: "item-1".into(),
6985 name: "bash".into(),
6986 input: serde_json::json!({}),
6987 }
6988 ));
6989
6990 // Wiring-level pin: applying a heartbeat to a running task leaves the
6991 // record unpersisted, while a real lifecycle edge still flushes. The
6992 // executor hangs so the task stays Running (default limits keep the
6993 // supervisor idle watchdog far away) while the events are applied.
6994 let root = tempfile::tempdir()?;
6995 let manager = TaskManager::start_with_executor(
6996 test_config(root.path().to_path_buf()),
6997 Arc::new(DeafHangExecutor),
6998 )
6999 .await?;
7000 let task = manager
7001 .add_task(NewTaskRequest::from_prompt("heartbeat persistence pin"))
7002 .await?;
7003 let running = tokio::time::timeout(Duration::from_secs(5), async {
7004 loop {
7005 match manager.get_task(&task.id).await {
7006 Ok(record) if record.status == TaskStatus::Running => break Some(record),
7007 Ok(_) => {}
7008 Err(_) => break None,
7009 }
7010 sleep(Duration::from_millis(10)).await;
7011 }
7012 })
7013 .await
7014 .context("task never reached Running for the persistence pin")?
7015 .context("task lookup failed before the persistence pin")?;
7016 assert_eq!(running.status, TaskStatus::Running);
7017
7018 let outcome = manager
7019 .apply_execution_event(&task.id, TaskExecutionEvent::ToolHeartbeat)
7020 .await?;
7021 assert!(
7022 !outcome.persisted,
7023 "a liveness-only heartbeat must not trigger a task-record write"
7024 );
7025
7026 let outcome = manager
7027 .apply_execution_event(
7028 &task.id,
7029 TaskExecutionEvent::ToolStarted {
7030 id: "item-1".into(),
7031 name: "bash".into(),
7032 input: serde_json::json!({}),
7033 },
7034 )
7035 .await?;
7036 assert!(
7037 outcome.persisted,
7038 "a real tool lifecycle edge must still flush the record"
7039 );
7040 Ok(())
7041 }
7042
7043 #[tokio::test]
7044 async fn engine_turn_prefers_runtime_terminal_over_cancel_grace() -> Result<()> {
7045 let runtime = test_runtime_manager().await?;
7046 let thread = runtime
7047 .create_thread(CreateThreadRequest::default())
7048 .await?;
7049 let cancel = CancellationToken::new();
7050 cancel.cancel();
7051 runtime
7052 .emit_event_for_test(
7053 &thread.id,
7054 Some("turn_done"),
7055 "turn.completed",
7056 json!({ "turn": { "status": "interrupted" } }),
7057 )
7058 .await?;
7059 let (tx, rx) = mpsc::channel(64);
7060 tokio::spawn(drain_task_events(rx));
7061 let result = drive_engine_turn(
7062 &runtime,
7063 &thread.id,
7064 "turn_done",
7065 tx,
7066 cancel,
7067 TaskExecutionLimits::short_for_tests(),
7068 )
7069 .await;
7070 assert_eq!(result.status, TaskStatus::Canceled);
7071 assert_eq!(result.terminal_reason, TaskTerminalReason::Canceled);
7072 Ok(())
7073 }
7074
7075 #[tokio::test]
7076 async fn runtime_store_failure_event_reaches_the_task_timeline() -> Result<()> {
7077 // #5931: the runtime's own store fault lands in the task timeline,
7078 // and a terminal one stops the driver instead of idling it out.
7079 let (tx, mut rx) = mpsc::channel(8);
7080 let mut final_text = RuntimeTaskOutput::default();
7081 let path = "/tmp/runtime/turns/turn_store.json";
7082 let event = RuntimeEventRecord {
7083 schema_version: 1,
7084 seq: 7,
7085 timestamp: Utc::now(),
7086 thread_id: "thr_store".to_string(),
7087 turn_id: Some("turn_store".to_string()),
7088 item_id: None,
7089 event: RUNTIME_STORE_FAILURE_EVENT.to_string(),
7090 payload: json!({
7091 "operation": "read",
7092 "record_kind": "turn",
7093 "record_id": "turn_store",
7094 "path": path,
7095 "error": format!("Failed to read turn {path}: No such file"),
7096 "reason": "No such file",
7097 "next_action": format!("Move {path} aside (or delete it) and retry."),
7098 "message": format!(
7099 "Session runtime store: turn turn_store at {path} could not be read: No such file. Move {path} aside (or delete it) and retry."
7100 ),
7101 }),
7102 };
7103
7104 assert!(
7105 ingest_runtime_event(&event, &mut final_text, &tx)
7106 .await
7107 .is_none(),
7108 "a non-terminal store fault leaves the driver waiting"
7109 );
7110 let mut saw_error = false;
7111 while let Ok(received) = rx.try_recv() {
7112 if let TaskExecutionEvent::Error { message } = received {
7113 assert!(message.contains(path), "{message}");
7114 saw_error = true;
7115 }
7116 }
7117 assert!(saw_error, "store fault missing from the task timeline");
7118
7119 let mut terminal = event.clone();
7120 terminal.payload["terminal"] = json!(true);
7121 let (status, error) = ingest_runtime_event(&terminal, &mut final_text, &tx)
7122 .await
7123 .expect("a terminal store fault ends the turn");
7124 assert_eq!(status, RuntimeTurnStatus::Failed);
7125 assert!(error.is_some_and(|message| message.contains(path)));
7126 Ok(())
7127 }
7128 }
7129
7130 #[cfg(test)]
7131 mod ownership_tests;
7132
7132 lines RUST