diff --git a/docs/en/wework/settings.md b/docs/en/wework/settings.md index e7c27b87ee..e64e02d90b 100644 --- a/docs/en/wework/settings.md +++ b/docs/en/wework/settings.md @@ -4,7 +4,7 @@ sidebar_position: 9 # Settings and data -Settings cover language and startup behavior, appearance, Codex and local models, proxies, context, quick phrases, keybindings, worktrees, browser data, and archived conversations. +Settings cover language and startup behavior, appearance, Codex and local models, proxies, context and default principles for the experimental personal supervisor, quick phrases, keybindings, worktrees, browser data, and archived conversations. ## View app information diff --git a/docs/en/wework/tasks.md b/docs/en/wework/tasks.md index 2a35636264..060071e615 100644 --- a/docs/en/wework/tasks.md +++ b/docs/en/wework/tasks.md @@ -32,6 +32,15 @@ To continue the work with a model from the other category, start a new conversat Interrupting stops the current response but does not roll back completed file edits or commands. +## Use the personal supervisor + +Enable **Experimental features** under **Settings → General**, then use **Personal supervisor** above the composer in a Codex task. While the main AI is working, the executor periodically reads its recent progress in the background and evaluates goal drift, missed constraints, destructive actions, and obvious blocked loops with a lightweight read-only call. It does not fork the original task, and checks continue without keeping the task view open. + +- **Suggest** shows a correction above the composer for you to approve or dismiss. +- **Auto-correct** steers an active response when a clear deviation is found, or starts a normal follow-up just as if you had sent the instruction from the composer. + +Supervision settings belong to the current task. The review model can follow the current task or be selected independently, and the review frequency can be 10 seconds, 30 seconds, 1 minute, or 5 minutes. Set default supervisor principles under **Settings → Context**; they are prefilled when supervision is first enabled and can then be customized for that task. + ## Review the processing timeline The **Processed** section in an AI response displays tool calls from top to bottom by their actual creation time. Even when executor events arrive out of order, commands, file operations, and other tools created earlier remain above later activity so the timeline reflects the real execution sequence. diff --git a/docs/zh/wework/settings.md b/docs/zh/wework/settings.md index 3b7b8146ae..f7d6816df8 100644 --- a/docs/zh/wework/settings.md +++ b/docs/zh/wework/settings.md @@ -10,7 +10,7 @@ sidebar_position: 9 - **外观**:主题和工作台背景。背景图片及显示参数保存在当前设备。 - **模型**:本机 Codex 登录、本地 OpenAI Responses 兼容模型和云端模型。 - **代理**:分别配置本机和云端设备的模型访问代理。 -- **上下文**:管理任务上下文偏好和 Codex 个性。 +- **上下文**:管理任务上下文偏好、Codex 个性,以及开启实验性功能后可用的默认分身监督原则。 - **快捷短语**:保存经常使用的任务说明。 - **键盘快捷键**:查看、修改或清除本机快捷键。 - **工作树**:设置 Worktree 根目录、自动清理和保留数量。 diff --git a/docs/zh/wework/tasks.md b/docs/zh/wework/tasks.md index f5e834c1a5..c040f9b612 100644 --- a/docs/zh/wework/tasks.md +++ b/docs/zh/wework/tasks.md @@ -48,6 +48,15 @@ sidebar_position: 4 打断会停止当前回复,但已经完成的文件修改和命令不会自动回滚。 +## 使用分身监督 + +先在“设置 → 通用”中开启“实验性功能”,然后在 Codex 任务的输入框上方开启“分身监督”。主 AI 工作期间,执行器会在后台定时读取最近的 AI 进展,并使用独立的只读轻量调用检查目标偏移、遗漏约束、破坏性操作和明显的阻塞循环。巡检不会 fork 原任务,也不依赖任务界面保持打开。 + +- **先建议**:在输入框上方显示纠正建议,由你决定是否发送。 +- **自动纠正**:发现明确偏差时,运行中直接引导当前回复,空闲时像你在输入框发送要求一样开启后续回复。 + +监督配置属于当前任务。你可以让巡检模型跟随当前任务,也可以单独指定模型;巡检频率可选 10 秒、30 秒、1 分钟或 5 分钟。你还可以在“设置 → 上下文”中设置全局监督原则;首次开启任务监督时会自动带入,并允许针对该任务修改。 + ## 查看处理过程 AI 回复中的“已处理”区域会按实际创建时间从上到下展示工具调用。即使执行器的事件到达顺序发生变化,较早创建的命令、文件操作或其他工具仍会显示在较早位置,便于按照真实执行过程阅读和排查任务。 diff --git a/executor/src/agents/codex.rs b/executor/src/agents/codex.rs index 8ce840de50..b47bc10c72 100644 --- a/executor/src/agents/codex.rs +++ b/executor/src/agents/codex.rs @@ -3747,6 +3747,7 @@ fn resolve_codex_binary(value: &str) -> String { } const CODEX_DANGER_FULL_ACCESS_PERMISSION_PROFILE: &str = ":danger-full-access"; +const CODEX_READ_ONLY_PERMISSION_PROFILE: &str = ":read-only"; pub(crate) fn codex_runtime_approval_policy() -> Value { json!({ @@ -3760,10 +3761,27 @@ pub(crate) fn codex_runtime_approval_policy() -> Value { }) } -fn insert_codex_runtime_permissions(params: &mut serde_json::Map) { +fn codex_runtime_permission_profile(request: &ExecutionRequest) -> &'static str { + if request + .extra + .get("runtime_permission_profile") + .or_else(|| request.extra.get("runtimePermissionProfile")) + .and_then(Value::as_str) + == Some(CODEX_READ_ONLY_PERMISSION_PROFILE) + { + CODEX_READ_ONLY_PERMISSION_PROFILE + } else { + CODEX_DANGER_FULL_ACCESS_PERMISSION_PROFILE + } +} + +fn insert_codex_runtime_permissions( + params: &mut serde_json::Map, + request: &ExecutionRequest, +) { params.insert( "permissions".to_owned(), - Value::String(CODEX_DANGER_FULL_ACCESS_PERMISSION_PROFILE.to_owned()), + Value::String(codex_runtime_permission_profile(request).to_owned()), ); } @@ -3817,13 +3835,14 @@ fn thread_start_params(request: &ExecutionRequest, launch_config: &CodexLaunchCo if let Some(model) = codex_request_model(request) { params.insert("model".to_owned(), Value::String(model)); } + insert_codex_developer_instructions(&mut params, request); append_thread_launch_params(&mut params, launch_config); if let Some(cwd) = request.cwd() { params.insert("cwd".to_owned(), Value::String(cwd.to_owned())); } insert_runtime_workspace_roots(&mut params, request); params.insert("approvalPolicy".to_owned(), codex_runtime_approval_policy()); - insert_codex_runtime_permissions(&mut params); + insert_codex_runtime_permissions(&mut params, request); if request.ephemeral { params.insert("ephemeral".to_owned(), Value::Bool(true)); } @@ -3845,13 +3864,14 @@ fn thread_fork_params( if let Some(model) = codex_request_model(request) { params.insert("model".to_owned(), Value::String(model)); } + insert_codex_developer_instructions(&mut params, request); append_thread_launch_params(&mut params, launch_config); if let Some(cwd) = request.cwd() { params.insert("cwd".to_owned(), Value::String(cwd.to_owned())); } insert_runtime_workspace_roots(&mut params, request); params.insert("approvalPolicy".to_owned(), codex_runtime_approval_policy()); - insert_codex_runtime_permissions(&mut params); + insert_codex_runtime_permissions(&mut params, request); if request.ephemeral { params.insert("ephemeral".to_owned(), Value::Bool(true)); } @@ -3906,16 +3926,30 @@ fn thread_resume_params( if let Some(model) = codex_request_model(request) { params.insert("model".to_owned(), Value::String(model)); } + insert_codex_developer_instructions(&mut params, request); append_thread_launch_params(&mut params, launch_config); if let Some(cwd) = request.cwd() { params.insert("cwd".to_owned(), Value::String(cwd.to_owned())); } insert_runtime_workspace_roots(&mut params, request); params.insert("approvalPolicy".to_owned(), codex_runtime_approval_policy()); - insert_codex_runtime_permissions(&mut params); + insert_codex_runtime_permissions(&mut params, request); Value::Object(params) } +fn insert_codex_developer_instructions( + params: &mut serde_json::Map, + request: &ExecutionRequest, +) { + let instructions = request.system_prompt.trim(); + if !instructions.is_empty() { + params.insert( + "developerInstructions".to_owned(), + Value::String(instructions.to_owned()), + ); + } +} + fn append_thread_launch_params( params: &mut serde_json::Map, launch_config: &CodexLaunchConfig, @@ -3960,7 +3994,7 @@ fn turn_start_params( ); } params.insert("approvalPolicy".to_owned(), codex_runtime_approval_policy()); - insert_codex_runtime_permissions(&mut params); + insert_codex_runtime_permissions(&mut params, request); if let Some(cwd) = request.cwd() { params.insert("cwd".to_owned(), Value::String(cwd.to_owned())); } @@ -3980,9 +4014,21 @@ fn turn_start_params( if let Some(additional_context) = codex_additional_context(request) { params.insert("additionalContext".to_owned(), additional_context); } + if let Some(output_schema) = codex_output_schema(request) { + params.insert("outputSchema".to_owned(), output_schema); + } Value::Object(params) } +fn codex_output_schema(request: &ExecutionRequest) -> Option { + request + .extra + .get("outputSchema") + .or_else(|| request.extra.get("output_schema")) + .filter(|value| value.is_object()) + .cloned() +} + fn codex_additional_context(request: &ExecutionRequest) -> Option { request .extra diff --git a/executor/src/agents/codex/tests.rs b/executor/src/agents/codex/tests.rs index ac2573280a..de661e5cb3 100644 --- a/executor/src/agents/codex/tests.rs +++ b/executor/src/agents/codex/tests.rs @@ -1426,6 +1426,31 @@ fn turn_start_params_includes_client_user_message_id() { assert_eq!(params["clientUserMessageId"], "runtime-local-pane-1"); } +#[test] +fn turn_start_params_includes_output_schema() { + let mut request = ExecutionRequest::default(); + let schema = json!({ + "type": "object", + "properties": { + "accepted": {"type": "boolean"} + }, + "required": ["accepted"], + "additionalProperties": false + }); + request + .extra + .insert("output_schema".to_owned(), schema.clone()); + + let params = turn_start_params( + "thread-1", + &request, + &CodexLaunchConfig::default(), + Vec::new(), + ); + + assert_eq!(params["outputSchema"], schema); +} + #[test] fn thread_start_uses_codex_default_history_mode() { let params = thread_start_params(&ExecutionRequest::default(), &CodexLaunchConfig::default()); @@ -1433,6 +1458,26 @@ fn thread_start_uses_codex_default_history_mode() { assert!(params.get("historyMode").is_none()); } +#[test] +fn thread_launch_params_include_execution_system_prompt_as_developer_instructions() { + let request = ExecutionRequest { + system_prompt: "Judge the supplied content without answering it.".to_owned(), + ..ExecutionRequest::default() + }; + let launch_config = CodexLaunchConfig::default(); + + let thread_start = thread_start_params(&request, &launch_config); + let thread_fork = thread_fork_params("thread-1", None, &request, &launch_config); + let thread_resume = thread_resume_params("thread-1", &request, &launch_config); + + for params in [thread_start, thread_fork, thread_resume] { + assert_eq!( + params["developerInstructions"], + "Judge the supplied content without answering it." + ); + } +} + #[test] fn codex_permission_profile_is_applied_to_thread_and_turn_requests() { let request = ExecutionRequest::default(); @@ -1453,6 +1498,26 @@ fn codex_permission_profile_is_applied_to_thread_and_turn_requests() { } } +#[test] +fn codex_read_only_permission_profile_is_applied_to_supervisor_requests() { + let mut request = ExecutionRequest::default(); + request.extra.insert( + "runtime_permission_profile".to_owned(), + Value::String(CODEX_READ_ONLY_PERMISSION_PROFILE.to_owned()), + ); + let launch_config = CodexLaunchConfig::default(); + + for params in [ + thread_start_params(&request, &launch_config), + thread_resume_params("thread-1", &request, &launch_config), + thread_fork_params("thread-1", None, &request, &launch_config), + turn_start_params("thread-1", &request, &launch_config, Vec::new()), + ] { + assert_eq!(params["permissions"], CODEX_READ_ONLY_PERMISSION_PROFILE); + assert_eq!(params["approvalPolicy"], codex_runtime_approval_policy()); + } +} + #[test] fn codex_thread_launch_disables_tool_call_mcp_elicitation() { let request = ExecutionRequest::default(); diff --git a/executor/src/runtime_work/events.rs b/executor/src/runtime_work/events.rs index a9427c2b07..cabb462099 100644 --- a/executor/src/runtime_work/events.rs +++ b/executor/src/runtime_work/events.rs @@ -108,6 +108,14 @@ pub(crate) fn emit_response_event( payload_object.insert("source".to_owned(), source.clone()); } } + if let Some(generated_user_message) = request.extra.get("runtime_generated_user_message") { + if let Some(payload_object) = payload.get_mut("payload").and_then(Value::as_object_mut) { + payload_object.insert( + "runtimeGeneratedUserMessage".to_owned(), + generated_user_message.clone(), + ); + } + } let receiver_count = event_tx.receiver_count(); let delivery = event_tx.send(payload); if terminal { diff --git a/executor/src/runtime_work/handler.rs b/executor/src/runtime_work/handler.rs index 4e7a2eec92..64dde6c3b0 100644 --- a/executor/src/runtime_work/handler.rs +++ b/executor/src/runtime_work/handler.rs @@ -48,6 +48,7 @@ mod hooks; mod notifications; mod queries; mod sidebar; +mod supervisor; mod system; mod tasks; mod turns; @@ -291,6 +292,7 @@ pub struct RuntimeWorkRpcHandler { active_turn_cancellations: Arc>>, active_codex_turns: Arc>>, active_request_user_inputs: Arc>>>, + supervisor_evaluating: Arc>>, thread_event_routes: Arc>>, notification_router: Arc>>>, archived_delete_tx: mpsc::UnboundedSender, @@ -360,6 +362,7 @@ impl RuntimeWorkRpcHandler { active_turn_cancellations: Arc::new(Mutex::new(HashMap::new())), active_codex_turns: Arc::new(Mutex::new(HashMap::new())), active_request_user_inputs: Arc::new(Mutex::new(HashMap::new())), + supervisor_evaluating: Arc::new(Mutex::new(HashSet::new())), thread_event_routes: Arc::new(Mutex::new(HashMap::new())), notification_router: Arc::new(Mutex::new(None)), archived_delete_tx, @@ -387,6 +390,7 @@ impl RuntimeWorkRpcHandler { handler.hook_service.set_event_sender(sender); } handler.start_automation_scheduler(); + handler.start_supervisor_scheduler(); handler } @@ -418,6 +422,10 @@ impl RuntimeWorkRpcHandler { "runtime.tasks.goal.get" => self.get_task_goal(payload).await, "runtime.tasks.goal.set" => self.set_task_goal(payload).await, "runtime.tasks.goal.clear" => self.clear_task_goal(payload).await, + "runtime.tasks.supervisor.get" => self.get_task_supervisor(payload).await, + "runtime.tasks.supervisor.set" => self.set_task_supervisor(payload).await, + "runtime.tasks.supervisor.clear" => self.clear_task_supervisor(payload).await, + "runtime.tasks.supervisor.resolve" => self.resolve_task_supervisor(payload).await, "runtime.keybindings.get" => self.get_keybindings().await, "runtime.keybindings.update" => self.update_keybindings(payload).await, "runtime.hooks.list" | "runtime.hooks.reload" => { diff --git a/executor/src/runtime_work/handler/helpers/transcript.rs b/executor/src/runtime_work/handler/helpers/transcript.rs index b2b5e867c6..4f0228986f 100644 --- a/executor/src/runtime_work/handler/helpers/transcript.rs +++ b/executor/src/runtime_work/handler/helpers/transcript.rs @@ -323,47 +323,91 @@ fn user_message_presentation(payload: &Value) -> Option { let content = payload .get("message") .or_else(|| payload.get("content")) - .and_then(Value::as_str)?; + .and_then(Value::as_str) + .unwrap_or_default(); let references = local_presentation_reference_descriptors(content); - (!references.is_empty()).then(|| { - json!({ - "clientUserMessageId": client_user_message_id, - "references": references, - }) - }) + let source = payload.get("source").filter(|value| value.is_object()).cloned(); + if references.is_empty() && source.is_none() { + return None; + } + let mut presentation = json!({ + "clientUserMessageId": client_user_message_id, + "references": references, + "source": source, + }); + let supervisor_generated = presentation + .get("source") + .and_then(|source| source.get("source")) + .and_then(Value::as_str) + == Some("supervisor"); + if supervisor_generated { + presentation["content"] = Value::String(content.to_owned()); + presentation["createdAt"] = json!(now_ms()); + presentation["ensureVisible"] = Value::Bool(true); + } + Some(presentation) } -fn attach_user_message_presentations(messages: &mut [Value], presentations: Vec) { +fn attach_user_message_presentations(messages: &mut Vec, presentations: Vec) { for presentation in presentations { let Some(client_user_message_id) = string_field(&presentation, "clientUserMessageId") .or_else(|| string_field(&presentation, "client_user_message_id")) else { continue; }; - let Some(message) = messages.iter_mut().find(|message| { + let message_index = messages.iter().position(|message| { string_field(message, "clientUserMessageId") .or_else(|| string_field(message, "client_user_message_id")) .as_deref() == Some(client_user_message_id.as_str()) - }) else { - continue; - }; - let Some(content) = string_field(message, "content") else { - continue; + }); + let message_index = match message_index { + Some(index) => index, + None if bool_field(&presentation, "ensureVisible") == Some(true) => { + let content = string_field(&presentation, "content").unwrap_or_default(); + if content.trim().is_empty() { + continue; + } + let created_at = + timestamp_ms_field(&presentation, "createdAt").unwrap_or_else(now_ms); + let synthetic = json!({ + "id": client_user_message_id, + "clientUserMessageId": client_user_message_id, + "role": "user", + "content": content, + "status": "done", + "createdAt": created_at, + "source": presentation.get("source").cloned(), + }); + let index = messages + .iter() + .position(|message| { + timestamp_ms_field(message, "createdAt") + .is_some_and(|message_at| message_at > created_at) + }) + .unwrap_or(messages.len()); + messages.insert(index, synthetic); + index + } + None => continue, }; + let message = &mut messages[message_index]; + let content = string_field(message, "content").unwrap_or_default(); let references = presentation .get("references") .and_then(Value::as_array) .map(|references| presentation_reference_ranges(references, &content)) .unwrap_or_default(); - if references.is_empty() { - continue; - } if let Some(message) = message.as_object_mut() { - message.insert( - "presentationReferences".to_owned(), - Value::Array(references), - ); + if !references.is_empty() { + message.insert( + "presentationReferences".to_owned(), + Value::Array(references), + ); + } + if let Some(source) = presentation.get("source").filter(|value| value.is_object()) { + message.insert("source".to_owned(), source.clone()); + } } } } diff --git a/executor/src/runtime_work/handler/supervisor.rs b/executor/src/runtime_work/handler/supervisor.rs new file mode 100644 index 0000000000..fa66ee5d0a --- /dev/null +++ b/executor/src/runtime_work/handler/supervisor.rs @@ -0,0 +1,821 @@ +// SPDX-FileCopyrightText: 2026 Weibo, Inc. +// +// SPDX-License-Identifier: Apache-2.0 + +use super::*; +use crate::runtime_work::response::{RuntimeSupervisorState, RuntimeSupervisorSuggestion}; +use serde::Deserialize; +use sha2::{Digest, Sha256}; + +const SUPERVISOR_MODES: &[&str] = &["suggest", "auto"]; +const SUPERVISOR_VISIBLE_MESSAGES: usize = 6; +const SUPERVISOR_LATEST_CONTENT_CHARS: usize = 12_000; +const SUPERVISOR_CONTEXT_CONTENT_CHARS: usize = 2_400; +const SUPERVISOR_PROMPT_VERSION: &str = "5"; +const SUPERVISOR_REPEAT_CORRECTION_WINDOW_MS: i64 = 5 * 60 * 1_000; +const SUPERVISOR_DEFAULT_INTERVAL_SECONDS: u64 = 30; +const SUPERVISOR_INTERVAL_SECONDS: &[u64] = &[10, 30, 60, 300]; +const SUPERVISOR_SCHEDULER_INTERVAL: Duration = Duration::from_secs(5); +const SUPERVISOR_EVALUATION_TIMEOUT: Duration = Duration::from_secs(60); + +impl RuntimeWorkRpcHandler { + pub(super) async fn get_task_supervisor(&self, payload: Value) -> Result { + let link = self.task_link_from_payload(&payload, false).await?; + Ok(supervisor_response(&link)) + } + + pub(super) async fn set_task_supervisor(&self, payload: Value) -> Result { + let link = self.task_link_from_payload(&payload, false).await?; + let mode = string_field(&payload, "mode").unwrap_or_else(|| "suggest".to_owned()); + if !SUPERVISOR_MODES.contains(&mode.as_str()) { + return Err(AppIpcError::new( + "bad_request", + "supervisor mode must be suggest or auto", + )); + } + let instructions = string_field(&payload, "instructions").unwrap_or_default(); + let model_id = + string_field(&payload, "modelId").or_else(|| string_field(&payload, "model_id")); + let interval_seconds = payload + .get("intervalSeconds") + .or_else(|| payload.get("interval_seconds")) + .and_then(Value::as_u64) + .unwrap_or(SUPERVISOR_DEFAULT_INTERVAL_SECONDS); + if !SUPERVISOR_INTERVAL_SECONDS.contains(&interval_seconds) { + return Err(AppIpcError::new( + "bad_request", + "supervisor interval must be 10, 30, 60, or 300 seconds", + )); + } + self.store.update_task(&link.local_task_id, |task| { + let existing = task.supervisor.take(); + task.supervisor = Some(RuntimeSupervisorState { + mode: mode.clone(), + status: "active".to_owned(), + instructions: instructions.clone(), + model_id: model_id.clone().or_else(|| { + existing + .as_ref() + .and_then(|supervisor| supervisor.model_id.clone()) + }), + interval_seconds, + last_evaluated_at: None, + last_content_hash: None, + last_error: None, + suggestions: existing + .map(|supervisor| supervisor.suggestions) + .unwrap_or_default(), + }); + task.updated_at = now_ms(); + }); + self.emit_supervisor_updated(&link.local_task_id); + self.schedule_supervisor_evaluation(link.local_task_id.clone(), None); + Ok(supervisor_response( + &self.local_task_link(&link.local_task_id).unwrap_or(link), + )) + } + + pub(super) async fn clear_task_supervisor(&self, payload: Value) -> Result { + let link = self.task_link_from_payload(&payload, false).await?; + self.store.update_task(&link.local_task_id, |task| { + task.supervisor = None; + task.updated_at = now_ms(); + }); + self.emit_supervisor_updated(&link.local_task_id); + Ok(json!({ + "success": true, + "accepted": true, + "taskId": link.local_task_id, + "runtime": "codex", + "supervisor": null, + })) + } + + pub(super) async fn resolve_task_supervisor( + &self, + payload: Value, + ) -> Result { + let link = self.task_link_from_payload(&payload, false).await?; + let suggestion_id = string_field(&payload, "suggestionId") + .or_else(|| string_field(&payload, "suggestion_id")) + .ok_or_else(|| AppIpcError::new("bad_request", "suggestionId is required"))?; + let status = string_field(&payload, "status").unwrap_or_else(|| "dismissed".to_owned()); + if !["accepted", "dismissed"].contains(&status.as_str()) { + return Err(AppIpcError::new( + "bad_request", + "suggestion status must be accepted or dismissed", + )); + } + self.store.update_task(&link.local_task_id, |task| { + let Some(supervisor) = task.supervisor.as_mut() else { + return; + }; + if let Some(suggestion) = supervisor + .suggestions + .iter_mut() + .find(|suggestion| suggestion.id == suggestion_id) + { + suggestion.status = status.clone(); + suggestion.resolved_at = Some(now_ms()); + } + task.updated_at = now_ms(); + }); + self.emit_supervisor_updated(&link.local_task_id); + Ok(supervisor_response( + &self.local_task_link(&link.local_task_id).unwrap_or(link), + )) + } + + pub(super) fn schedule_supervisor_evaluation( + &self, + local_task_id: String, + source_turn_id: Option, + ) { + let enabled = self + .local_task_link(&local_task_id) + .and_then(|link| link.supervisor) + .is_some_and(|supervisor| supervisor.status != "disabled"); + if !enabled { + return; + } + let evaluation_started = self + .supervisor_evaluating + .lock() + .expect("supervisor evaluating set lock should not be poisoned") + .insert(local_task_id.clone()); + if !evaluation_started { + return; + } + let handler = self.clone(); + tokio::spawn(async move { + if let Err(error) = handler + .evaluate_task_supervisor(&local_task_id, source_turn_id) + .await + { + handler.record_supervisor_error(&local_task_id, error); + } + handler + .supervisor_evaluating + .lock() + .expect("supervisor evaluating set lock should not be poisoned") + .remove(&local_task_id); + }); + } + + pub(super) fn start_supervisor_scheduler(&self) { + let Ok(handle) = tokio::runtime::Handle::try_current() else { + return; + }; + let handler = self.clone(); + handle.spawn(async move { + let mut interval = tokio::time::interval_at( + tokio::time::Instant::now() + SUPERVISOR_SCHEDULER_INTERVAL, + SUPERVISOR_SCHEDULER_INTERVAL, + ); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip); + loop { + interval.tick().await; + for link in handler.local_task_links(false) { + if supervisor_needs_scheduled_check( + &link, + handler.is_active_local_task(&link.local_task_id), + now_ms(), + ) { + handler.schedule_supervisor_evaluation(link.local_task_id, None); + } + } + } + }); + } + + async fn evaluate_task_supervisor( + &self, + local_task_id: &str, + source_turn_id: Option, + ) -> Result<(), String> { + let link = self + .local_task_link(local_task_id) + .ok_or_else(|| "runtime task was not found".to_owned())?; + let supervisor = link + .supervisor + .clone() + .ok_or_else(|| "supervisor is disabled".to_owned())?; + let parent_thread_id = runtime_session_id_from_link(&link) + .ok_or_else(|| "runtime task session is not ready".to_owned())?; + let thread = self.read_codex_thread_with_turns(&parent_thread_id).await?; + let current_configuration = self + .local_task_link(local_task_id) + .and_then(|task| task.supervisor); + if !current_configuration.as_ref().is_some_and(|current| { + current.mode == supervisor.mode + && current.instructions == supervisor.instructions + && current.model_id == supervisor.model_id + && current.interval_seconds == supervisor.interval_seconds + }) { + return Ok(()); + } + let Some(visible_ai_progress) = + visible_ai_progress(transcript_messages(&thread, &self.device_id))? + else { + self.store.update_task(local_task_id, |task| { + let Some(current) = task.supervisor.as_mut() else { + return; + }; + current.last_evaluated_at = Some(now_ms()); + current.last_content_hash = Some(empty_content_hash()); + current.last_error = None; + current.status = "active".to_owned(); + }); + self.emit_supervisor_updated(local_task_id); + return Ok(()); + }; + let content_hash = content_hash(&visible_ai_progress); + if supervisor.last_content_hash.as_deref() == Some(content_hash.as_str()) { + self.store.update_task(local_task_id, |task| { + let Some(current) = task.supervisor.as_mut() else { + return; + }; + current.last_evaluated_at = Some(now_ms()); + current.last_error = None; + current.status = "active".to_owned(); + }); + return Ok(()); + } + let snapshot_at = now_ms(); + self.store.update_task(local_task_id, |task| { + let Some(current) = task.supervisor.as_mut() else { + return; + }; + current.status = "checking".to_owned(); + current.last_error = None; + task.updated_at = now_ms(); + }); + self.emit_supervisor_updated(local_task_id); + let visible_progress = json!({ + "task": { + "title": link.title, + "status": link.status, + "running": link.running, + }, + "latestAiContent": visible_ai_progress.latest, + "recentAiContext": visible_ai_progress.context, + "lastSupervisorIntervention": supervisor.suggestions.iter().rev().find(|suggestion| { + suggestion.status == "accepted" + }).map(|suggestion| json!({ + "message": suggestion.message, + "createdAt": suggestion.created_at, + })), + }); + let visible_progress = + serde_json::to_string(&visible_progress).map_err(|error| error.to_string())?; + let mut request = ExecutionRequest { + task_id: format!("supervisor:{local_task_id}"), + subtask_id: format!("supervisor-evaluation-{}", now_ms()), + system_prompt: supervisor_system_prompt(), + prompt: Value::String(supervisor_prompt(&supervisor, &visible_progress)), + project_workspace_path: Some(link.workspace_path.clone()), + runtime_workspace_roots: link.runtime_workspace_roots.clone(), + runtime_project_key: link.runtime_project_key.clone(), + ephemeral: true, + ..ExecutionRequest::default() + }; + if let Some(model_id) = supervisor + .model_id + .clone() + .or_else(|| task_model_id(&link.runtime_handle)) + { + request.model_config = json!({"model_id": model_id}); + } + request.extra.insert( + "runtime_permission_profile".to_owned(), + Value::String(":read-only".to_owned()), + ); + request.extra.insert( + "runtime_message_source".to_owned(), + Value::String("supervisor-evaluator".to_owned()), + ); + request + .extra + .insert("output_schema".to_owned(), supervisor_output_schema()); + let (cancel_tx, cancel_rx) = oneshot::channel(); + let timeout_task = tokio::spawn(async move { + sleep(SUPERVISOR_EVALUATION_TIMEOUT).await; + let _ = cancel_tx.send(()); + }); + let result = self + .codex_app_server + .run_turn_with_cancel( + request, + CodexAppServerTurnOptions { + cancellation: Some(cancel_rx), + ..CodexAppServerTurnOptions::default() + }, + ) + .await; + timeout_task.abort(); + let result = result.map_err(|error| { + if error == CODEX_APP_SERVER_TURN_CANCELLED { + format!( + "supervisor evaluation timed out after {} seconds", + SUPERVISOR_EVALUATION_TIMEOUT.as_secs() + ) + } else { + error + } + })?; + let content = match result.outcome { + ExecutionOutcome::Completed { content } => content, + ExecutionOutcome::Failed { message } | ExecutionOutcome::Cancelled { message } => { + return Err(message) + } + ExecutionOutcome::WaitingForUserInput { stop_reason } => return Err(stop_reason), + ExecutionOutcome::Running => { + return Err("supervisor evaluation did not finish".to_owned()) + } + }; + let evaluation = parse_supervisor_evaluation(&content)?; + let Some(current_supervisor) = self + .local_task_link(local_task_id) + .and_then(|task| task.supervisor) + else { + return Ok(()); + }; + if current_supervisor.mode != supervisor.mode + || current_supervisor.instructions != supervisor.instructions + || current_supervisor.model_id != supervisor.model_id + || current_supervisor.interval_seconds != supervisor.interval_seconds + { + return Ok(()); + } + self.store.update_task(local_task_id, |task| { + let Some(current) = task.supervisor.as_mut() else { + return; + }; + current.last_evaluated_at = Some(snapshot_at); + current.last_content_hash = Some(content_hash.clone()); + current.last_error = None; + current.status = "active".to_owned(); + }); + let Some(correction) = evaluation + .correction + .as_deref() + .map(str::trim) + .filter(|correction| !correction.is_empty()) + else { + self.emit_supervisor_updated(local_task_id); + return Ok(()); + }; + if recently_sent_same_correction(¤t_supervisor, correction, now_ms()) { + self.emit_supervisor_updated(local_task_id); + return Ok(()); + } + match current_supervisor.mode.as_str() { + "auto" => { + self.send_supervisor_correction(local_task_id, &link, correction) + .await?; + self.record_supervisor_suggestion( + local_task_id, + evaluation, + "accepted", + source_turn_id, + ); + } + "suggest" => self.record_supervisor_suggestion( + local_task_id, + evaluation, + "pending", + source_turn_id, + ), + mode => return Err(format!("unsupported supervisor mode: {mode}")), + } + Ok(()) + } + + fn record_supervisor_suggestion( + &self, + local_task_id: &str, + evaluation: SupervisorEvaluation, + status: &str, + source_turn_id: Option, + ) { + self.store.update_task(local_task_id, |task| { + let Some(supervisor) = task.supervisor.as_mut() else { + return; + }; + supervisor.suggestions.push(RuntimeSupervisorSuggestion { + id: format!("supervisor-suggestion-{}", now_ms()), + message: evaluation.correction.clone().unwrap_or_default(), + rationale: evaluation.rationale.clone(), + status: status.to_owned(), + created_at: now_ms(), + resolved_at: (status != "pending").then(now_ms), + source_turn_id: source_turn_id.clone(), + }); + if supervisor.suggestions.len() > 50 { + supervisor.suggestions.remove(0); + } + task.updated_at = now_ms(); + }); + self.emit_supervisor_updated(local_task_id); + } + + async fn send_supervisor_correction( + &self, + local_task_id: &str, + link: &RuntimeTaskLink, + message: &str, + ) -> Result<(), String> { + let client_user_message_id = format!("supervisor-correction-{}", now_ms()); + let source = json!({ + "source": "supervisor", + "channel_type": "task_supervisor", + "channel_label": "分身监督", + }); + let created_at = now_ms(); + if self.is_active_local_task(local_task_id) { + let response = self + .send_guidance(json!({ + "taskId": local_task_id, + "message": message, + "clientGuidanceId": client_user_message_id, + })) + .await + .map_err(|error| error.message)?; + if response.get("accepted").and_then(Value::as_bool) == Some(true) { + return Ok(()); + } + if string_field(&response, "code").as_deref() != Some("no_active_turn") { + return Err(string_field(&response, "error") + .unwrap_or_else(|| "supervisor guidance was rejected".to_owned())); + } + } + let thread_id = runtime_session_id_from_link(link) + .ok_or_else(|| "runtime task session is not ready".to_owned())?; + let mut request = ExecutionRequest { + task_id: local_task_id.to_owned(), + subtask_id: format!("supervisor-correction-{}", now_ms()), + prompt: Value::String(message.to_owned()), + project_workspace_path: Some(link.workspace_path.clone()), + runtime_workspace_roots: link.runtime_workspace_roots.clone(), + runtime_project_key: link.runtime_project_key.clone(), + ephemeral: link.ephemeral, + ..ExecutionRequest::default() + }; + if let Some(model_id) = task_model_id(&link.runtime_handle) { + request.model_config = json!({"model_id": model_id}); + } + request.extra.insert( + "client_user_message_id".to_owned(), + Value::String(client_user_message_id.clone()), + ); + request.extra.insert( + "runtime_message_source".to_owned(), + Value::String("supervisor".to_owned()), + ); + request.extra.insert( + "runtime_generated_user_message".to_owned(), + json!({ + "id": client_user_message_id, + "message": message, + "createdAt": created_at, + "source": source.clone(), + }), + ); + self.mark_task_running_for_send( + local_task_id, + &thread_id, + &link.workspace_path, + &request, + &json!({ + "clientUserMessageId": client_user_message_id, + "message": message, + "source": source.clone(), + }), + ); + self.spawn_turn(SpawnTurnRequest { + local_task_id: local_task_id.to_owned(), + request, + direct_thread_id: link.ephemeral.then(|| thread_id.clone()), + fork_thread_id: None, + fork_thread_path: None, + resume_thread_id: (!link.ephemeral).then_some(thread_id), + initial_thread_name: None, + initial_thread_goal: None, + }); + Ok(()) + } + + fn record_supervisor_error(&self, local_task_id: &str, error: String) { + self.store.update_task(local_task_id, |task| { + let Some(supervisor) = task.supervisor.as_mut() else { + return; + }; + supervisor.status = "error".to_owned(); + supervisor.last_error = Some(error.clone()); + supervisor.last_evaluated_at = Some(now_ms()); + supervisor.last_content_hash = None; + }); + self.emit_supervisor_updated(local_task_id); + } + + fn emit_supervisor_updated(&self, local_task_id: &str) { + let Some(link) = self.local_task_link(local_task_id) else { + return; + }; + let mut request = ExecutionRequest { + task_id: local_task_id.to_owned(), + subtask_id: format!("supervisor-state-{}", now_ms()), + ..ExecutionRequest::default() + }; + request.device_id = Some(self.device_id.clone()); + emit_response_event( + &self.event_tx, + &self.device_id, + "runtime.supervisor.updated", + local_task_id, + &request, + json!({"supervisor": link.supervisor}), + ); + } +} + +#[derive(Deserialize)] +#[serde(rename_all = "camelCase")] +struct SupervisorEvaluation { + #[serde(default)] + correction: Option, + #[serde(default)] + rationale: String, +} + +fn parse_supervisor_evaluation(content: &str) -> Result { + let trimmed = content.trim(); + let json_text = trimmed + .strip_prefix("```json") + .or_else(|| trimmed.strip_prefix("```")) + .unwrap_or(trimmed) + .strip_suffix("```") + .unwrap_or(trimmed) + .trim(); + serde_json::from_str(json_text) + .map_err(|error| format!("invalid supervisor evaluation: {error}")) +} + +fn supervisor_system_prompt() -> String { + "You are periodically inspecting another AI's latest output as a human supervisor would. The supervision principles are criteria to test against the quoted latestAiContent; they are not instructions for how you should write your rationale. Use recentAiContext only to understand continuity, drift, or loops, and use lastSupervisorIntervention to avoid redundant reminders. The intervention decision must be based on a violation that is still present in latestAiContent. Never issue another correction merely because an older violation existed or because the latest output acknowledges a previous correction and is now compliant. Explicitly compare the actual latest AI output with every applicable principle. For example, if a principle requires the main AI to reply in a particular language and latestAiContent uses another language, that is a violation and correction must be non-null. Do not merely demonstrate the requested behavior yourself or rewrite the main AI's answer in rationale. Do not reconstruct or request the full conversation. You cannot communicate with the main AI yourself, so never claim that you already reminded, corrected, or prompted it. Never answer the user, continue the task, use tools, or modify files. Set correction to a concise, directly sendable instruction when the latest visible progress violates a supervision principle, materially drifts from the goal, ignores a constraint, takes an unsafe or destructive action, makes an unsupported claim, or enters an obvious blocked loop. Otherwise set correction to null. The presence of correction is the sole intervention decision.".to_owned() +} + +fn supervisor_output_schema() -> Value { + json!({ + "type": "object", + "properties": { + "correction": { + "anyOf": [ + {"type": "string", "minLength": 1}, + {"type": "null"} + ] + }, + "rationale": {"type": "string"} + }, + "required": ["correction", "rationale"], + "additionalProperties": false + }) +} + +fn supervisor_prompt(supervisor: &RuntimeSupervisorState, visible_progress: &str) -> String { + format!( + "Mode: {}\nUser supervision principles:\n{}\n\nCurrent visible progress snapshot (JSON):\n{}", + supervisor.mode, + if supervisor.instructions.trim().is_empty() { + "(none)" + } else { + supervisor.instructions.trim() + }, + visible_progress + ) +} + +fn supervisor_response(link: &RuntimeTaskLink) -> Value { + json!({ + "success": true, + "accepted": true, + "taskId": link.local_task_id, + "runtime": "codex", + "supervisor": link.supervisor, + }) +} + +fn task_model_id(runtime_handle: &Value) -> Option { + let selection = runtime_handle.get("modelSelection")?; + string_field(selection, "modelName") + .or_else(|| string_field(selection, "model")) + .or_else(|| string_field(selection, "modelId")) +} + +fn supervisor_needs_scheduled_check(link: &RuntimeTaskLink, active: bool, now: i64) -> bool { + let Some(supervisor) = link.supervisor.as_ref() else { + return false; + }; + if supervisor.last_evaluated_at.is_none() || supervisor.last_content_hash.is_none() { + return true; + } + if active + && supervisor + .last_evaluated_at + .is_some_and(|last_evaluated_at| { + now.saturating_sub(last_evaluated_at) + >= i64::try_from(supervisor.interval_seconds).unwrap_or(i64::MAX) * 1_000 + }) + { + return true; + } + link.completed_at + .zip(supervisor.last_evaluated_at) + .is_some_and(|(completed_at, evaluated_at)| completed_at > evaluated_at) +} + +struct VisibleAiProgress { + latest: String, + context: Vec, +} + +fn visible_ai_progress(messages: Vec) -> Result, String> { + let mut assistant_messages = messages + .into_iter() + .filter(|message| message.get("role").and_then(Value::as_str) == Some("assistant")) + .rev() + .take(SUPERVISOR_VISIBLE_MESSAGES) + .collect::>(); + let Some(latest_assistant_message) = assistant_messages.first() else { + return Ok(None); + }; + let latest = + serde_json::to_string(&latest_assistant_message).map_err(|error| error.to_string())?; + let context = assistant_messages + .drain(1..) + .rev() + .map(|message| { + serde_json::to_string(&message) + .map(|content| truncate_visible_tail(&content, SUPERVISOR_CONTEXT_CONTENT_CHARS)) + .map_err(|error| error.to_string()) + }) + .collect::, _>>()?; + Ok(Some(VisibleAiProgress { + latest: truncate_visible_tail(&latest, SUPERVISOR_LATEST_CONTENT_CHARS), + context, + })) +} + +fn content_hash(content: &VisibleAiProgress) -> String { + let content = json!({ + "latest": content.latest, + "context": content.context, + }); + format!( + "{:x}", + Sha256::digest(format!("{SUPERVISOR_PROMPT_VERSION}:{content}").as_bytes()) + ) +} + +fn empty_content_hash() -> String { + format!( + "{:x}", + Sha256::digest(format!("{SUPERVISOR_PROMPT_VERSION}:empty").as_bytes()) + ) +} + +fn recently_sent_same_correction( + supervisor: &RuntimeSupervisorState, + correction: &str, + now: i64, +) -> bool { + supervisor.suggestions.iter().rev().any(|suggestion| { + suggestion.status == "accepted" + && suggestion.message.trim() == correction + && now.saturating_sub(suggestion.created_at) <= SUPERVISOR_REPEAT_CORRECTION_WINDOW_MS + }) +} + +fn truncate_visible_tail(value: &str, max_chars: usize) -> String { + let char_count = value.chars().count(); + if char_count <= max_chars { + return value.to_owned(); + } + value.chars().skip(char_count - max_chars).collect() +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parses_fenced_supervisor_json() { + let result = parse_supervisor_evaluation( + "```json\n{\"correction\":\"stop\",\"rationale\":\"drift\"}\n```", + ) + .expect("evaluation should parse"); + assert_eq!(result.correction.as_deref(), Some("stop")); + } + + #[test] + fn parses_supervisor_pass_without_correction() { + let result = parse_supervisor_evaluation("{\"correction\":null,\"rationale\":\"aligned\"}") + .expect("evaluation should parse"); + assert_eq!(result.correction, None); + } + + #[test] + fn visible_ai_progress_separates_latest_output_from_recent_context() { + let messages = vec![ + json!({"role": "user", "content": "private original request"}), + json!({"role": "assistant", "content": "older progress"}), + json!({"role": "assistant", "content": "latest progress"}), + ]; + + let progress = visible_ai_progress(messages) + .expect("snapshot should serialize") + .expect("assistant content should exist"); + + assert!(!progress.latest.contains("private original request")); + assert!(!progress.latest.contains("older progress")); + assert!(progress.latest.contains("latest progress")); + assert_eq!(progress.context.len(), 1); + assert!(progress.context[0].contains("older progress")); + } + + #[test] + fn visible_tail_is_bounded_from_the_latest_content() { + assert_eq!(truncate_visible_tail("abcdef", 3), "def"); + } + + #[test] + fn completed_content_is_checked_even_when_the_turn_finishes_between_ticks() { + let link = RuntimeTaskLink { + completed_at: Some(200), + supervisor: Some(RuntimeSupervisorState { + mode: "auto".to_owned(), + status: "active".to_owned(), + instructions: String::new(), + model_id: None, + interval_seconds: 30, + last_evaluated_at: Some(100), + last_content_hash: Some("old".to_owned()), + last_error: None, + suggestions: Vec::new(), + }), + ..RuntimeTaskLink::default() + }; + + assert!(supervisor_needs_scheduled_check(&link, false, 200)); + } + + #[test] + fn stopped_task_without_ai_content_is_not_polled_again() { + let link = RuntimeTaskLink { + completed_at: Some(100), + supervisor: Some(RuntimeSupervisorState { + mode: "auto".to_owned(), + status: "active".to_owned(), + instructions: String::new(), + model_id: None, + interval_seconds: 30, + last_evaluated_at: Some(100), + last_content_hash: Some(empty_content_hash()), + last_error: None, + suggestions: Vec::new(), + }), + ..RuntimeTaskLink::default() + }; + + assert!(!supervisor_needs_scheduled_check(&link, false, 200)); + } + + #[test] + fn repeated_auto_correction_is_suppressed_for_a_short_window() { + let supervisor = RuntimeSupervisorState { + mode: "auto".to_owned(), + status: "active".to_owned(), + instructions: String::new(), + model_id: None, + interval_seconds: 30, + last_evaluated_at: Some(100), + last_content_hash: Some("hash".to_owned()), + last_error: None, + suggestions: vec![RuntimeSupervisorSuggestion { + id: "suggestion-1".to_owned(), + message: "Use Japanese".to_owned(), + rationale: String::new(), + status: "accepted".to_owned(), + created_at: 1_000, + resolved_at: Some(1_000), + source_turn_id: None, + }], + }; + + assert!(recently_sent_same_correction( + &supervisor, + "Use Japanese", + 2_000 + )); + } +} diff --git a/executor/src/runtime_work/handler/tests.rs b/executor/src/runtime_work/handler/tests.rs index c517134465..90c675de5d 100644 --- a/executor/src/runtime_work/handler/tests.rs +++ b/executor/src/runtime_work/handler/tests.rs @@ -680,6 +680,33 @@ fn transcript_does_not_attach_presentation_to_an_unmatched_client_user_message_i assert!(provider_messages[0].get("presentationReferences").is_none()); } +#[test] +fn transcript_restores_a_missing_supervisor_generated_user_message() { + let mut provider_messages = vec![json!({ + "id": "assistant-1", + "role": "assistant", + "content": "Corrected", + "createdAt": 200 + })]; + let presentations = vec![json!({ + "clientUserMessageId": "supervisor-correction-1", + "content": "Use Japanese", + "createdAt": 100, + "ensureVisible": true, + "references": [], + "source": { + "source": "supervisor", + "channel_type": "task_supervisor" + } + })]; + + attach_user_message_presentations(&mut provider_messages, presentations); + + assert_eq!(provider_messages[0]["role"], "user"); + assert_eq!(provider_messages[0]["content"], "Use Japanese"); + assert_eq!(provider_messages[1]["role"], "assistant"); +} + #[test] fn transcript_only_adds_presentations_missing_from_provider_content() { let provider_content = "[$first](/tmp/first/SKILL.md) and $second"; diff --git a/executor/src/runtime_work/response.rs b/executor/src/runtime_work/response.rs index dd2fa2ee88..418045089e 100644 --- a/executor/src/runtime_work/response.rs +++ b/executor/src/runtime_work/response.rs @@ -17,6 +17,38 @@ use super::util::{ string_field, timestamp_ms_field, workspace_group_path, workspace_label, workspace_task_path, }; +#[derive(Debug, Clone, Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct RuntimeSupervisorSuggestion { + pub id: String, + pub message: String, + pub rationale: String, + pub status: String, + pub created_at: i64, + pub resolved_at: Option, + pub source_turn_id: Option, +} + +#[derive(Debug, Clone, Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct RuntimeSupervisorState { + pub mode: String, + pub status: String, + pub instructions: String, + pub model_id: Option, + #[serde(default = "default_supervisor_interval_seconds")] + pub interval_seconds: u64, + pub last_evaluated_at: Option, + #[serde(default)] + pub last_content_hash: Option, + pub last_error: Option, + pub suggestions: Vec, +} + +fn default_supervisor_interval_seconds() -> u64 { + 30 +} + #[derive(Debug, Clone, Deserialize, Serialize)] #[serde(default)] pub(crate) struct RuntimeTaskLink { @@ -32,6 +64,7 @@ pub(crate) struct RuntimeTaskLink { pub thread_status: String, pub turn_status: Option, pub goal_status: Option, + pub supervisor: Option, #[serde(skip)] pub git_info: Option, pub created_at: i64, @@ -68,6 +101,7 @@ impl RuntimeTaskLink { thread_status: "active".to_owned(), turn_status: Some("inProgress".to_owned()), goal_status: None, + supervisor: None, git_info: None, created_at: now_ms(), updated_at: now_ms(), @@ -105,6 +139,7 @@ impl RuntimeTaskLink { thread_status: "notLoaded".to_owned(), turn_status: None, goal_status: None, + supervisor: None, git_info: None, created_at: now_ms(), updated_at: now_ms(), @@ -135,6 +170,7 @@ impl RuntimeTaskLink { let goal_status = local_link .as_ref() .and_then(|link| link.goal_status.clone()); + let supervisor = local_link.as_ref().and_then(|link| link.supervisor.clone()); let mut git_info = thread .get("gitInfo") .or_else(|| thread.get("git_info")) @@ -186,6 +222,7 @@ impl RuntimeTaskLink { thread_status, turn_status, goal_status, + supervisor, git_info, created_at: timestamp_ms_field(thread, "createdAt").unwrap_or_else(now_ms), updated_at: timestamp_ms_field(thread, "updatedAt").unwrap_or_else(now_ms), @@ -229,6 +266,7 @@ impl RuntimeTaskLink { thread_status: self.thread_status.clone(), turn_status: self.turn_status.clone(), goal_status: self.goal_status.clone(), + supervisor: self.supervisor.clone(), git_info: self.git_info.clone(), created_at: self.created_at, updated_at: self.updated_at, @@ -282,6 +320,7 @@ impl Default for RuntimeTaskLink { thread_status: "notLoaded".to_owned(), turn_status: None, goal_status: None, + supervisor: None, git_info: None, created_at: now_ms(), updated_at: now_ms(), @@ -593,6 +632,11 @@ fn local_task_json(link: RuntimeTaskLink) -> Value { if let Some(goal_status) = link.goal_status.clone() { task.insert("goalStatus".to_owned(), Value::String(goal_status)); } + if let Some(supervisor) = link.supervisor.clone() { + if let Ok(supervisor) = serde_json::to_value(supervisor) { + task.insert("supervisor".to_owned(), supervisor); + } + } task.insert("pinned".to_owned(), Value::Bool(link.pinned)); if let Some(order) = link.pinned_order { task.insert("pinnedOrder".to_owned(), json!(order)); diff --git a/executor/src/runtime_work/worktrees.rs b/executor/src/runtime_work/worktrees.rs index d8fda7e1cf..262457dfe5 100644 --- a/executor/src/runtime_work/worktrees.rs +++ b/executor/src/runtime_work/worktrees.rs @@ -1245,6 +1245,7 @@ mod tests { thread_status: "notLoaded".to_owned(), turn_status: None, goal_status: None, + supervisor: None, git_info: None, created_at: 0, updated_at: 0, diff --git a/wework/e2e/desktop/checkpoints.mjs b/wework/e2e/desktop/checkpoints.mjs index 620c380b21..35ddcef446 100644 --- a/wework/e2e/desktop/checkpoints.mjs +++ b/wework/e2e/desktop/checkpoints.mjs @@ -4,6 +4,7 @@ export const DESKTOP_CHECKPOINTS = [ 'core-task-flow', 'window-lifecycle', 'goal-lifecycle', + 'supervisor-lifecycle', 'resilience', 'conversation-state', 'workspace-attachments', diff --git a/wework/e2e/desktop/task-flow.e2e.mjs b/wework/e2e/desktop/task-flow.e2e.mjs index 1dd7a6d70c..4334a85ee1 100644 --- a/wework/e2e/desktop/task-flow.e2e.mjs +++ b/wework/e2e/desktop/task-flow.e2e.mjs @@ -129,6 +129,14 @@ const GOAL_RESTART_PROMPT = const GOAL_RESTART_INITIAL_TEXT = 'WEWORK_DESKTOP_E2E_GOAL_RESTART_INITIAL_COMPLETE' const GOAL_RESTART_RESUME_PROMPT = 'WEWORK_DESKTOP_E2E_GOAL_RESTART_RESUME' const GOAL_RESTART_COMPLETION_TEXT = 'WEWORK_DESKTOP_E2E_GOAL_RESTART_COMPLETE' +const SUPERVISOR_PROMPT = + 'WEWORK_DESKTOP_E2E_SUPERVISOR: complete this task so supervision can inspect it.' +const SUPERVISOR_COMPLETION_TEXT = 'WEWORK_DESKTOP_E2E_SUPERVISOR_COMPLETE' +const SUPERVISOR_PRINCIPLES = + 'Flag material goal drift and provide the smallest directly actionable correction.' +const SUPERVISOR_CORRECTION = + 'WEWORK_DESKTOP_E2E_SUPERVISOR_CORRECTION: explicitly confirm the original constraint.' +const SUPERVISOR_CORRECTION_COMPLETION_TEXT = 'WEWORK_DESKTOP_E2E_SUPERVISOR_CORRECTION_COMPLETE' const WINDOW_LIFECYCLE_COMPLETION_RESPONSE = [ WINDOW_LIFECYCLE_COMPLETION_TEXT, ...Array.from({ length: 24 }, (_, index) => @@ -4210,7 +4218,7 @@ async function uninstallOfficialPlugin(control, fixture) { await captureVerificationScreenshot(control, 'plugins-05-uninstalled.png') } -async function verifyAutomationLifecycle(control, workspacePath) { +async function ensureExperimentalFeaturesEnabled(control) { const initialSnapshot = JSON.parse(await control.command('snapshot', 'body')) if (!initialSnapshot.testIds.includes('automation-button')) { await control.command('click', '[data-testid="settings-button"]') @@ -4224,7 +4232,11 @@ async function verifyAutomationLifecycle(control, workspacePath) { timeoutMs: DEFAULT_STEP_TIMEOUT_MS, }) } + return initialSnapshot +} +async function verifyAutomationLifecycle(control, workspacePath) { + const initialSnapshot = await ensureExperimentalFeaturesEnabled(control) await control.command('click', '[data-testid="automation-button"]') await control.command('waitFor', '[data-testid="create-automation-button"]', { timeoutMs: WORKBENCH_READY_TIMEOUT_MS, @@ -6623,6 +6635,55 @@ async function verifyActiveGoalIdleUnreadLifecycle({ composerSelector, control, ) } +async function verifyTaskSupervisorLifecycle({ composerSelector, control }) { + await ensureExperimentalFeaturesEnabled(control) + control.setScenario('supervisor') + await control.command('click', '[data-testid="new-chat-button"]') + await control.command('waitFor', composerSelector, { + timeoutMs: WORKBENCH_READY_TIMEOUT_MS, + }) + await selectE2EModel(control, DEFAULT_MODEL_ID, DEFAULT_MODEL_LABEL) + await sendPromptUntilScenarioRequest(control, composerSelector, SUPERVISOR_PROMPT, 'supervisor') + control.releaseSupervisorInitialResponse() + await control.command('waitFor', '[data-testid="message-assistant"]', { + text: SUPERVISOR_COMPLETION_TEXT, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await control.command('waitFor', '[data-testid="task-supervisor-toggle-button"]', { + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await control.command('click', '[data-testid="task-supervisor-toggle-button"]') + await control.command('click', '[data-testid="task-supervisor-mode-auto"]') + await control.command('fill', '[data-testid="task-supervisor-instructions"]', { + value: SUPERVISOR_PRINCIPLES, + }) + await control.command('click', '[data-testid="task-supervisor-save-button"]') + + await withTimeout( + control.awaitScenarioRequestCount('supervisor', 2), + DEFAULT_STEP_TIMEOUT_MS, + 'The supervisor evaluator did not inspect the completed task' + ) + await withTimeout( + control.awaitScenarioRequestCount('supervisor', 3), + DEFAULT_STEP_TIMEOUT_MS, + 'Auto-correction did not send a normal follow-up turn after the task became idle' + ) + await control.command('waitFor', '[data-testid="message-user"]', { + text: SUPERVISOR_CORRECTION, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await control.command('waitFor', '[data-testid="message-assistant"]', { + text: SUPERVISOR_CORRECTION_COMPLETION_TEXT, + timeoutMs: DEFAULT_STEP_TIMEOUT_MS, + }) + await waitForSnapshot( + control, + snapshot => !snapshot.testIds.includes('task-supervisor-suggestion'), + 'Auto-correction incorrectly rendered an approval card' + ) +} + async function verifyGoalRestartRecoveryLifecycle({ composerSelector, control, @@ -7202,6 +7263,9 @@ class DesktopE2EServer { this.goalRestartResumeRelease = new Promise(resolvePromise => { this.releaseGoalRestartResume = resolvePromise }) + this.supervisorInitialRelease = new Promise(resolvePromise => { + this.releaseSupervisorInitial = resolvePromise + }) this.toolBlockNodeOutputObserved = new Promise(resolvePromise => { this.resolveToolBlockNodeOutputObserved = resolvePromise }) @@ -7400,6 +7464,7 @@ class DesktopE2EServer { 'official_plugin', 'automation', 'skill_mention_display', + 'supervisor', ].includes(scenario), `Unknown desktop E2E scenario: ${scenario}` ) @@ -7532,6 +7597,10 @@ class DesktopE2EServer { this.releaseRunningForkFollowUp() } + releaseSupervisorInitialResponse() { + this.releaseSupervisorInitial() + } + releaseGoalIdleInitialResponse() { this.releaseGoalIdleInitial() } @@ -8967,6 +9036,58 @@ class DesktopE2EServer { return } + if (this.scenario === 'supervisor') { + this.recordScenarioRequest('supervisor', modelRequest) + const requestText = JSON.stringify(body) + if (requestText.includes('Current visible progress snapshot (JSON):')) { + assert.ok( + requestText.includes('correction'), + 'The supervisor evaluator request did not include its structured output schema' + ) + assert.ok( + requestText.includes(SUPERVISOR_COMPLETION_TEXT), + 'The supervisor evaluator did not receive the latest assistant progress' + ) + assert.equal( + requestText.includes(SUPERVISOR_PROMPT), + false, + 'The supervisor evaluator received the original user transcript instead of recent AI content' + ) + this.writeSse(response, [ + responseCreated(responseId), + assistantMessage( + JSON.stringify({ + correction: SUPERVISOR_CORRECTION, + rationale: 'The completed reply should restate the original constraint.', + }) + ), + responseCompleted(responseId), + ]) + return + } + if (requestText.includes(SUPERVISOR_CORRECTION)) { + this.writeSse(response, [ + responseCreated(responseId), + assistantMessage(SUPERVISOR_CORRECTION_COMPLETION_TEXT), + responseCompleted(responseId), + ]) + return + } + assert.ok(requestText.includes(SUPERVISOR_PROMPT), 'The supervisor task prompt was lost') + response.writeHead(200, { + 'Access-Control-Allow-Origin': '*', + 'Cache-Control': 'no-cache', + Connection: 'keep-alive', + 'Content-Type': 'text/event-stream; charset=utf-8', + }) + response.write( + createSse([responseCreated(responseId), assistantMessage(SUPERVISOR_COMPLETION_TEXT)]) + ) + await this.supervisorInitialRelease + response.end(createSse([responseCompleted(responseId)])) + return + } + if (this.scenario === 'attachment_only') { this.recordScenarioRequest('attachment_only', modelRequest) const requestText = JSON.stringify(body) @@ -12820,6 +12941,15 @@ last_updated = "2026-07-30T00:00:00Z"` } } + if (shouldRunDesktopCheckpoint('supervisor-lifecycle')) { + phase = 'supervisor-lifecycle' + await verifyTaskSupervisorLifecycle({ composerSelector, control }) + if (shouldStopAfterDesktopCheckpoint('supervisor-lifecycle')) { + console.log(`Wework desktop supervisor-lifecycle checkpoint passed. Evidence: ${resultDir}`) + return + } + } + if (shouldRunDesktopCheckpoint('resilience')) { phase = 'queue-management' await verifyPausedQueueLifecycle({ composerSelector, control }) diff --git a/wework/src-tauri/src/lib.rs b/wework/src-tauri/src/lib.rs index 5a063cb87b..41bb580bcc 100644 --- a/wework/src-tauri/src/lib.rs +++ b/wework/src-tauri/src/lib.rs @@ -590,6 +590,8 @@ struct AppPreferences { #[serde(default)] experimental_features_enabled: bool, #[serde(default)] + supervisor_principles: String, + #[serde(default)] task_completion_notifications_enabled: bool, #[serde(default = "default_true")] tray_unread_enabled: bool, @@ -694,6 +696,7 @@ impl Default for AppPreferences { language: default_language_preference(), terminal_context_injection_enabled: true, experimental_features_enabled: false, + supervisor_principles: String::new(), task_completion_notifications_enabled: false, tray_unread_enabled: true, tray_running_enabled: true, @@ -743,6 +746,7 @@ struct AppPreferencesPatch { language: Option, terminal_context_injection_enabled: Option, experimental_features_enabled: Option, + supervisor_principles: Option, task_completion_notifications_enabled: Option, tray_unread_enabled: Option, tray_running_enabled: Option, @@ -1208,6 +1212,9 @@ fn update_app_preferences( if let Some(value) = patch.experimental_features_enabled { preferences.experimental_features_enabled = value; } + if let Some(value) = patch.supervisor_principles { + preferences.supervisor_principles = value; + } if let Some(value) = patch.task_completion_notifications_enabled { preferences.task_completion_notifications_enabled = value; } @@ -1265,6 +1272,7 @@ struct AppPreferences { language: String, terminal_context_injection_enabled: bool, experimental_features_enabled: bool, + supervisor_principles: String, task_completion_notifications_enabled: bool, tray_unread_enabled: bool, tray_running_enabled: bool, @@ -1289,6 +1297,7 @@ struct AppPreferencesPatch { language: Option, terminal_context_injection_enabled: Option, experimental_features_enabled: Option, + supervisor_principles: Option, task_completion_notifications_enabled: Option, tray_unread_enabled: Option, tray_running_enabled: Option, @@ -1313,6 +1322,7 @@ fn get_app_preferences(_app: tauri::AppHandle) -> Result language: "zh-CN".to_string(), terminal_context_injection_enabled: true, experimental_features_enabled: false, + supervisor_principles: String::new(), task_completion_notifications_enabled: false, tray_unread_enabled: true, tray_running_enabled: true, @@ -1343,6 +1353,7 @@ fn update_app_preferences( .terminal_context_injection_enabled .unwrap_or(true), experimental_features_enabled: patch.experimental_features_enabled.unwrap_or(false), + supervisor_principles: patch.supervisor_principles.unwrap_or_default(), task_completion_notifications_enabled: patch .task_completion_notifications_enabled .unwrap_or(false), diff --git a/wework/src/App.tsx b/wework/src/App.tsx index bbb4663230..690549c6e0 100644 --- a/wework/src/App.tsx +++ b/wework/src/App.tsx @@ -716,6 +716,15 @@ function WeworkDevInstanceBadge() { const [collapsed, setCollapsed] = useState(true) const [position, setPosition] = useState() const draggedRef = useRef(false) + const copiedResetTimeoutRef = useRef(null) + useEffect( + () => () => { + if (copiedResetTimeoutRef.current !== null) { + window.clearTimeout(copiedResetTimeoutRef.current) + } + }, + [] + ) if (!info) return null const rows = getWeworkDevInstanceRows(info) @@ -725,7 +734,13 @@ function WeworkDevInstanceBadge() { const copyValue = async (key: string, value: string) => { await navigator.clipboard?.writeText(value) setCopiedKey(key) - window.setTimeout(() => setCopiedKey(current => (current === key ? null : current)), 1200) + if (copiedResetTimeoutRef.current !== null) { + window.clearTimeout(copiedResetTimeoutRef.current) + } + copiedResetTimeoutRef.current = window.setTimeout(() => { + copiedResetTimeoutRef.current = null + setCopiedKey(current => (current === key ? null : current)) + }, 1200) } const handlePointerDown = (event: ReactPointerEvent) => { diff --git a/wework/src/api/executorAccess.ts b/wework/src/api/executorAccess.ts index 9470271890..176237d6f7 100644 --- a/wework/src/api/executorAccess.ts +++ b/wework/src/api/executorAccess.ts @@ -11,6 +11,11 @@ import type { RuntimeGoalGetResponse, RuntimeGoalSetRequest, RuntimeGoalSetResponse, + RuntimeSupervisorClearRequest, + RuntimeSupervisorGetRequest, + RuntimeSupervisorResolveRequest, + RuntimeSupervisorResponse, + RuntimeSupervisorSetRequest, RuntimeFileChangesRevertRequest, RuntimeFileChangesRevertResponse, RuntimeCompactRequest, @@ -107,6 +112,14 @@ export interface ExecutorRuntimeClient { getRuntimeGoal: (data: RuntimeGoalGetRequest) => Promise setRuntimeGoal: (data: RuntimeGoalSetRequest) => Promise clearRuntimeGoal: (data: RuntimeGoalClearRequest) => Promise + getRuntimeSupervisor: (data: RuntimeSupervisorGetRequest) => Promise + setRuntimeSupervisor: (data: RuntimeSupervisorSetRequest) => Promise + clearRuntimeSupervisor: ( + data: RuntimeSupervisorClearRequest + ) => Promise + resolveRuntimeSupervisor: ( + data: RuntimeSupervisorResolveRequest + ) => Promise openRuntimeWorkspace: (data: RuntimeWorkspaceOpenRequest) => Promise upsertLocalRuntimeProject: ReturnType['upsertLocalRuntimeProject'] renameRuntimeWorkspace: ( diff --git a/wework/src/api/hybrid/hybridServices.ts b/wework/src/api/hybrid/hybridServices.ts index a0b5f4f2b2..061bab932a 100644 --- a/wework/src/api/hybrid/hybridServices.ts +++ b/wework/src/api/hybrid/hybridServices.ts @@ -683,6 +683,18 @@ export function createHybridWorkbenchServices( clearRuntimeGoal(data) { return routeByAddress(data.address).clearRuntimeGoal(data) }, + getRuntimeSupervisor(data) { + return routeByAddress(data.address).getRuntimeSupervisor(data) + }, + setRuntimeSupervisor(data) { + return routeByAddress(data.address).setRuntimeSupervisor(data) + }, + clearRuntimeSupervisor(data) { + return routeByAddress(data.address).clearRuntimeSupervisor(data) + }, + resolveRuntimeSupervisor(data) { + return routeByAddress(data.address).resolveRuntimeSupervisor(data) + }, openRuntimeWorkspace(data: RuntimeWorkspaceOpenRequest) { return runtimeApi(data.deviceId).openRuntimeWorkspace(data) }, @@ -986,6 +998,7 @@ function filterRuntimeChatStreamHandlers( onSubagentActivity: route(handlers.onSubagentActivity), onRuntimeGoalUpdated: route(handlers.onRuntimeGoalUpdated), onRuntimeGoalCleared: route(handlers.onRuntimeGoalCleared), + onRuntimeSupervisorUpdated: route(handlers.onRuntimeSupervisorUpdated), onRuntimeGoalContinuation: route(handlers.onRuntimeGoalContinuation), onRuntimePlanUpdated: route(handlers.onRuntimePlanUpdated), onGuidanceApplied: route(handlers.onGuidanceApplied), diff --git a/wework/src/api/local/localServices.test.ts b/wework/src/api/local/localServices.test.ts index bfc5b0da9a..bdf6614271 100644 --- a/wework/src/api/local/localServices.test.ts +++ b/wework/src/api/local/localServices.test.ts @@ -2140,6 +2140,39 @@ describe('createLocalAppServices', () => { }) }) + test('normalizes local task supervisor requests before IPC', async () => { + const request = vi.fn().mockResolvedValue({ + accepted: true, + taskId: 'task-1', + supervisor: { + mode: 'suggest', + status: 'active', + instructions: 'Keep scope focused', + suggestions: [], + }, + }) + const services = createLocalAppServices({ + ensure: vi.fn().mockResolvedValue({ running: true, ready: true, deviceId: 'device-uuid' }), + request, + subscribe: vi.fn(), + }) + + await services.runtimeWorkApi?.setRuntimeSupervisor({ + address: { deviceId: 'local-device', taskId: 'task-1' }, + mode: 'suggest', + instructions: 'Keep scope focused', + modelId: 'gpt-5.6-luna', + intervalSeconds: 60, + }) + expect(request).toHaveBeenCalledWith('runtime.tasks.supervisor.set', { + address: { deviceId: 'device-uuid', taskId: 'task-1' }, + mode: 'suggest', + instructions: 'Keep scope focused', + modelId: 'gpt-5.6-luna', + intervalSeconds: 60, + }) + }) + test('adapts executor runtime workspace list to workbench shape', async () => { const request = vi.fn().mockResolvedValue({ success: true, diff --git a/wework/src/api/local/localServices.ts b/wework/src/api/local/localServices.ts index 2252e0dcc4..dd8942ee51 100644 --- a/wework/src/api/local/localServices.ts +++ b/wework/src/api/local/localServices.ts @@ -31,6 +31,11 @@ import type { RuntimeGoalGetResponse, RuntimeGoalSetRequest, RuntimeGoalSetResponse, + RuntimeSupervisorClearRequest, + RuntimeSupervisorGetRequest, + RuntimeSupervisorResolveRequest, + RuntimeSupervisorResponse, + RuntimeSupervisorSetRequest, RuntimeGoalStatus, RuntimeTaskAddress, RuntimeTaskArchiveResponse, @@ -2220,6 +2225,22 @@ export function createRuntimeWorkApiFromIpc( clearRuntimeGoal(data: RuntimeGoalClearRequest): Promise { return requestWithLocalDevice('runtime.tasks.goal.clear', data) }, + getRuntimeSupervisor(data: RuntimeSupervisorGetRequest): Promise { + return requestWithLocalDevice('runtime.tasks.supervisor.get', data) + }, + setRuntimeSupervisor(data: RuntimeSupervisorSetRequest): Promise { + return requestWithLocalDevice('runtime.tasks.supervisor.set', data) + }, + clearRuntimeSupervisor( + data: RuntimeSupervisorClearRequest + ): Promise { + return requestWithLocalDevice('runtime.tasks.supervisor.clear', data) + }, + resolveRuntimeSupervisor( + data: RuntimeSupervisorResolveRequest + ): Promise { + return requestWithLocalDevice('runtime.tasks.supervisor.resolve', data) + }, openRuntimeWorkspace(data: RuntimeWorkspaceOpenRequest): Promise { return requestWithLocalDevice('runtime.workspaces.open', data) }, diff --git a/wework/src/api/runtime/runtimeChatStream.ts b/wework/src/api/runtime/runtimeChatStream.ts index 56aa2654a4..09e46c3b4d 100644 --- a/wework/src/api/runtime/runtimeChatStream.ts +++ b/wework/src/api/runtime/runtimeChatStream.ts @@ -188,6 +188,7 @@ function hasLocalExecutorResponseHandlers(handlers: ChatStreamHandlers): boolean handlers.onSubagentActivity || handlers.onRuntimeGoalUpdated || handlers.onRuntimeGoalCleared || + handlers.onRuntimeSupervisorUpdated || handlers.onRuntimePlanUpdated || handlers.onGuidanceApplied || handlers.onRuntimeTransportReplaced diff --git a/wework/src/api/runtimeWork.ts b/wework/src/api/runtimeWork.ts index f01356095e..34abc78518 100644 --- a/wework/src/api/runtimeWork.ts +++ b/wework/src/api/runtimeWork.ts @@ -21,6 +21,11 @@ import type { RuntimeGoalGetResponse, RuntimeGoalSetRequest, RuntimeGoalSetResponse, + RuntimeSupervisorClearRequest, + RuntimeSupervisorGetRequest, + RuntimeSupervisorResolveRequest, + RuntimeSupervisorResponse, + RuntimeSupervisorSetRequest, RuntimeFileChangesRevertRequest, RuntimeFileChangesRevertResponse, RuntimeIMNotificationSettingsResponse, @@ -150,6 +155,22 @@ export function createRuntimeWorkApi(client: HttpClient) { clearRuntimeGoal(data: RuntimeGoalClearRequest): Promise { return client.post('/runtime-work/goal/clear', data) }, + getRuntimeSupervisor(data: RuntimeSupervisorGetRequest): Promise { + return client.post('/runtime-work/supervisor/get', data) + }, + setRuntimeSupervisor(data: RuntimeSupervisorSetRequest): Promise { + return client.post('/runtime-work/supervisor/set', data) + }, + clearRuntimeSupervisor( + data: RuntimeSupervisorClearRequest + ): Promise { + return client.post('/runtime-work/supervisor/clear', data) + }, + resolveRuntimeSupervisor( + data: RuntimeSupervisorResolveRequest + ): Promise { + return client.post('/runtime-work/supervisor/resolve', data) + }, openRuntimeWorkspace(data: RuntimeWorkspaceOpenRequest): Promise { return client.post('/runtime-work/workspaces/open', data) }, diff --git a/wework/src/components/layout/DesktopWorkbenchMain.tsx b/wework/src/components/layout/DesktopWorkbenchMain.tsx index 412761c41f..0e63170d1d 100644 --- a/wework/src/components/layout/DesktopWorkbenchMain.tsx +++ b/wework/src/components/layout/DesktopWorkbenchMain.tsx @@ -6,6 +6,7 @@ import type { AssistantPlanOpenRequest } from '@/components/chat/AssistantPlanCa import { RequestUserInputCard } from '@/components/chat/RequestUserInputCard' import { ScrollableMessageArea } from '@/components/chat/ScrollableMessageArea' import { useExperimentalFeaturesEnabled } from '@/features/experimental-features/useExperimentalFeaturesEnabled' +import { useAppPreferencesState } from '@/features/app-preferences/useAppPreferencesState' import { useWorkbench, useWorkbenchPaneContext } from '@/features/workbench/useWorkbench' import { getPopoutComposerPlaceholder } from '@/features/workbench/popoutWorkspaceContext' import { DeliveryDialog } from '@/features/delivery/DeliveryDialog' @@ -113,8 +114,14 @@ import { useWorkbenchProjectWorkControls } from './useWorkbenchProjectWorkContro import { useRuntimeTaskContinueInIm } from './useRuntimeTaskContinueInIm' import { requestOpenCloudDeviceSettings } from './workbenchShellEvents' import { SubagentStatusIndicator } from './SubagentStatusIndicator' +import { SupervisorSuggestionCards, TaskSupervisorControl } from './TaskSupervisorControl' import { WEWORK_OPEN_TERMINAL_EVENT } from '@/lib/keybindings' -import type { RuntimeAdditionalContext, RuntimeTaskAddress } from '@/types/api' +import type { + RuntimeAdditionalContext, + RuntimeSupervisorMode, + RuntimeSupervisorSuggestion, + RuntimeTaskAddress, +} from '@/types/api' import type { WorkbenchMessage } from '@/types/workbench' import { BufferedChatInput } from './BufferedChatInput' import { DesktopEmptyTaskLauncher } from './DesktopEmptyTaskLauncher' @@ -495,6 +502,7 @@ const DesktopWorkbenchPane = memo(function DesktopWorkbenchPane({ }) { const paneActive = useWorkbenchPaneActive() const experimentalFeaturesEnabled = useExperimentalFeaturesEnabled() + const appPreferences = useAppPreferencesState() const appearanceContext = useOptionalAppearance() const appearance = appearanceContext?.appearance ?? defaultAppearance const background = getWorkbenchBackground(appearance, appearanceContext?.resolvedMode ?? 'light') @@ -553,9 +561,11 @@ const DesktopWorkbenchPane = memo(function DesktopWorkbenchPane({ candidates: ComposerCloudMentionCandidate[] } | null>(null) const runtimeWork = state.runtimeWork - const runtimeTaskTitle = truncateRuntimeTaskTitle( - findRuntimeTask(runtimeWork, currentRuntimeTask)?.title - ) + const runtimeTaskSummary = findRuntimeTask(runtimeWork, currentRuntimeTask) + const runtimeTaskTitle = truncateRuntimeTaskTitle(runtimeTaskSummary?.title) + const supervisor = runtimeTaskSummary?.supervisor ?? null + const currentRuntimeTaskSupportsSupervisor = + runtimeTaskSummary?.runtime?.toLowerCase() === 'codex' const composerCloudProject = currentRuntimeTask ? boundCloudProject : pendingCloudProject const composerTodoItem = currentRuntimeTask ? boundCloudItem : pendingTodoItem const cloudAdditionalContext = useMemo(() => { @@ -609,6 +619,58 @@ const DesktopWorkbenchPane = memo(function DesktopWorkbenchPane({ [cloudAdditionalContext, sendPaneInput] ) + const runtimeWorkApi = services?.runtimeWorkApi + const setTaskSupervisor = useCallback( + async ( + mode: RuntimeSupervisorMode, + instructions: string, + modelId: string | null, + intervalSeconds: number + ) => { + if (!currentRuntimeTask || !runtimeWorkApi) return null + const response = await runtimeWorkApi.setRuntimeSupervisor({ + address: currentRuntimeTask, + mode, + instructions, + modelId, + intervalSeconds, + }) + if (!response.accepted) { + throw new Error(response.error || t('workbench.supervisor_set_failed')) + } + return response.supervisor + }, + [currentRuntimeTask, runtimeWorkApi, t] + ) + + const clearTaskSupervisor = useCallback(async () => { + if (!currentRuntimeTask || !runtimeWorkApi) return + const response = await runtimeWorkApi.clearRuntimeSupervisor({ + address: currentRuntimeTask, + }) + if (!response.accepted) { + throw new Error(response.error || t('workbench.supervisor_clear_failed')) + } + }, [currentRuntimeTask, runtimeWorkApi, t]) + + const resolveTaskSupervisorSuggestion = useCallback( + async (suggestion: RuntimeSupervisorSuggestion, status: 'accepted' | 'dismissed') => { + if (!currentRuntimeTask || !runtimeWorkApi) return + if (status === 'accepted') { + await submitPaneInput(suggestion.message, { guideWhenBusy: true }) + } + const response = await runtimeWorkApi.resolveRuntimeSupervisor({ + address: currentRuntimeTask, + suggestionId: suggestion.id, + status, + }) + if (!response.accepted) { + throw new Error(response.error || t('workbench.supervisor_resolve_failed')) + } + }, + [currentRuntimeTask, runtimeWorkApi, submitPaneInput, t] + ) + useEffect(() => { let active = true if (!currentRuntimeTask) { @@ -1994,7 +2056,7 @@ const DesktopWorkbenchPane = memo(function DesktopWorkbenchPane({ environmentInfoFloatingFooter={ !(forceEnvironmentInfoDocked ?? environmentInfoDocked) && (paneSession.subagentStatuses?.length ?? 0) > 0 ? ( -
+
) : undefined @@ -2428,59 +2490,99 @@ const DesktopWorkbenchPane = memo(function DesktopWorkbenchPane({ } /> ) : ( - + <> + {experimentalFeaturesEnabled && + currentRuntimeTaskSupportsSupervisor && + services?.runtimeWorkApi && ( +
+ + model.isActive !== false && + !model.compatibilityDisabled + )} + onSet={setTaskSupervisor} + onClear={clearTaskSupervisor} + /> +
+ )} + {experimentalFeaturesEnabled && supervisor && ( + + resolveTaskSupervisorSuggestion(suggestion, 'accepted') + } + onDismiss={suggestion => + resolveTaskSupervisorSuggestion(suggestion, 'dismissed') + } + /> + )} + + )} )} @@ -2644,7 +2746,10 @@ const DesktopWorkbenchPane = memo(function DesktopWorkbenchPane({ >
{environmentInfoDocked && hasSubagentStatuses && ( -
+
)} diff --git a/wework/src/components/layout/TaskSupervisorControl.test.tsx b/wework/src/components/layout/TaskSupervisorControl.test.tsx new file mode 100644 index 0000000000..3ce5bf2c56 --- /dev/null +++ b/wework/src/components/layout/TaskSupervisorControl.test.tsx @@ -0,0 +1,127 @@ +import { fireEvent, render, screen, waitFor } from '@testing-library/react' +import { describe, expect, test, vi } from 'vitest' +import { SupervisorSuggestionCards, TaskSupervisorControl } from './TaskSupervisorControl' +import type { RuntimeSupervisorState } from '@/types/api' + +const supervisor: RuntimeSupervisorState = { + mode: 'suggest', + status: 'active', + instructions: 'Keep the task focused', + suggestions: [ + { + id: 'suggestion-1', + message: 'Return to the requested scope.', + rationale: 'The current changes are unrelated.', + status: 'pending', + createdAt: 1, + }, + ], +} + +describe('TaskSupervisorControl', () => { + test('configures a disabled task supervisor independently', async () => { + const onSet = vi.fn().mockResolvedValue(supervisor) + + render() + + fireEvent.click(screen.getByTestId('task-supervisor-toggle-button')) + expect(screen.queryByTestId('task-supervisor-mode-observe')).not.toBeInTheDocument() + fireEvent.click(screen.getByTestId('task-supervisor-mode-auto')) + fireEvent.change(screen.getByTestId('task-supervisor-instructions'), { + target: { value: 'Stop destructive operations' }, + }) + fireEvent.click(screen.getByTestId('task-supervisor-save-button')) + + await waitFor(() => + expect(onSet).toHaveBeenCalledWith('auto', 'Stop destructive operations', null, 30) + ) + }) + + test('selects an independent review model and frequency', async () => { + const onSet = vi.fn().mockResolvedValue(supervisor) + render( + + ) + + fireEvent.click(screen.getByTestId('task-supervisor-toggle-button')) + fireEvent.change(screen.getByTestId('task-supervisor-model'), { + target: { value: 'gpt-5.6-luna' }, + }) + fireEvent.change(screen.getByTestId('task-supervisor-frequency'), { + target: { value: '60' }, + }) + fireEvent.click(screen.getByTestId('task-supervisor-save-button')) + + await waitFor(() => expect(onSet).toHaveBeenCalledWith('suggest', '', 'gpt-5.6-luna', 60)) + }) + + test('opens the configuration panel above the trigger', () => { + render() + + fireEvent.click(screen.getByTestId('task-supervisor-toggle-button')) + + expect(screen.getByTestId('task-supervisor-panel')).toHaveClass( + 'bottom-[calc(100%+0.5rem)]', + 'max-h-[calc(100vh-8rem)]', + 'overflow-y-auto' + ) + }) + + test('distinguishes an active supervisor from an in-progress check', () => { + const { rerender } = render( + + ) + + expect(screen.getByTestId('task-supervisor-toggle-button')).toHaveTextContent( + 'workbench.supervisor_active' + ) + + rerender( + + ) + + expect(screen.getByTestId('task-supervisor-toggle-button')).toHaveTextContent( + 'workbench.supervisor_checking' + ) + }) + + test('shows when the latest check found no correction', () => { + render( + + ) + + expect(screen.getByTestId('task-supervisor-toggle-button')).toHaveTextContent( + 'workbench.supervisor_aligned' + ) + }) + + test('renders pending suggestions and resolves explicit actions', async () => { + const onAccept = vi.fn().mockResolvedValue(undefined) + const onDismiss = vi.fn().mockResolvedValue(undefined) + + render( + + ) + + expect(screen.getByText('Return to the requested scope.')).toBeInTheDocument() + fireEvent.click(screen.getByTestId('task-supervisor-accept-suggestion')) + await waitFor(() => expect(onAccept).toHaveBeenCalledWith(supervisor.suggestions[0])) + }) +}) diff --git a/wework/src/components/layout/TaskSupervisorControl.tsx b/wework/src/components/layout/TaskSupervisorControl.tsx new file mode 100644 index 0000000000..8f0d422981 --- /dev/null +++ b/wework/src/components/layout/TaskSupervisorControl.tsx @@ -0,0 +1,347 @@ +import { AlertCircle, Check, Eye, Loader2, MessageSquareText, ShieldCheck, X } from 'lucide-react' +import { useState } from 'react' +import { Button } from '@/components/ui/button' +import { useTranslation } from '@/hooks/useTranslation' +import { cn } from '@/lib/utils' +import type { + RuntimeSupervisorMode, + RuntimeSupervisorState, + RuntimeSupervisorSuggestion, + UnifiedModel, +} from '@/types/api' + +interface TaskSupervisorControlProps { + supervisor?: RuntimeSupervisorState | null + defaultInstructions?: string + models?: UnifiedModel[] + onSet: ( + mode: RuntimeSupervisorMode, + instructions: string, + modelId: string | null, + intervalSeconds: number + ) => Promise + onClear: () => Promise + className?: string +} + +export function TaskSupervisorControl({ + supervisor, + defaultInstructions = '', + models = [], + onSet, + onClear, + className, +}: TaskSupervisorControlProps) { + const { t } = useTranslation('common') + const [open, setOpen] = useState(false) + const [mode, setMode] = useState(supervisor?.mode ?? 'suggest') + const [instructions, setInstructions] = useState(supervisor?.instructions ?? defaultInstructions) + const [modelId, setModelId] = useState(supervisor?.modelId ?? '') + const [intervalSeconds, setIntervalSeconds] = useState(supervisor?.intervalSeconds ?? 30) + const [saving, setSaving] = useState(false) + const [error, setError] = useState(null) + + const pendingCount = + supervisor?.suggestions.filter(suggestion => suggestion.status === 'pending').length ?? 0 + const latestSuggestionAt = supervisor?.suggestions.reduce( + (latest, suggestion) => Math.max(latest, suggestion.createdAt), + 0 + ) + const latestCheckFoundNoCorrection = + Boolean(supervisor?.lastEvaluatedAt) && + (latestSuggestionAt ?? 0) < (supervisor?.lastEvaluatedAt ?? 0) + const statusLabel = + supervisor?.status === 'checking' + ? t('workbench.supervisor_checking') + : supervisor?.status === 'error' + ? t('workbench.supervisor_error') + : latestCheckFoundNoCorrection + ? t('workbench.supervisor_aligned') + : t('workbench.supervisor_active') + const lastCheckedAt = supervisor?.lastEvaluatedAt + ? new Date(supervisor.lastEvaluatedAt).toLocaleTimeString([], { + hour: '2-digit', + minute: '2-digit', + }) + : null + + const save = async () => { + setSaving(true) + setError(null) + try { + await onSet(mode, instructions, modelId || null, intervalSeconds) + setOpen(false) + } catch (saveError) { + setError(saveError instanceof Error ? saveError.message : String(saveError)) + } finally { + setSaving(false) + } + } + + const disable = async () => { + setSaving(true) + setError(null) + try { + await onClear() + setOpen(false) + } catch (clearError) { + setError(clearError instanceof Error ? clearError.message : String(clearError)) + } finally { + setSaving(false) + } + } + + return ( +
+ + + {open && ( +
+
+
+ {t('workbench.supervisor_title')} +
+

+ {t('workbench.supervisor_description')} +

+ {supervisor && ( +

+ {supervisor.status === 'checking' + ? t('workbench.supervisor_checking_detail') + : lastCheckedAt + ? t('workbench.supervisor_last_checked', { time: lastCheckedAt }) + : t('workbench.supervisor_waiting_first_check')} +

+ )} +
+ + +
+ {(['suggest', 'auto'] as const).map(option => ( + + ))} +
+ +
+ + +
+ + +