From 589917ad7d828e9c0567426c4cbc08a3a14f276e Mon Sep 17 00:00:00 2001 From: asto Date: Sat, 26 Sep 2026 13:37:29 +0800 Subject: [PATCH 1/5] Revert "revert(tasks): move the human-wait clock pause to its own PR" This reverts commit ce5a01fe93fde5a17e94f44b9127d3c00d0782c2. Signed-off-by: asto --- crates/tui/src/task_manager.rs | 375 +++++++++++++++++++++++++++++++-- docs/SUBAGENTS.md | 15 +- docs/zh_hans/SUBAGENTS.md | 2 +- 3 files changed, 364 insertions(+), 28 deletions(-) diff --git a/crates/tui/src/task_manager.rs b/crates/tui/src/task_manager.rs index f6c42e716f..708054b632 100644 --- a/crates/tui/src/task_manager.rs +++ b/crates/tui/src/task_manager.rs @@ -102,6 +102,7 @@ pub enum TaskTerminalReason { CancelTimeout, Shutdown, WallTimeout, + HumanWaitTimeout, IdleTimeout, Failed, } @@ -115,6 +116,7 @@ impl TaskTerminalReason { Self::CancelTimeout => "cancel_timeout", Self::Shutdown => "shutdown", Self::WallTimeout => "wall_timeout", + Self::HumanWaitTimeout => "human_wait_timeout", Self::IdleTimeout => "idle_timeout", Self::Failed => "failed", } @@ -125,7 +127,9 @@ impl TaskTerminalReason { match self { Self::Completed => TaskStatus::Completed, Self::Canceled | Self::CancelTimeout | Self::Shutdown => TaskStatus::Canceled, - Self::WallTimeout | Self::IdleTimeout | Self::Failed => TaskStatus::Failed, + Self::WallTimeout | Self::HumanWaitTimeout | Self::IdleTimeout | Self::Failed => { + TaskStatus::Failed + } } } @@ -141,6 +145,9 @@ impl TaskTerminalReason { Self::WallTimeout => { "Task exceeded its wall-time deadline without completing".to_string() } + Self::HumanWaitTimeout => { + "Task stayed parked on a human answer past the fail-safe cap".to_string() + } Self::IdleTimeout => { "Task made no model or tool progress before the idle deadline".to_string() } @@ -464,6 +471,11 @@ pub struct TaskExecutionLimits { pub idle_progress: Duration, pub cancel_grace: Duration, pub persist_debounce: Duration, + /// Fail-safe ceiling for one continuous human wait (a pending approval + /// or user-input prompt). Human-paced waits are excluded from + /// `wall_time` — a slow answer must not end the task — but a prompt + /// nobody will ever answer must still release the worker eventually. + pub human_wait_cap: Duration, } impl Default for TaskExecutionLimits { @@ -473,6 +485,7 @@ impl Default for TaskExecutionLimits { idle_progress: Duration::from_secs(2 * 60), cancel_grace: Duration::from_secs(5), persist_debounce: Duration::from_millis(250), + human_wait_cap: Duration::from_secs(24 * 60 * 60), } } } @@ -485,6 +498,7 @@ impl TaskExecutionLimits { idle_progress: Duration::from_millis(150), cancel_grace: Duration::from_millis(50), persist_debounce: Duration::from_millis(10), + human_wait_cap: Duration::from_millis(300), } } } @@ -495,6 +509,12 @@ struct ExecutionGuard { last_progress_at: Instant, interrupt_at: Option, interrupt_reason: Option, + /// Open human wait (pending approval or user-input prompt), if any. + /// Human-paced wall time is excluded from `wall_time` exactly like the + /// engine-side turn budget excludes it — an answer submitted past the + /// budget must complete the task, not be collected and dropped. + human_wait_since: Option, + human_wait_total: Duration, limits: TaskExecutionLimits, } @@ -512,10 +532,50 @@ impl ExecutionGuard { last_progress_at: now, interrupt_at: None, interrupt_reason: None, + human_wait_since: None, + human_wait_total: Duration::ZERO, limits, } } + /// Mark the start of a human-paced wait (idempotent). Only the first + /// begin of an unmatched begin/end pair takes effect. + /// + /// Windows merge rather than stack: the first begin sets the origin and + /// any end closes the whole window. That is sound because the engine + /// blocks inside a single prompt at a time — approval and user-input + /// waits are serial within a turn, and parallel batches exclude + /// interactive tools — so overlapping windows never reach this guard + /// today. If a future change lets nested prompts surface concurrently, + /// switch to paired (counted) windows first. + fn begin_human_wait(&mut self, now: Instant) { + if self.human_wait_since.is_none() { + self.human_wait_since = Some(now); + } + } + + /// Mark the end of a human-paced wait (idempotent). The idle budget + /// restarts fresh: the engine is expected to think and act again after + /// an answer, and window silence was expected, not idleness. + fn end_human_wait(&mut self, now: Instant) { + if let Some(since) = self.human_wait_since.take() { + self.human_wait_total += now.saturating_duration_since(since); + self.last_progress_at = now; + } + } + + /// Wall time attributable to the engine itself: excludes every closed + /// human-wait window and the currently open one. + fn engine_wall_elapsed(&self, now: Instant) -> Duration { + let open = self + .human_wait_since + .map(|since| now.saturating_duration_since(since)) + .unwrap_or_default(); + now.saturating_duration_since(self.started_at) + .saturating_sub(self.human_wait_total) + .saturating_sub(open) + } + fn note_progress(&mut self, now: Instant) { self.last_progress_at = now; } @@ -542,12 +602,24 @@ impl ExecutionGuard { }; } - let wall_elapsed = now.saturating_duration_since(self.started_at); + let wall_elapsed = self.engine_wall_elapsed(now); + // The human-wait exclusion is bounded by a fail-safe cap: a prompt + // that sits unanswered far past any plausible human pace still + // releases the worker instead of parking it forever. While the + // window is open both clocks pause — silence during a human decision + // is expected, not idleness — and the idle budget restarts when the + // window closes (see `end_human_wait`). let idle_elapsed = now.saturating_duration_since(self.last_progress_at); let pending = if shutdown { Some(TaskTerminalReason::Shutdown) } else if cancel { Some(TaskTerminalReason::Canceled) + } else if let Some(since) = self.human_wait_since { + if now.saturating_duration_since(since) >= self.limits.human_wait_cap { + Some(TaskTerminalReason::HumanWaitTimeout) + } else { + None + } } else if wall_elapsed >= self.limits.wall_time { Some(TaskTerminalReason::WallTimeout) } else if idle_elapsed >= self.limits.idle_progress { @@ -559,12 +631,18 @@ impl ExecutionGuard { return GuardAction::Interrupt { reason }; } - let wait = self + let mut wait = self .limits .wall_time .saturating_sub(wall_elapsed) - .min(self.limits.idle_progress.saturating_sub(idle_elapsed)) - .min(EVENT_CATCHUP_POLL); + .min(self.limits.idle_progress.saturating_sub(idle_elapsed)); + if self.human_wait_since.is_some() { + // Both clocks are paused; wake at the fail-safe cadence. + wait = self.limits.human_wait_cap.saturating_sub( + now.saturating_duration_since(self.human_wait_since.unwrap_or(now)), + ); + } + let wait = wait.min(EVENT_CATCHUP_POLL); GuardAction::Run { wait: wait.max(Duration::from_millis(1)), } @@ -575,9 +653,11 @@ impl ExecutionGuard { return result; } match self.interrupt_reason { - Some(reason @ (TaskTerminalReason::WallTimeout | TaskTerminalReason::IdleTimeout)) => { - TaskExecutionResult::from_reason(reason, result.result_text) - } + Some( + reason @ (TaskTerminalReason::WallTimeout + | TaskTerminalReason::HumanWaitTimeout + | TaskTerminalReason::IdleTimeout), + ) => TaskExecutionResult::from_reason(reason, result.result_text), _ => result, } } @@ -742,6 +822,12 @@ pub enum TaskExecutionEvent { /// Journal item id of one in-flight tool. id: String, }, + /// A human-paced wait (pending approval or user-input prompt) opened or + /// closed, so the worker supervisor's wall clock can exclude it exactly + /// like the turn loop's own guard does. Supervisor-side signal only: + /// never persisted and never shown on the task timeline. + HumanWaitStarted, + HumanWaitEnded, ToolCompleted { id: String, name: String, @@ -959,10 +1045,14 @@ async fn drive_engine_turn( match event.event.as_str() { "item.started" => { // Journal tool items carry a `tool` payload key; - // non-tool items do not. A silent context compaction is - // therefore invisible to this set (disclosed gap: a - // background task mid-compaction can still hit the idle - // deadline — compaction has no heartbeat of its own). + // non-tool items do not. Two disclosed blind spots + // remain: a silent context compaction (no heartbeat of + // its own, so a background task mid-compaction can hit + // the idle deadline), and the pre-execution window + // before the first tool item starts (approval + // scheduling, MCP discovery) — a wait there is bounded + // by the idle/wall deadlines rather than the human-wait + // exclusion below, which only sees journaled prompts. if event.payload.get("tool").is_some() && let Some(item_id) = journal_item_id(&event) { @@ -974,6 +1064,26 @@ async fn drive_engine_turn( running_tools.remove(&item_id); } } + // A pending approval or user-input prompt is a human-paced + // wait: exclude it from the wall clock exactly like the + // engine-side turn budget does, and publish the window so + // the worker supervisor's guard follows in lockstep. Every + // wait exit emits one of the paired closing events (the + // decision, the interrupt deny, the legacy replayed + // timeout, or the cancellation), so the window cannot stick + // open; the guard's fail-safe cap bounds the pathological + // case regardless. + "approval.required" | "user_input.required" => { + guard.begin_human_wait(Instant::now()); + emit_task_event(&events, TaskExecutionEvent::HumanWaitStarted).await; + } + "approval.decided" + | "approval.timeout" + | "user_input.answered" + | "user_input.canceled" => { + guard.end_human_wait(Instant::now()); + emit_task_event(&events, TaskExecutionEvent::HumanWaitEnded).await; + } _ => {} } if runtime_event_is_progress(&event) { @@ -2060,13 +2170,27 @@ impl TaskManager { if execution_event_is_progress(&event) { guard.note_progress(Instant::now()); } - // Liveness-only heartbeats never mutate the record: short-circuit - // before the state lock so a silent build's ~5 ticks/s do not take - // the manager-wide lock for a no-op, and so they do not mark the - // record dirty (a spurious dirty would only arm a redundant - // debounce flush). - if matches!(event, TaskExecutionEvent::ToolHeartbeat { .. }) { - return; + // Liveness-only signals never mutate the record: short-circuit + // before the state lock so a silent build's ~5 heartbeat ticks/s do + // not take the manager-wide lock for a no-op, and so they do not + // mark the record dirty (a spurious dirty would only arm a redundant + // debounce flush). Note this does not change when real mutations + // persist: the run loop rebuilds the debounce timer on every event, + // so while a tool streams heartbeats the debounce stays starved and + // persistence waits for the stream to quiet either way. + // The human-wait pair moves the supervisor's wall clock in lockstep + // with the turn loop's own guard. + match event { + TaskExecutionEvent::ToolHeartbeat { .. } => return, + TaskExecutionEvent::HumanWaitStarted => { + guard.begin_human_wait(Instant::now()); + return; + } + TaskExecutionEvent::HumanWaitEnded => { + guard.end_human_wait(Instant::now()); + return; + } + _ => {} } append_message_delta(accumulated_result_text, &event); match self.apply_execution_event(task_id, event).await { @@ -2184,6 +2308,8 @@ impl TaskManager { // Supervisor-side liveness only: recording it would put a // timeline entry behind every poll tick of a silent build. TaskExecutionEvent::ToolHeartbeat { .. } => {} + // Supervisor-side wall-clock signal only (see the variant doc). + TaskExecutionEvent::HumanWaitStarted | TaskExecutionEvent::HumanWaitEnded => {} TaskExecutionEvent::ToolCompleted { id, name, @@ -2340,6 +2466,7 @@ impl TaskManager { | TaskTerminalReason::Shutdown => result.terminal_reason.receipt_message(), TaskTerminalReason::Failed | TaskTerminalReason::WallTimeout + | TaskTerminalReason::HumanWaitTimeout | TaskTerminalReason::IdleTimeout | TaskTerminalReason::CancelTimeout => format!( "{}: {}", @@ -2709,6 +2836,10 @@ fn execution_event_persist_urgent(event: &TaskExecutionEvent) -> bool { // rewrite the whole task record on every tick while holding the // manager-wide state lock. | TaskExecutionEvent::ToolHeartbeat { .. } + // Supervisor-side wall-clock signal for human waits (see the + // variant doc): carries no record state. + | TaskExecutionEvent::HumanWaitStarted + | TaskExecutionEvent::HumanWaitEnded ) } @@ -4804,6 +4935,7 @@ mod tests { idle_progress: Duration::from_millis(80), cancel_grace: Duration::from_millis(500), persist_debounce: Duration::from_millis(10), + human_wait_cap: Duration::from_secs(24 * 60 * 60), }, ) .await @@ -4860,6 +4992,7 @@ mod tests { idle_progress: Duration::from_millis(80), cancel_grace: Duration::from_millis(500), persist_debounce: Duration::from_millis(10), + human_wait_cap: Duration::from_secs(24 * 60 * 60), }, ) .await @@ -4935,6 +5068,208 @@ mod tests { assert_eq!(result.terminal_reason, TaskTerminalReason::IdleTimeout); Ok(()) } + #[tokio::test] + async fn human_wait_is_excluded_from_the_task_wall_clock() -> Result<()> { + // The engine journals the tool item before its approval/user-input + // gate, so the production shape is: tool in flight, then a prompt + // opens a human-paced window. An answer submitted past the wall + // budget must complete the task, exactly like the engine-side turn + // budget excludes the same window. + let runtime = Arc::new(test_runtime_manager().await?); + let thread = runtime + .create_thread(CreateThreadRequest::default()) + .await?; + let (tx, mut rx) = mpsc::channel(64); + let runtime_for_drive = Arc::clone(&runtime); + let thread_id = thread.id.clone(); + let drive = tokio::spawn(async move { + drive_engine_turn( + runtime_for_drive.as_ref(), + &thread_id, + "turn_human_wait", + tx, + CancellationToken::new(), + TaskExecutionLimits { + wall_time: Duration::from_millis(400), + idle_progress: Duration::from_secs(10), + cancel_grace: Duration::from_millis(500), + persist_debounce: Duration::from_millis(10), + human_wait_cap: Duration::from_secs(24 * 60 * 60), + }, + ) + .await + }); + + runtime + .emit_event_for_test( + &thread.id, + Some("turn_human_wait"), + "item.started", + json!({ + "item": { "id": "item_human_wait", "status": "in_progress" }, + "tool": { "id": "call_1", "name": "request_user_input", "input": {} } + }), + ) + .await?; + runtime + .emit_event_for_test( + &thread.id, + Some("turn_human_wait"), + "user_input.required", + json!({ "id": "ui_1", "request": {} }), + ) + .await?; + + // Far past the wall budget the decision is still pending: no + // wall-time interrupt may fire, and the window must be published + // for the worker supervisor's guard. + tokio::time::sleep(Duration::from_millis(900)).await; + let mut saw_window_signal = false; + while let Ok(event) = rx.try_recv() { + if matches!(event, TaskExecutionEvent::HumanWaitStarted) { + saw_window_signal = true; + } + assert!( + !matches!(event, TaskExecutionEvent::Status { message } + if message.contains("wall-time")), + "wall-time interrupt fired while a human decision was pending" + ); + } + assert!( + saw_window_signal, + "the human-wait window must be published on the task event channel" + ); + + // The answer arrives past the budget: the window closes, the tool + // completes, and the turn completes with no wall-time anywhere. + runtime + .emit_event_for_test( + &thread.id, + Some("turn_human_wait"), + "user_input.answered", + json!({ "id": "ui_1" }), + ) + .await?; + runtime + .emit_event_for_test( + &thread.id, + Some("turn_human_wait"), + "item.completed", + json!({ "item": { "id": "item_human_wait", "status": "completed" } }), + ) + .await?; + runtime + .emit_event_for_test( + &thread.id, + Some("turn_human_wait"), + "turn.completed", + json!({ "turn": { "status": "completed" } }), + ) + .await?; + let result = drive.await?; + assert_eq!(result.status, TaskStatus::Completed); + assert_eq!(result.terminal_reason, TaskTerminalReason::Completed); + Ok(()) + } + + #[tokio::test] + async fn human_wait_fail_safe_cap_releases_the_worker() -> Result<()> { + // A prompt nobody will ever answer must not park the worker + // forever: past the fail-safe cap the task ends with the dedicated + // terminal reason even though wall and idle never elapsed. + let runtime = Arc::new(test_runtime_manager().await?); + let thread = runtime + .create_thread(CreateThreadRequest::default()) + .await?; + let (tx, mut rx) = mpsc::channel(64); + let runtime_for_drive = Arc::clone(&runtime); + let thread_id = thread.id.clone(); + let drive = tokio::spawn(async move { + drive_engine_turn( + runtime_for_drive.as_ref(), + &thread_id, + "turn_unanswered", + tx, + CancellationToken::new(), + TaskExecutionLimits { + wall_time: Duration::from_secs(30), + idle_progress: Duration::from_secs(30), + cancel_grace: Duration::from_millis(50), + persist_debounce: Duration::from_millis(10), + human_wait_cap: Duration::from_millis(200), + }, + ) + .await + }); + + runtime + .emit_event_for_test( + &thread.id, + Some("turn_unanswered"), + "approval.required", + json!({ "id": "ap_1", "tool_name": "exec_command" }), + ) + .await?; + + let result = tokio::time::timeout(Duration::from_secs(5), drive) + .await + .expect("the fail-safe cap must release the turn loop")?; + assert_eq!(result.status, TaskStatus::Failed); + assert_eq!(result.terminal_reason, TaskTerminalReason::HumanWaitTimeout); + let _ = rx.try_recv(); + Ok(()) + } + + /// Executor that opens a supervisor-side human wait for `delay` and + /// only then completes — the task survives a wall budget far shorter + /// than the delay if and only if the supervisor's guard honors the + /// HumanWaitStarted/Ended signals. + struct SupervisedHumanWaitExecutor { + delay: Duration, + } + + #[async_trait] + impl TaskExecutor for SupervisedHumanWaitExecutor { + async fn execute( + &self, + _task: ExecutionTask, + events: mpsc::Sender, + _cancel: CancellationToken, + ) -> TaskExecutionResult { + let _ = events.send(TaskExecutionEvent::HumanWaitStarted).await; + tokio::time::sleep(self.delay).await; + let _ = events.send(TaskExecutionEvent::HumanWaitEnded).await; + TaskExecutionResult { + status: TaskStatus::Completed, + result_text: None, + error: None, + terminal_reason: TaskTerminalReason::Completed, + } + } + } + + #[tokio::test] + async fn supervisor_wall_clock_follows_human_wait_signals() -> Result<()> { + let root = std::env::temp_dir().join(format!("deepseek-task-test-{}", Uuid::new_v4())); + let mut config = short_test_config(root); + // Short config's fail-safe cap (300ms) is shorter than the wait + // below; this test pins the wall exclusion, not the fail-safe. + config.execution_limits.human_wait_cap = Duration::from_secs(30); + let manager = TaskManager::start_with_executor( + config, + Arc::new(SupervisedHumanWaitExecutor { + delay: Duration::from_millis(900), + }), + ) + .await?; + let task = manager + .add_task(NewTaskRequest::from_prompt("human wait past wall")) + .await?; + let finished = wait_for_terminal_state(&manager, &task.id, Duration::from_secs(10)).await?; + assert_eq!(finished.status, TaskStatus::Completed); + assert_eq!(finished.terminal_reason.as_deref(), Some("completed")); + Ok(()) + } #[tokio::test] async fn silent_running_tool_heartbeats_the_worker_watchdog() -> Result<()> { @@ -4961,6 +5296,7 @@ mod tests { idle_progress: Duration::from_millis(80), cancel_grace: Duration::from_millis(500), persist_debounce: Duration::from_millis(10), + human_wait_cap: Duration::from_secs(24 * 60 * 60), }, ) .await @@ -5036,6 +5372,7 @@ mod tests { idle_progress: Duration::from_millis(80), cancel_grace: Duration::from_millis(500), persist_debounce: Duration::from_millis(10), + human_wait_cap: Duration::from_secs(24 * 60 * 60), }, ) .await diff --git a/docs/SUBAGENTS.md b/docs/SUBAGENTS.md index e40c7a4ca8..6f1af98b27 100644 --- a/docs/SUBAGENTS.md +++ b/docs/SUBAGENTS.md @@ -596,14 +596,13 @@ progress at step boundaries, not mid-tool, so a silent build or MCP call must not be cleaned up mid-run), so neither a configured long model request nor a long in-flight tool is cancelled before its own timeout can fire. -These floors only keep the heartbeat from firing early; they do not outrank -the wall clocks above them. Every sub-agent runs under its own wall time -(`default_wall_time_secs`, 1800 seconds by default), and inside a durable task -the task's `wall_time` (default 30 minutes, measured from task start, and -running while a prompt waits on a human) bounds the whole run. A maximal -1800-second tool started late in a child or a task can therefore still be -interrupted by either deadline; raising `execute_timeout` above the remaining -wall time has no effect. +Inside a durable task these budgets are additionally bounded by the task's +`wall_time` (default 30 minutes, measured from task start): the wall clock +pauses while a pending approval or user-input prompt waits on a human (with a +24-hour fail-safe cap per prompt), but it otherwise wins over an in-flight +sub-agent, so a maximal 1800-second tool started late in a task can still be +interrupted by the task deadline. Raising `execute_timeout` or the tool +budget above the remaining task wall has no effect inside a task. ## Lifecycle diff --git a/docs/zh_hans/SUBAGENTS.md b/docs/zh_hans/SUBAGENTS.md index fd1d73b86b..e698f1f1ba 100644 --- a/docs/zh_hans/SUBAGENTS.md +++ b/docs/zh_hans/SUBAGENTS.md @@ -322,7 +322,7 @@ heartbeat_timeout_secs = 300 # 钳制到 30..=3600 有效心跳至少保持在解析后的 `api_timeout_secs` 与内置子代理工具超时之上各 30 秒(子代理只在步骤边界记录进度,不会在工具执行中途记录,因此静默的构建或 MCP 调用不能在运行中被清理),所以配置的长模型请求与长时间在途工具都不会在自己的超时触发之前被取消。 -这些下限只保证心跳不会提前触发,并不凌驾于其上的墙钟:每个子代理都受自身墙钟(`default_wall_time_secs`,默认 1800 秒)约束;在持久任务内部,任务 `wall_time`(默认 30 分钟,自任务启动起算,等待人工回应提示期间照常计时)约束整个运行。因此在子代理或任务后段启动的满额 1800 秒工具仍可能被任一截止时间中断;把 `execute_timeout` 调到超过剩余墙钟不会生效。 +在持久任务内部,这些预算还受任务 `wall_time`(默认 30 分钟,自任务启动起算)约束:等待审批或用户输入时墙钟会暂停(每个提示有 24 小时的兜底上限),其余情况下墙钟优先于在途子代理——任务后段启动的满额 1800 秒工具仍可能被任务截止时间中断。在任务内部把 `execute_timeout` 或工具预算调到超过剩余任务墙钟不会生效。 ## 生命周期 From 424e695642831b492d2d19e33fb5830f06a6d388 Mon Sep 17 00:00:00 2001 From: asto Date: Sat, 26 Sep 2026 14:29:58 +0800 Subject: [PATCH 2/5] feat(tasks): pause the wall and idle clocks for human waits, answerably MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A durable task whose turn hit an approval or user-input prompt was still bounded by the 30-minute task wall_time, which interrupted the turn at expiry and denied the pending decision on the human's behalf. The window is now excluded from the wall clock (like the engine-side turn budget), the idle clock pauses and restarts fresh when the window closes, and the exclusion is bounded by a fail-safe cap. Three design holes close the gaps the first landing shipped with: - Answerability: the pause applies only where a host can actually deliver the decision. The runtime API's task threads are served by HTTP decide_approval/submit_user_input and pause; the TUI's private task runtime has no such channel, so a prompt there never arms the window and the wall clock releases the worker at its own deadline (before, it parked the worker for the 24h cap). - Aggregate bound: the cap now bounds the windows summed across the whole task, not only each open window — a task chaining short prompts can no longer sum its way past the cap once per prompt. - Park-poll backoff: while a window is parked the journal poll backs off from the 200ms catch-up cadence, doubling up to 2s; an arriving decision broadcasts on the subscription and wakes the loop early, so the backoff only removes the empty polls. Pinned by red-green tests: an unanswerable surface keeps the wall clock running (mutation check fails without the gate), chained short waits sum to the cap and terminalize between windows (fails without the hoisted aggregate check), and every wait-writer surface that arms the window still completes on a late answer. Signed-off-by: asto --- crates/tui/src/runtime_api.rs | 10 +- crates/tui/src/task_manager.rs | 209 +++++++++++++++++++++++++++++++-- docs/SUBAGENTS.md | 13 +- docs/zh_hans/SUBAGENTS.md | 2 +- 4 files changed, 217 insertions(+), 17 deletions(-) diff --git a/crates/tui/src/runtime_api.rs b/crates/tui/src/runtime_api.rs index c6a82afbed..e27fe7c2ce 100644 --- a/crates/tui/src/runtime_api.rs +++ b/crates/tui/src/runtime_api.rs @@ -855,9 +855,13 @@ pub async fn run_http_server( RuntimeThreadManagerConfig::from_task_data_dir(task_cfg.data_dir.clone()), plugin_discovery.registry_for_workspace(&workspace), )?; - let task_manager = - TaskManager::start_with_runtime_manager(task_cfg, config.clone(), runtime_threads.clone()) - .await?; + let task_manager = TaskManager::start_with_runtime_manager( + task_cfg, + config.clone(), + runtime_threads.clone(), + true, + ) + .await?; let automations = Arc::new(Mutex::new(AutomationManager::default_location()?)); runtime_threads.attach_automation_manager(automations.clone()); let scheduler_cancel = CancellationToken::new(); diff --git a/crates/tui/src/task_manager.rs b/crates/tui/src/task_manager.rs index 708054b632..a93c9023b2 100644 --- a/crates/tui/src/task_manager.rs +++ b/crates/tui/src/task_manager.rs @@ -38,6 +38,12 @@ const ARTIFACT_THRESHOLD: usize = 1200; const TASK_EVENT_CHANNEL_CAPACITY: usize = 256; const EVENT_CURSOR_BATCH: usize = 256; const EVENT_CATCHUP_POLL: Duration = Duration::from_millis(200); + +/// Ceiling for the park-poll backoff: while a prompt waits on a human the +/// turn loop backs its journal poll off from the catch-up cadence up to +/// this bound; an arriving decision wakes it through the subscription, so +/// this only bounds the empty polls. +const PARKED_POLL_CAP: Duration = Duration::from_secs(2); // v3 adds a provider pin; older executors must not silently ignore it. const CURRENT_TASK_SCHEMA_VERSION: u32 = 3; // Pinvou v0.9.0 persisted an additive v4 schema. Its additional fields are @@ -610,12 +616,22 @@ impl ExecutionGuard { // is expected, not idleness — and the idle budget restarts when the // window closes (see `end_human_wait`). let idle_elapsed = now.saturating_duration_since(self.last_progress_at); + // The aggregate bound: each window may close under the cap, but a + // task chaining prompts must not sum its way past the cap once per + // prompt. Checked whether or not a window is open right now, so the + // accumulated total terminalizes the task even between windows. + if self.human_wait_total >= self.limits.human_wait_cap { + return GuardAction::Interrupt { + reason: TaskTerminalReason::HumanWaitTimeout, + }; + } let pending = if shutdown { Some(TaskTerminalReason::Shutdown) } else if cancel { Some(TaskTerminalReason::Canceled) } else if let Some(since) = self.human_wait_since { - if now.saturating_duration_since(since) >= self.limits.human_wait_cap { + let parked = now.saturating_duration_since(since); + if self.human_wait_total.saturating_add(parked) >= self.limits.human_wait_cap { Some(TaskTerminalReason::HumanWaitTimeout) } else { None @@ -636,11 +652,20 @@ impl ExecutionGuard { .wall_time .saturating_sub(wall_elapsed) .min(self.limits.idle_progress.saturating_sub(idle_elapsed)); - if self.human_wait_since.is_some() { - // Both clocks are paused; wake at the fail-safe cadence. - wait = self.limits.human_wait_cap.saturating_sub( - now.saturating_duration_since(self.human_wait_since.unwrap_or(now)), - ); + if let Some(since) = self.human_wait_since { + // Both clocks are paused; the only deadline left is the + // fail-safe cap. A parked window can last hours, so polling + // the journal at the catch-up cadence for that whole span is + // waste: back the poll off, doubling with each parked second + // up to PARKED_POLL_CAP. An arriving decision broadcasts on + // the subscription and wakes the loop early, so the backoff + // never delays an answer, only the empty poll. + let parked = now.saturating_duration_since(since); + let cap_wait = self.limits.human_wait_cap.saturating_sub(parked); + let backoff = EVENT_CATCHUP_POLL + .saturating_mul(1u32 << parked.as_secs().min(4)) + .min(PARKED_POLL_CAP); + wait = cap_wait.min(backoff); } let wait = wait.min(EVENT_CATCHUP_POLL); GuardAction::Run { @@ -893,14 +918,26 @@ pub trait TaskExecutor: Send + Sync { pub struct EngineTaskExecutor { runtime_threads: SharedRuntimeThreadManager, limits: TaskExecutionLimits, + /// Whether a host exists that can actually deliver an approval or + /// user-input decision to this runtime surface. Only the runtime + /// API's HTTP delivery reaches a task thread's engine turn; the + /// TUI's private task runtime has no such host, so a prompt there + /// can never be answered, and arming the human-wait window for it + /// would park the worker on the fail-safe cap instead of the wall. + answerable_waits: bool, } impl EngineTaskExecutor { #[must_use] - pub fn new(runtime_threads: SharedRuntimeThreadManager, limits: TaskExecutionLimits) -> Self { + pub fn new( + runtime_threads: SharedRuntimeThreadManager, + limits: TaskExecutionLimits, + answerable_waits: bool, + ) -> Self { Self { runtime_threads, limits, + answerable_waits, } } } @@ -980,11 +1017,18 @@ impl TaskExecutor for EngineTaskExecutor { events, cancel, self.limits, + self.answerable_waits, ) .await } } +/// Drive one engine turn to its terminal state under the supervisor guard. +/// +/// `answerable_waits` carries [`EngineTaskExecutor`]'s answerability +/// decision: with no host able to deliver a decision, a pending prompt is +/// not a human-paced wait but a prompt nobody will answer, so the wall and +/// idle clocks keep running and release the worker at their own deadlines. async fn drive_engine_turn( runtime_threads: &RuntimeThreadManager, thread_id: &str, @@ -992,6 +1036,7 @@ async fn drive_engine_turn( events: mpsc::Sender, cancel: CancellationToken, limits: TaskExecutionLimits, + answerable_waits: bool, ) -> TaskExecutionResult { let mut subscription = runtime_threads.subscribe_events(); let mut guard = ExecutionGuard::new(limits, Instant::now()); @@ -1073,7 +1118,7 @@ async fn drive_engine_turn( // timeout, or the cancellation), so the window cannot stick // open; the guard's fail-safe cap bounds the pathological // case regardless. - "approval.required" | "user_input.required" => { + "approval.required" | "user_input.required" if answerable_waits => { guard.begin_human_wait(Instant::now()); emit_task_event(&events, TaskExecutionEvent::HumanWaitStarted).await; } @@ -1403,18 +1448,26 @@ impl TaskManager { RuntimeThreadManagerConfig::for_session(cfg.data_dir.clone(), session_id), plugin_registry, )?); - Self::start_with_runtime_manager(cfg, api_config, runtime_threads).await + Self::start_with_runtime_manager(cfg, api_config, runtime_threads, false).await } /// Start the manager with an injected runtime thread manager. + /// + /// `human_waits_answerable` is true only when the given manager is + /// served by the runtime API, whose HTTP delivery is the one channel + /// that can answer a task thread's prompt. The TUI's `start` builds a + /// private manager and passes false; a task that prompts there must + /// be released by the wall clock, not parked on the human-wait window. pub async fn start_with_runtime_manager( cfg: TaskManagerConfig, _api_config: Config, runtime_threads: SharedRuntimeThreadManager, + human_waits_answerable: bool, ) -> Result { let executor: Arc = Arc::new(EngineTaskExecutor::new( runtime_threads.clone(), cfg.execution_limits, + human_waits_answerable, )); let manager = Self::start_with_executor(cfg, executor).await?; runtime_threads.attach_task_manager(manager.clone()); @@ -4907,6 +4960,7 @@ mod tests { tx, CancellationToken::new(), TaskExecutionLimits::short_for_tests(), + false, ) .await; assert_eq!(result.status, TaskStatus::Failed); @@ -4937,6 +4991,7 @@ mod tests { persist_debounce: Duration::from_millis(10), human_wait_cap: Duration::from_secs(24 * 60 * 60), }, + false, ) .await }); @@ -4994,6 +5049,7 @@ mod tests { persist_debounce: Duration::from_millis(10), human_wait_cap: Duration::from_secs(24 * 60 * 60), }, + false, ) .await }); @@ -5096,6 +5152,7 @@ mod tests { persist_debounce: Duration::from_millis(10), human_wait_cap: Duration::from_secs(24 * 60 * 60), }, + true, ) .await }); @@ -5198,6 +5255,7 @@ mod tests { persist_debounce: Duration::from_millis(10), human_wait_cap: Duration::from_millis(200), }, + true, ) .await }); @@ -5271,6 +5329,135 @@ mod tests { Ok(()) } + #[tokio::test] + async fn unanswerable_surface_keeps_the_wall_clock_running_through_a_prompt() -> Result<()> { + // The TUI regression this feature shipped with: the TUI's private + // task runtime has no host that can deliver an approval or + // user-input decision, so a prompt there would park the worker for + // the fail-safe cap (24h by default) instead of releasing it at + // the 30-minute wall. On an unanswerable surface the prompt must + // not arm the human-wait window at all: the wall clock keeps + // running and interrupts the turn at its own deadline. + let runtime = Arc::new(test_runtime_manager().await?); + let thread = runtime + .create_thread(CreateThreadRequest::default()) + .await?; + let (tx, mut rx) = mpsc::channel(64); + let runtime_for_drive = Arc::clone(&runtime); + let thread_id = thread.id.clone(); + let drive = tokio::spawn(async move { + drive_engine_turn( + runtime_for_drive.as_ref(), + &thread_id, + "turn_unanswerable", + tx, + CancellationToken::new(), + TaskExecutionLimits { + wall_time: Duration::from_millis(300), + idle_progress: Duration::from_secs(30), + cancel_grace: Duration::from_millis(500), + persist_debounce: Duration::from_millis(10), + // The fail-safe cap is the default: if the prompt were + // (wrongly) treated as an answerable wait, the task + // would survive the wall and end much later with + // HumanWaitTimeout instead of WallTimeout. + human_wait_cap: Duration::from_secs(24 * 60 * 60), + }, + false, + ) + .await + }); + + runtime + .emit_event_for_test( + &thread.id, + Some("turn_unanswerable"), + "approval.required", + json!({ "id": "ap_1", "tool_name": "exec_command" }), + ) + .await?; + + let result = tokio::time::timeout(Duration::from_secs(5), drive) + .await + .expect("the wall clock must release the worker")?; + assert_eq!(result.terminal_reason, TaskTerminalReason::WallTimeout); + let mut saw_window_signal = false; + while let Ok(event) = rx.try_recv() { + if matches!(event, TaskExecutionEvent::HumanWaitStarted) { + saw_window_signal = true; + } + } + assert!( + !saw_window_signal, + "an unanswerable surface must not publish a human-wait window" + ); + Ok(()) + } + + #[tokio::test] + async fn chained_short_waits_sum_to_the_fail_safe_cap() -> Result<()> { + // Each window is well under the cap, so the per-window check alone + // would let a task chain prompts forever. The aggregate bound must + // release the worker once the closed windows alone reach the cap. + let runtime = Arc::new(test_runtime_manager().await?); + let thread = runtime + .create_thread(CreateThreadRequest::default()) + .await?; + let (tx, _rx) = mpsc::channel(64); + let runtime_for_drive = Arc::clone(&runtime); + let thread_id = thread.id.clone(); + let drive = tokio::spawn(async move { + drive_engine_turn( + runtime_for_drive.as_ref(), + &thread_id, + "turn_chained", + tx, + CancellationToken::new(), + TaskExecutionLimits { + wall_time: Duration::from_secs(30), + idle_progress: Duration::from_secs(30), + cancel_grace: Duration::from_millis(50), + persist_debounce: Duration::from_millis(10), + // Two 120ms windows sum past this; neither alone does. + human_wait_cap: Duration::from_millis(200), + }, + true, + ) + .await + }); + + for round in 0..2 { + runtime + .emit_event_for_test( + &thread.id, + Some("turn_chained"), + "user_input.required", + json!({ "id": format!("ui_{round}"), "request": {} }), + ) + .await?; + tokio::time::sleep(Duration::from_millis(120)).await; + runtime + .emit_event_for_test( + &thread.id, + Some("turn_chained"), + "user_input.answered", + json!({ "id": format!("ui_{round}") }), + ) + .await?; + tokio::time::sleep(Duration::from_millis(50)).await; + } + + let result = tokio::time::timeout(Duration::from_secs(5), drive) + .await + .expect("the aggregate bound must release the worker")?; + assert_eq!( + result.terminal_reason, + TaskTerminalReason::HumanWaitTimeout, + "closed windows past the cap must terminalize, not park forever" + ); + Ok(()) + } + #[tokio::test] async fn silent_running_tool_heartbeats_the_worker_watchdog() -> Result<()> { // The supervisor's guard (run_task) only sees task events, so the @@ -5298,6 +5485,7 @@ mod tests { persist_debounce: Duration::from_millis(10), human_wait_cap: Duration::from_secs(24 * 60 * 60), }, + false, ) .await }); @@ -5374,6 +5562,7 @@ mod tests { persist_debounce: Duration::from_millis(10), human_wait_cap: Duration::from_secs(24 * 60 * 60), }, + false, ) .await }); @@ -5445,6 +5634,7 @@ mod tests { tx, CancellationToken::new(), TaskExecutionLimits::short_for_tests(), + false, ) .await; assert_eq!(result.status, TaskStatus::Completed); @@ -5477,6 +5667,7 @@ mod tests { tx, cancel, TaskExecutionLimits::short_for_tests(), + false, ) .await; assert_eq!(result.status, TaskStatus::Canceled); diff --git a/docs/SUBAGENTS.md b/docs/SUBAGENTS.md index 6f1af98b27..fdd5a357d9 100644 --- a/docs/SUBAGENTS.md +++ b/docs/SUBAGENTS.md @@ -599,10 +599,15 @@ long in-flight tool is cancelled before its own timeout can fire. Inside a durable task these budgets are additionally bounded by the task's `wall_time` (default 30 minutes, measured from task start): the wall clock pauses while a pending approval or user-input prompt waits on a human (with a -24-hour fail-safe cap per prompt), but it otherwise wins over an in-flight -sub-agent, so a maximal 1800-second tool started late in a task can still be -interrupted by the task deadline. Raising `execute_timeout` or the tool -budget above the remaining task wall has no effect inside a task. +24-hour fail-safe cap that bounds each open window and the windows summed +across the whole task), but it otherwise wins over an in-flight sub-agent, so +a maximal 1800-second tool started late in a task can still be interrupted by +the task deadline. Raising `execute_timeout` or the tool budget above the +remaining task wall has no effect inside a task. The pause applies only where +a host can actually deliver the decision: the runtime API serves task threads +with HTTP `decide_approval`/`submit_user_input`, so its tasks pause; the +TUI's private task runtime has no such channel, so there a prompt is not +treated as a human-paced wait at all and the wall clock keeps running. ## Lifecycle diff --git a/docs/zh_hans/SUBAGENTS.md b/docs/zh_hans/SUBAGENTS.md index e698f1f1ba..40ab13e0fe 100644 --- a/docs/zh_hans/SUBAGENTS.md +++ b/docs/zh_hans/SUBAGENTS.md @@ -322,7 +322,7 @@ heartbeat_timeout_secs = 300 # 钳制到 30..=3600 有效心跳至少保持在解析后的 `api_timeout_secs` 与内置子代理工具超时之上各 30 秒(子代理只在步骤边界记录进度,不会在工具执行中途记录,因此静默的构建或 MCP 调用不能在运行中被清理),所以配置的长模型请求与长时间在途工具都不会在自己的超时触发之前被取消。 -在持久任务内部,这些预算还受任务 `wall_time`(默认 30 分钟,自任务启动起算)约束:等待审批或用户输入时墙钟会暂停(每个提示有 24 小时的兜底上限),其余情况下墙钟优先于在途子代理——任务后段启动的满额 1800 秒工具仍可能被任务截止时间中断。在任务内部把 `execute_timeout` 或工具预算调到超过剩余任务墙钟不会生效。 +在持久任务内部,这些预算还受任务 `wall_time`(默认 30 分钟,自任务启动起算)约束:等待审批或用户输入时墙钟会暂停(24 小时兜底上限既约束每个打开的窗口,也约束整个任务内各窗口之和),其余情况下墙钟优先于在途子代理——任务后段启动的满额 1800 秒工具仍可能被任务截止时间中断。在任务内部把 `execute_timeout` 或工具预算调到超过剩余任务墙钟不会生效。暂停只在确有宿主能送达决策的场合生效:runtime API 的任务线程可通过 HTTP `decide_approval`/`submit_user_input` 应答,其任务会暂停;TUI 的私有任务运行时没有这条通道,那里的提示不会被视为一个人工等待,墙钟继续运行。 ## 生命周期 From 77c09b194c0e2e7ec10370ec156a7996bf8bf5d2 Mon Sep 17 00:00:00 2001 From: asto Date: Mon, 28 Sep 2026 01:58:41 +0800 Subject: [PATCH 3/5] fix(tasks): make the parked poll backoff real and the answerability opt-in explicit The park-poll backoff shipped inert: evaluate() clamped the loop wait to the 200ms catch-up cadence unconditionally, so the doubling and PARKED_POLL_CAP never took effect and the PR's own comments described a behavior that did not exist. The clamp now applies only on the non-parked path, and the parked cap wait excludes the already-closed windows so the aggregate budget, not the open window alone, bounds the sleep. The aggregate HumanWaitTimeout check moved into the pending chain after shutdown and cancel (it previously jumped the queue and mislabeled a canceled or shutting-down task as human_wait_timeout), and the window- closing event arm is now gated like the opening arm, so unanswerable surfaces no longer emit unpaired HumanWaitEnded signals. The answerability flag moved from a positional bool parameter to the named TaskManagerConfig::human_waits_answerable field, defaulting to false (the TUI's safe default); the runtime API opts in through server_task_config, which carries the contract in its name and doc. New deterministic pins cover the backoff cadence and its 2s ceiling, the aggregate-budget wake, the between-windows terminalization, cancel and shutdown outranking a spent cap, the idle restart on window close, the from_runtime default, and the server opt-in. Signed-off-by: asto --- crates/tui/src/automation_manager.rs | 1 + crates/tui/src/runtime_api.rs | 41 +++-- crates/tui/src/runtime_api/tests.rs | 14 ++ crates/tui/src/task_manager.rs | 232 +++++++++++++++++++++++---- docs/SUBAGENTS.md | 3 + docs/zh_hans/SUBAGENTS.md | 2 +- 6 files changed, 247 insertions(+), 46 deletions(-) diff --git a/crates/tui/src/automation_manager.rs b/crates/tui/src/automation_manager.rs index c48233634b..daaf5520e4 100644 --- a/crates/tui/src/automation_manager.rs +++ b/crates/tui/src/automation_manager.rs @@ -2250,6 +2250,7 @@ mod tests { allow_shell: true, trust_mode: true, execution_limits: crate::task_manager::TaskExecutionLimits::default(), + human_waits_answerable: false, } } diff --git a/crates/tui/src/runtime_api.rs b/crates/tui/src/runtime_api.rs index e27fe7c2ce..f07492dd53 100644 --- a/crates/tui/src/runtime_api.rs +++ b/crates/tui/src/runtime_api.rs @@ -833,6 +833,29 @@ fn open_runtime_threads_for_server( Ok((manager, workshop_activation)) } +/// Build the task manager config for the runtime API server. +/// +/// This is where the server opts its tasks into answerable human waits: +/// the tasks it runs are served by this server's HTTP decision endpoints +/// (`/v1/approvals/`, `/v1/user-input/`), so a pending prompt on +/// them is a human-paced wait the wall clock may pause for — see +/// [`TaskManagerConfig::human_waits_answerable`]. +#[must_use] +pub fn server_task_config( + config: &Config, + workspace: PathBuf, + workers: usize, +) -> TaskManagerConfig { + let mut task_cfg = TaskManagerConfig::from_runtime( + config, + workspace, + Some(config.default_model()), + Some(workers), + ); + task_cfg.human_waits_answerable = true; + task_cfg +} + /// Start the runtime API server. pub async fn run_http_server( config: Config, @@ -842,26 +865,16 @@ pub async fn run_http_server( ) -> Result<()> { validate_runtime_listener_security(&options)?; - let task_default_model = config.default_model(); - let task_cfg = TaskManagerConfig::from_runtime( - &config, - workspace.clone(), - Some(task_default_model), - Some(options.workers), - ); + let task_cfg = server_task_config(&config, workspace.clone(), options.workers); let (runtime_threads, _workshop_activation) = open_runtime_threads_for_server( &config, workspace.clone(), RuntimeThreadManagerConfig::from_task_data_dir(task_cfg.data_dir.clone()), plugin_discovery.registry_for_workspace(&workspace), )?; - let task_manager = TaskManager::start_with_runtime_manager( - task_cfg, - config.clone(), - runtime_threads.clone(), - true, - ) - .await?; + let task_manager = + TaskManager::start_with_runtime_manager(task_cfg, config.clone(), runtime_threads.clone()) + .await?; let automations = Arc::new(Mutex::new(AutomationManager::default_location()?)); runtime_threads.attach_automation_manager(automations.clone()); let scheduler_cancel = CancellationToken::new(); diff --git a/crates/tui/src/runtime_api/tests.rs b/crates/tui/src/runtime_api/tests.rs index 2dd3549ebc..b3d998cb3a 100644 --- a/crates/tui/src/runtime_api/tests.rs +++ b/crates/tui/src/runtime_api/tests.rs @@ -51,6 +51,19 @@ fn runtime_session_fallback_retains_non_unicode_explicit_home_boundary() { assert_eq!(fallback_sessions_dir(), explicit.join("sessions")); } +#[test] +fn server_task_config_opts_tasks_into_answerable_human_waits() { + // The HTTP server is the one host that can deliver a decision to its + // tasks (`/v1/approvals/`, `/v1/user-input/`), so its config must opt + // in explicitly; everywhere else (`from_runtime`, the TUI's private + // runtime) the safe default is unanswerable. + let cfg = server_task_config(&Config::default(), PathBuf::from("."), 2); + assert!( + cfg.human_waits_answerable, + "server task config must arm answerable human waits" + ); +} + #[test] fn thread_route_credential_error_is_bad_request_not_not_found() { let credential = map_thread_err(anyhow::anyhow!("DeepSeek API key not found")); @@ -969,6 +982,7 @@ async fn spawn_test_server_with_root_token_mobile_workspace_and_overrides( allow_shell: false, trust_mode: false, execution_limits: crate::task_manager::TaskExecutionLimits::default(), + human_waits_answerable: false, }, Arc::new(MockExecutor), ) diff --git a/crates/tui/src/task_manager.rs b/crates/tui/src/task_manager.rs index a93c9023b2..db000e5bd7 100644 --- a/crates/tui/src/task_manager.rs +++ b/crates/tui/src/task_manager.rs @@ -468,6 +468,18 @@ pub struct TaskManagerConfig { pub allow_shell: bool, pub trust_mode: bool, pub execution_limits: TaskExecutionLimits, + /// Whether a host is attached that can actually deliver an approval or + /// user-input decision to tasks run by this manager — today only the + /// runtime API, whose HTTP `decide_approval`/`submit_user_input` + /// endpoints reach a task thread's engine turn. A pending prompt on an + /// answerable manager is a human-paced wait: the wall and idle clocks + /// pause for it (bounded by `human_wait_cap`). A prompt on an + /// unanswerable manager — the TUI's private task runtime has no + /// delivery channel — can never be answered, so arming the window + /// there would park the worker on the fail-safe cap; the clocks keep + /// running instead. Defaults to `false`: opt in only where the + /// delivery channel exists. + pub human_waits_answerable: bool, } /// Deadlines and persistence cadence for one durable execution. @@ -477,10 +489,11 @@ pub struct TaskExecutionLimits { pub idle_progress: Duration, pub cancel_grace: Duration, pub persist_debounce: Duration, - /// Fail-safe ceiling for one continuous human wait (a pending approval - /// or user-input prompt). Human-paced waits are excluded from - /// `wall_time` — a slow answer must not end the task — but a prompt - /// nobody will ever answer must still release the worker eventually. + /// Fail-safe ceiling for the human-wait windows of one task, applied + /// both to each open window and to the windows summed across the whole + /// execution. Human-paced waits are excluded from `wall_time` — a slow + /// answer must not end the task — but prompts that nobody will ever + /// answer must still release the worker eventually. pub human_wait_cap: Duration, } @@ -616,15 +629,6 @@ impl ExecutionGuard { // is expected, not idleness — and the idle budget restarts when the // window closes (see `end_human_wait`). let idle_elapsed = now.saturating_duration_since(self.last_progress_at); - // The aggregate bound: each window may close under the cap, but a - // task chaining prompts must not sum its way past the cap once per - // prompt. Checked whether or not a window is open right now, so the - // accumulated total terminalizes the task even between windows. - if self.human_wait_total >= self.limits.human_wait_cap { - return GuardAction::Interrupt { - reason: TaskTerminalReason::HumanWaitTimeout, - }; - } let pending = if shutdown { Some(TaskTerminalReason::Shutdown) } else if cancel { @@ -636,6 +640,13 @@ impl ExecutionGuard { } else { None } + } else if self.human_wait_total >= self.limits.human_wait_cap { + // The aggregate bound: each window may close under the cap, but + // a task chaining prompts must not sum its way past the cap once + // per prompt. With no window open right now, the closed windows + // alone can cross it, so the accumulated total terminalizes the + // task even between windows. + Some(TaskTerminalReason::HumanWaitTimeout) } else if wall_elapsed >= self.limits.wall_time { Some(TaskTerminalReason::WallTimeout) } else if idle_elapsed >= self.limits.idle_progress { @@ -659,15 +670,23 @@ impl ExecutionGuard { // waste: back the poll off, doubling with each parked second // up to PARKED_POLL_CAP. An arriving decision broadcasts on // the subscription and wakes the loop early, so the backoff - // never delays an answer, only the empty poll. + // never delays an answer — it only thins the empty polls, and + // only while the event bus is quiet. The cap wait excludes the + // already-closed windows too: the aggregate budget, not this + // window alone, bounds the sleep. let parked = now.saturating_duration_since(since); - let cap_wait = self.limits.human_wait_cap.saturating_sub(parked); + let cap_wait = self + .limits + .human_wait_cap + .saturating_sub(self.human_wait_total) + .saturating_sub(parked); let backoff = EVENT_CATCHUP_POLL .saturating_mul(1u32 << parked.as_secs().min(4)) .min(PARKED_POLL_CAP); wait = cap_wait.min(backoff); + } else { + wait = wait.min(EVENT_CATCHUP_POLL); } - let wait = wait.min(EVENT_CATCHUP_POLL); GuardAction::Run { wait: wait.max(Duration::from_millis(1)), } @@ -705,6 +724,7 @@ impl TaskManagerConfig { allow_shell: config.allow_shell(), trust_mode: false, execution_limits: TaskExecutionLimits::default(), + human_waits_answerable: false, } } } @@ -1112,12 +1132,14 @@ async fn drive_engine_turn( // A pending approval or user-input prompt is a human-paced // wait: exclude it from the wall clock exactly like the // engine-side turn budget does, and publish the window so - // the worker supervisor's guard follows in lockstep. Every - // wait exit emits one of the paired closing events (the - // decision, the interrupt deny, the legacy replayed - // timeout, or the cancellation), so the window cannot stick - // open; the guard's fail-safe cap bounds the pathological - // case regardless. + // the worker supervisor's guard follows in lockstep. Both + // arms are gated on `answerable_waits`, so the window pair + // only flows where a decision can actually be delivered. + // Every wait exit emits one of the paired closing events + // (the decision, the interrupt deny, the legacy replayed + // timeout, or the cancellation), so the window cannot + // stick open; the guard's fail-safe cap bounds the + // pathological case regardless. "approval.required" | "user_input.required" if answerable_waits => { guard.begin_human_wait(Instant::now()); emit_task_event(&events, TaskExecutionEvent::HumanWaitStarted).await; @@ -1125,7 +1147,9 @@ async fn drive_engine_turn( "approval.decided" | "approval.timeout" | "user_input.answered" - | "user_input.canceled" => { + | "user_input.canceled" + if answerable_waits => + { guard.end_human_wait(Instant::now()); emit_task_event(&events, TaskExecutionEvent::HumanWaitEnded).await; } @@ -1448,26 +1472,23 @@ impl TaskManager { RuntimeThreadManagerConfig::for_session(cfg.data_dir.clone(), session_id), plugin_registry, )?); - Self::start_with_runtime_manager(cfg, api_config, runtime_threads, false).await + Self::start_with_runtime_manager(cfg, api_config, runtime_threads).await } /// Start the manager with an injected runtime thread manager. /// - /// `human_waits_answerable` is true only when the given manager is - /// served by the runtime API, whose HTTP delivery is the one channel - /// that can answer a task thread's prompt. The TUI's `start` builds a - /// private manager and passes false; a task that prompts there must - /// be released by the wall clock, not parked on the human-wait window. + /// Whether prompts on this manager's tasks are answerable human-paced + /// waits is `cfg.human_waits_answerable` — the TUI's `start` leaves it + /// `false`, the runtime API opts in. See that field for the contract. pub async fn start_with_runtime_manager( cfg: TaskManagerConfig, _api_config: Config, runtime_threads: SharedRuntimeThreadManager, - human_waits_answerable: bool, ) -> Result { let executor: Arc = Arc::new(EngineTaskExecutor::new( runtime_threads.clone(), cfg.execution_limits, - human_waits_answerable, + cfg.human_waits_answerable, )); let manager = Self::start_with_executor(cfg, executor).await?; runtime_threads.attach_task_manager(manager.clone()); @@ -3251,6 +3272,7 @@ mod tests { allow_shell: false, trust_mode: false, execution_limits: TaskExecutionLimits::default(), + human_waits_answerable: false, } } @@ -4608,6 +4630,154 @@ mod tests { } } + /// Limits with generous wall/idle budgets so a test exercises exactly + /// one guard mechanism (the human-wait arithmetic) in isolation. + fn execution_guard_limits(human_wait_cap: Duration) -> TaskExecutionLimits { + TaskExecutionLimits { + wall_time: Duration::from_secs(30), + idle_progress: Duration::from_secs(30), + cancel_grace: Duration::from_millis(50), + persist_debounce: Duration::from_millis(10), + human_wait_cap, + } + } + + #[test] + fn execution_guard_parked_backoff_grows_past_the_catchup_poll() { + // Two seconds parked: the doubling has run twice, so the parked + // wait must exceed the 200ms catch-up cadence the backoff exists + // to escape. + let start = Instant::now(); + let mut guard = ExecutionGuard::new( + execution_guard_limits(Duration::from_secs(24 * 60 * 60)), + start, + ); + guard.begin_human_wait(start); + match guard.evaluate(start + Duration::from_secs(2), false, false) { + GuardAction::Run { wait } => { + assert_eq!(wait, EVENT_CATCHUP_POLL * 4); + } + other => panic!("expected a parked run, got {other:?}"), + } + } + + #[test] + fn execution_guard_parked_backoff_caps_at_two_seconds() { + let start = Instant::now(); + let mut guard = ExecutionGuard::new( + execution_guard_limits(Duration::from_secs(24 * 60 * 60)), + start, + ); + guard.begin_human_wait(start); + match guard.evaluate(start + Duration::from_secs(3600), false, false) { + GuardAction::Run { wait } => { + assert_eq!(wait, PARKED_POLL_CAP); + } + other => panic!("expected a parked run, got {other:?}"), + } + } + + #[test] + fn execution_guard_parked_wait_respects_the_aggregate_budget() { + // The parked wake must respect the aggregate budget, not just the + // open window: 10s cap, 7s already closed, 2.5s parked this window + // leaves 500ms — not the 800ms the doubling alone would suggest. + let start = Instant::now(); + let mut guard = ExecutionGuard::new(execution_guard_limits(Duration::from_secs(10)), start); + guard.begin_human_wait(start); + guard.end_human_wait(start + Duration::from_secs(7)); + guard.begin_human_wait(start + Duration::from_secs(7)); + match guard.evaluate(start + Duration::from_millis(9500), false, false) { + GuardAction::Run { wait } => { + assert_eq!(wait, Duration::from_millis(500)); + } + other => panic!("expected a parked run, got {other:?}"), + } + } + + #[test] + fn execution_guard_spent_aggregate_cap_terminalizes_between_windows() { + // Closed windows alone reach the cap with no window open: the + // guard must terminalize between prompts even though wall and idle + // have barely elapsed. + let start = Instant::now(); + let mut guard = + ExecutionGuard::new(execution_guard_limits(Duration::from_millis(200)), start); + guard.begin_human_wait(start); + guard.end_human_wait(start + Duration::from_millis(200)); + match guard.evaluate(start + Duration::from_millis(250), false, false) { + GuardAction::Interrupt { reason } => { + assert_eq!(reason, TaskTerminalReason::HumanWaitTimeout); + } + other => panic!("expected the aggregate cap to fire, got {other:?}"), + } + } + + #[test] + fn execution_guard_cancel_outranks_a_spent_aggregate_cap() { + let start = Instant::now(); + let mut guard = + ExecutionGuard::new(execution_guard_limits(Duration::from_millis(200)), start); + guard.begin_human_wait(start); + guard.end_human_wait(start + Duration::from_millis(200)); + match guard.evaluate(start + Duration::from_millis(250), true, false) { + GuardAction::Interrupt { reason } => { + assert_eq!(reason, TaskTerminalReason::Canceled); + } + other => panic!("cancel must outrank the spent cap, got {other:?}"), + } + } + + #[test] + fn execution_guard_shutdown_outranks_a_spent_aggregate_cap() { + let start = Instant::now(); + let mut guard = + ExecutionGuard::new(execution_guard_limits(Duration::from_millis(200)), start); + guard.begin_human_wait(start); + guard.end_human_wait(start + Duration::from_millis(200)); + match guard.evaluate(start + Duration::from_millis(250), false, true) { + GuardAction::Interrupt { reason } => { + assert_eq!(reason, TaskTerminalReason::Shutdown); + } + other => panic!("shutdown must outrank the spent cap, got {other:?}"), + } + } + + #[test] + fn execution_guard_human_wait_close_restarts_the_idle_budget() { + // Ending a window restarts idle from zero: 100ms after the close + // the guard must still run, though 1.1s have passed since start. + let start = Instant::now(); + let limits = TaskExecutionLimits { + wall_time: Duration::from_millis(400), + idle_progress: Duration::from_millis(150), + cancel_grace: Duration::from_millis(500), + persist_debounce: Duration::from_millis(10), + human_wait_cap: Duration::from_secs(24 * 60 * 60), + }; + let mut guard = ExecutionGuard::new(limits, start); + guard.begin_human_wait(start); + guard.end_human_wait(start + Duration::from_secs(1)); + match guard.evaluate(start + Duration::from_millis(1100), false, false) { + GuardAction::Run { .. } => {} + other => panic!("a closed window must restart idle, got {other:?}"), + } + } + + #[test] + fn from_runtime_defaults_to_unanswerable_human_waits() { + // The safe default: the TUI builds its private task runtime through + // `from_runtime` and must never arm the human-wait window — no host + // can answer a prompt there. Only the runtime API opts in, via + // `server_task_config`. + let cfg = + TaskManagerConfig::from_runtime(&Config::default(), PathBuf::from("."), None, None); + assert!( + !cfg.human_waits_answerable, + "from_runtime must default to unanswerable human waits" + ); + } + #[test] fn consecutive_message_deltas_coalesce_on_the_timeline() { let mut task = sample_task_record(); diff --git a/docs/SUBAGENTS.md b/docs/SUBAGENTS.md index fdd5a357d9..065a07c2a1 100644 --- a/docs/SUBAGENTS.md +++ b/docs/SUBAGENTS.md @@ -608,6 +608,9 @@ a host can actually deliver the decision: the runtime API serves task threads with HTTP `decide_approval`/`submit_user_input`, so its tasks pause; the TUI's private task runtime has no such channel, so there a prompt is not treated as a human-paced wait at all and the wall clock keeps running. +While parked, the task keeps holding its worker slot (two workers by +default), so several simultaneously parked tasks can stall the durable-task +queue until their answers arrive or the fail-safe cap fires. ## Lifecycle diff --git a/docs/zh_hans/SUBAGENTS.md b/docs/zh_hans/SUBAGENTS.md index 40ab13e0fe..780d8e1f31 100644 --- a/docs/zh_hans/SUBAGENTS.md +++ b/docs/zh_hans/SUBAGENTS.md @@ -322,7 +322,7 @@ heartbeat_timeout_secs = 300 # 钳制到 30..=3600 有效心跳至少保持在解析后的 `api_timeout_secs` 与内置子代理工具超时之上各 30 秒(子代理只在步骤边界记录进度,不会在工具执行中途记录,因此静默的构建或 MCP 调用不能在运行中被清理),所以配置的长模型请求与长时间在途工具都不会在自己的超时触发之前被取消。 -在持久任务内部,这些预算还受任务 `wall_time`(默认 30 分钟,自任务启动起算)约束:等待审批或用户输入时墙钟会暂停(24 小时兜底上限既约束每个打开的窗口,也约束整个任务内各窗口之和),其余情况下墙钟优先于在途子代理——任务后段启动的满额 1800 秒工具仍可能被任务截止时间中断。在任务内部把 `execute_timeout` 或工具预算调到超过剩余任务墙钟不会生效。暂停只在确有宿主能送达决策的场合生效:runtime API 的任务线程可通过 HTTP `decide_approval`/`submit_user_input` 应答,其任务会暂停;TUI 的私有任务运行时没有这条通道,那里的提示不会被视为一个人工等待,墙钟继续运行。 +在持久任务内部,这些预算还受任务 `wall_time`(默认 30 分钟,自任务启动起算)约束:等待审批或用户输入时墙钟会暂停(24 小时兜底上限既约束每个打开的窗口,也约束整个任务内各窗口之和),其余情况下墙钟优先于在途子代理——任务后段启动的满额 1800 秒工具仍可能被任务截止时间中断。在任务内部把 `execute_timeout` 或工具预算调到超过剩余任务墙钟不会生效。暂停只在确有宿主能送达决策的场合生效:runtime API 的任务线程可通过 HTTP `decide_approval`/`submit_user_input` 应答,其任务会暂停;TUI 的私有任务运行时没有这条通道,那里的提示不会被视为一个人工等待,墙钟继续运行。暂停期间任务仍占用其工作者槽位(默认两个工作者),因此多个同时暂停的任务可能让持久任务队列停摆,直到应答抵达或兜底上限触发。 ## 生命周期 From 1c145c94d5d19f5eed83db701b652e0872c5f6ec Mon Sep 17 00:00:00 2001 From: asto Date: Mon, 28 Sep 2026 10:06:49 +0800 Subject: [PATCH 4/5] test(tasks): pin the human-wait guards at their claimed mechanisms A fresh from-scratch re-review confirmed the guard behavior but caught the chained-waits integration test claiming the wrong mechanism: with back-to-back windows the open-window check fires mid-window (window one's spend is already in the aggregate total), so deleting the between-windows aggregate arm leaves the test green. The test is renamed and recommented to say what it actually pins; the between-windows arm itself stays pinned by execution_guard_spent_aggregate_cap_terminalizes_between_windows, which does go red under that mutation. Also from the review: - the backoff pins assert literal durations (800ms, 2s) instead of restating the constants, so drifting EVENT_CATCHUP_POLL or PARKED_POLL_CAP off their documented 200ms / 2s values fails them; - the supervised wall-exclusion executor delay grows from 900ms to 2s, clearing the thin 500ms scheduling-starvation margin on loaded CI; - the shared evaluate() comment now names both wake paths (the turn loop's event subscription and the supervisor's forwarded task event); - the human-wait event and config docs state that a custom executor emitting the HumanWaitStarted/Ended pair self-declares its waits answerable, independent of human_waits_answerable. Signed-off-by: asto --- crates/tui/src/task_manager.rs | 60 ++++++++++++++++++++++++---------- 1 file changed, 42 insertions(+), 18 deletions(-) diff --git a/crates/tui/src/task_manager.rs b/crates/tui/src/task_manager.rs index db000e5bd7..55fa6db485 100644 --- a/crates/tui/src/task_manager.rs +++ b/crates/tui/src/task_manager.rs @@ -478,7 +478,11 @@ pub struct TaskManagerConfig { /// delivery channel — can never be answered, so arming the window /// there would park the worker on the fail-safe cap; the clocks keep /// running instead. Defaults to `false`: opt in only where the - /// delivery channel exists. + /// delivery channel exists. The gate arms the built-in engine + /// executor's journal tracking; a custom executor that emits the + /// `HumanWaitStarted`/`HumanWaitEnded` pair is self-declaring its + /// waits answerable, and the supervisor honors the pair regardless of + /// this flag. pub human_waits_answerable: bool, } @@ -668,12 +672,14 @@ impl ExecutionGuard { // fail-safe cap. A parked window can last hours, so polling // the journal at the catch-up cadence for that whole span is // waste: back the poll off, doubling with each parked second - // up to PARKED_POLL_CAP. An arriving decision broadcasts on - // the subscription and wakes the loop early, so the backoff - // never delays an answer — it only thins the empty polls, and - // only while the event bus is quiet. The cap wait excludes the - // already-closed windows too: the aggregate budget, not this - // window alone, bounds the sleep. + // up to PARKED_POLL_CAP. An arriving decision wakes this loop + // early — the turn loop through the runtime event + // subscription, the worker supervisor through the forwarded + // task event — so the backoff never delays an answer; it only + // thins the empty polls, and only while the event bus is + // quiet. The cap wait excludes the already-closed windows too: + // the aggregate budget, not this window alone, bounds the + // sleep. let parked = now.saturating_duration_since(since); let cap_wait = self .limits @@ -870,7 +876,11 @@ pub enum TaskExecutionEvent { /// A human-paced wait (pending approval or user-input prompt) opened or /// closed, so the worker supervisor's wall clock can exclude it exactly /// like the turn loop's own guard does. Supervisor-side signal only: - /// never persisted and never shown on the task timeline. + /// never persisted and never shown on the task timeline. The built-in + /// engine executor emits the pair only where `human_waits_answerable` + /// armed the window; a custom executor that emits it is self-declaring + /// its waits answerable, and the supervisor honors the pair regardless + /// of that flag. HumanWaitStarted, HumanWaitEnded, ToolCompleted { @@ -4646,7 +4656,9 @@ mod tests { fn execution_guard_parked_backoff_grows_past_the_catchup_poll() { // Two seconds parked: the doubling has run twice, so the parked // wait must exceed the 200ms catch-up cadence the backoff exists - // to escape. + // to escape. Asserted as a literal so drifting EVENT_CATCHUP_POLL + // off its documented 200ms fails here instead of silently + // tracking the constant. let start = Instant::now(); let mut guard = ExecutionGuard::new( execution_guard_limits(Duration::from_secs(24 * 60 * 60)), @@ -4655,7 +4667,7 @@ mod tests { guard.begin_human_wait(start); match guard.evaluate(start + Duration::from_secs(2), false, false) { GuardAction::Run { wait } => { - assert_eq!(wait, EVENT_CATCHUP_POLL * 4); + assert_eq!(wait, Duration::from_millis(800)); } other => panic!("expected a parked run, got {other:?}"), } @@ -4663,6 +4675,9 @@ mod tests { #[test] fn execution_guard_parked_backoff_caps_at_two_seconds() { + // Asserted as a literal so moving PARKED_POLL_CAP off its + // documented 2s ceiling fails here instead of silently tracking + // the constant. let start = Instant::now(); let mut guard = ExecutionGuard::new( execution_guard_limits(Duration::from_secs(24 * 60 * 60)), @@ -4671,7 +4686,7 @@ mod tests { guard.begin_human_wait(start); match guard.evaluate(start + Duration::from_secs(3600), false, false) { GuardAction::Run { wait } => { - assert_eq!(wait, PARKED_POLL_CAP); + assert_eq!(wait, Duration::from_secs(2)); } other => panic!("expected a parked run, got {other:?}"), } @@ -5486,7 +5501,11 @@ mod tests { let manager = TaskManager::start_with_executor( config, Arc::new(SupervisedHumanWaitExecutor { - delay: Duration::from_millis(900), + // Two seconds: far past the short wall, with enough margin + // that scheduler starvation on a loaded CI runner cannot + // eat the exclusion before the supervisor sees the + // HumanWaitStarted signal. + delay: Duration::from_secs(2), }), ) .await?; @@ -5565,10 +5584,15 @@ mod tests { } #[tokio::test] - async fn chained_short_waits_sum_to_the_fail_safe_cap() -> Result<()> { - // Each window is well under the cap, so the per-window check alone - // would let a task chain prompts forever. The aggregate bound must - // release the worker once the closed windows alone reach the cap. + async fn chained_short_waits_still_terminate_at_the_fail_safe_cap() -> Result<()> { + // Each window is well under the cap, so chaining prompts must not + // park the worker forever: once the windows sum past the cap the + // guard terminalizes. Which arm observes the crossing first is a + // wake-timing detail this test does not pin — with these + // back-to-back windows the open-window check fires mid-window + // (window one's spend is already in the aggregate total); the + // between-windows arm alone is pinned deterministically by + // execution_guard_spent_aggregate_cap_terminalizes_between_windows. let runtime = Arc::new(test_runtime_manager().await?); let thread = runtime .create_thread(CreateThreadRequest::default()) @@ -5619,11 +5643,11 @@ mod tests { let result = tokio::time::timeout(Duration::from_secs(5), drive) .await - .expect("the aggregate bound must release the worker")?; + .expect("the fail-safe cap must release the worker")?; assert_eq!( result.terminal_reason, TaskTerminalReason::HumanWaitTimeout, - "closed windows past the cap must terminalize, not park forever" + "chained sub-cap windows must not park the worker forever" ); Ok(()) } From 5ae3511fdf9e031d135e04952b013e1bfb057d3b Mon Sep 17 00:00:00 2001 From: asto Date: Fri, 2 Oct 2026 13:55:43 +0800 Subject: [PATCH 5/5] test(tasks): pin the server answerability wiring end to end MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The fresh review's MAJOR: the server opt-in had no test below the config builder, so reverting run_http_server's call site or dropping the flag's forward into EngineTaskExecutor kept the whole suite green — the one mutation the PR's own red-green standard did not cover. run_http_server's task-runtime wiring moves into open_server_task_runtime, the single wiring point the new pin also drives. The pin reads the flag where it lands (the executor field the manager runs tasks on) through a cfg(test) as_engine_executor seam, and goes red if any of the three hops (builder, forward, call site) regresses to the unanswerable default. Verified by mutation: inlining the builder in open_server_task_runtime and hardcoding the forward to false each turn only this pin red (the builder pin stays green — it never saw those hops); dropping the builder's flag assignment turns both red. Signed-off-by: asto --- crates/tui/src/runtime_api.rs | 41 ++++++++++++++++++++++------ crates/tui/src/task_manager.rs | 49 ++++++++++++++++++++++++++++++++++ 2 files changed, 82 insertions(+), 8 deletions(-) diff --git a/crates/tui/src/runtime_api.rs b/crates/tui/src/runtime_api.rs index f07492dd53..7a38b32ab7 100644 --- a/crates/tui/src/runtime_api.rs +++ b/crates/tui/src/runtime_api.rs @@ -856,6 +856,34 @@ pub fn server_task_config( task_cfg } +/// Open the runtime threads and the task manager the server serves, wiring +/// the server's answerable human-wait opt-in ([`server_task_config`]) into +/// the task executor. +/// +/// This is the single wiring point for that opt-in: `run_http_server` builds +/// its runtime here, and the wiring pin reads the flag where it lands — the +/// executor the manager runs tasks on — so dropping the opt-in, or its +/// forward into the executor, turns the pin red instead of silently keeping +/// every served task's wall clock running through prompts. +pub(crate) async fn open_server_task_runtime( + config: &Config, + workspace: &FsPath, + plugin_registry: Arc, + workers: usize, +) -> Result<(SharedRuntimeThreadManager, SharedTaskManager)> { + let task_cfg = server_task_config(config, workspace.to_path_buf(), workers); + let (runtime_threads, _workshop_activation) = open_runtime_threads_for_server( + config, + workspace.to_path_buf(), + RuntimeThreadManagerConfig::from_task_data_dir(task_cfg.data_dir.clone()), + plugin_registry, + )?; + let task_manager = + TaskManager::start_with_runtime_manager(task_cfg, config.clone(), runtime_threads.clone()) + .await?; + Ok((runtime_threads, task_manager)) +} + /// Start the runtime API server. pub async fn run_http_server( config: Config, @@ -865,16 +893,13 @@ pub async fn run_http_server( ) -> Result<()> { validate_runtime_listener_security(&options)?; - let task_cfg = server_task_config(&config, workspace.clone(), options.workers); - let (runtime_threads, _workshop_activation) = open_runtime_threads_for_server( + let (runtime_threads, task_manager) = open_server_task_runtime( &config, - workspace.clone(), - RuntimeThreadManagerConfig::from_task_data_dir(task_cfg.data_dir.clone()), + &workspace, plugin_discovery.registry_for_workspace(&workspace), - )?; - let task_manager = - TaskManager::start_with_runtime_manager(task_cfg, config.clone(), runtime_threads.clone()) - .await?; + options.workers, + ) + .await?; let automations = Arc::new(Mutex::new(AutomationManager::default_location()?)); runtime_threads.attach_automation_manager(automations.clone()); let scheduler_cancel = CancellationToken::new(); diff --git a/crates/tui/src/task_manager.rs b/crates/tui/src/task_manager.rs index 55fa6db485..0daa4baa99 100644 --- a/crates/tui/src/task_manager.rs +++ b/crates/tui/src/task_manager.rs @@ -942,6 +942,15 @@ pub trait TaskExecutor: Send + Sync { events: mpsc::Sender, cancel: CancellationToken, ) -> TaskExecutionResult; + + /// Test seam for the server wiring pin: the built-in executor as the + /// manager actually received it, so the pin reads the answerability flag + /// where `start_with_runtime_manager` forwarded it. Custom executors + /// keep the default — there is nothing to expose. + #[cfg(test)] + fn as_engine_executor(&self) -> Option<&EngineTaskExecutor> { + None + } } /// Executor backed by the shared runtime and its canonical provider resolver. @@ -1051,6 +1060,11 @@ impl TaskExecutor for EngineTaskExecutor { ) .await } + + #[cfg(test)] + fn as_engine_executor(&self) -> Option<&EngineTaskExecutor> { + Some(self) + } } /// Drive one engine turn to its terminal state under the supervisor guard. @@ -4793,6 +4807,41 @@ mod tests { ); } + #[tokio::test] + async fn server_task_runtime_wires_answerable_human_waits_into_the_executor() -> Result<()> { + // The server opt-in must survive every hop of its wiring: the + // `server_task_config` builder and `start_with_runtime_manager`'s + // forward into `EngineTaskExecutor`. Reverting any one hop to the + // unanswerable default left the whole suite green — the review gap + // this pin closes — so the assertion reads the flag where it lands: + // the executor field the manager runs tasks on, built through the + // same `open_server_task_runtime` entry `run_http_server` uses. + let _env_lock = crate::test_support::lock_test_env(); + let tasks_dir = tempfile::tempdir()?.keep(); + let _tasks_dir = + crate::test_support::EnvVarGuard::set("CODEWHALE_TASKS_DIR", tasks_dir.as_os_str()); + let workspace = tempfile::tempdir()?.keep(); + let registry = crate::plugins::PluginDiscoveryContext::capture_pre_dotenv() + .registry_for_workspace(&workspace); + let (_runtime_threads, manager) = crate::runtime_api::open_server_task_runtime( + &Config::default(), + &workspace, + registry, + 1, + ) + .await?; + let engine = manager + .executor + .as_engine_executor() + .expect("the server task runtime must use the built-in engine executor"); + assert!( + engine.answerable_waits, + "the server task runtime must forward the answerable opt-in into its executor" + ); + manager.shutdown(); + Ok(()) + } + #[test] fn consecutive_message_deltas_coalesce_on_the_timeline() { let mut task = sample_task_record();