From 412896580449e38de427325b41b9e3fbc370ecda Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E7=8E=8B=E6=80=A7=E9=A9=8A?= Date: Sun, 13 Sep 2026 15:56:46 +0800 Subject: [PATCH] fix session --- apps/web/src/App.tsx | 1 + apps/web/src/locales/en.ts | 1 + apps/web/src/locales/zh-TW.ts | 2 +- crates/api/src/computer.rs | 157 ++++++++++++++++++++++++++++++++-- crates/api/src/operations.rs | 140 +++++++++++++++++++++++++++++- crates/api/src/runs.rs | 18 +++- tests/frontend.test.mjs | 1 + 7 files changed, 310 insertions(+), 10 deletions(-) diff --git a/apps/web/src/App.tsx b/apps/web/src/App.tsx index 01cfefe..af9c891 100644 --- a/apps/web/src/App.tsx +++ b/apps/web/src/App.tsx @@ -107,6 +107,7 @@ function localizeError(message:string){ if(message==="AI 回應逾時(150 秒)")return t("aiTimeout"); if(message==="run is not retryable")return t("errorRetryFailed"); if(message==="Stop the task first")return t("takeoverBusy"); + if(message==="Takeover incomplete: an operation is still in flight or has unknown effects; reconcile it before taking control")return t("takeoverInFlight"); return message; } function hudLabel(computer:ComputerStatus,connecting:boolean,handingOff:boolean){ diff --git a/apps/web/src/locales/en.ts b/apps/web/src/locales/en.ts index 6507ace..cf4ff4b 100644 --- a/apps/web/src/locales/en.ts +++ b/apps/web/src/locales/en.ts @@ -346,6 +346,7 @@ export const en: { [K in keyof typeof zhTW]: string } = { sharedPowerHint: "This is a shared computer. The action affects every agent using it.", stopAndTakeOver: "Stop and take over", takeoverBusy: "Stop the task first, then take the mouse.", + takeoverInFlight: "Something is still running on this computer. Wait a moment, then take over again.", hudBooting: "Computer starting…", hudWaking: "Waking…", hudConnecting: "Connecting…", diff --git a/apps/web/src/locales/zh-TW.ts b/apps/web/src/locales/zh-TW.ts index 4d2105f..939978f 100644 --- a/apps/web/src/locales/zh-TW.ts +++ b/apps/web/src/locales/zh-TW.ts @@ -146,7 +146,7 @@ export const zhTW = { mcpNoResults: "找不到符合的 MCP。換個關鍵字,或改用自訂接入。", mcpCustom: "自訂接入", mcpBackToList: "回到列表", mcpConnectNamed: "接入 {name}", mcpKeyHint: "這個 MCP 需要憑證才能連。", toolsCount: "{count} 個工具", disabled: "已關閉", disconnected: "未連線", noTools: "沒有可用工具", reconnect: "重新連線", disable: "停用", enable: "啟用", - preparingDesktop: "正在準備 Agent 的獨立桌面…", computerPreviewHint: "開啟電腦後,畫面會顯示在這裡。", bootingProgress: "啟動中…", restartComputer: "重啟", shutDownComputer: "關閉電腦", computerPower: "電腦選單", sharedPowerHint: "此為共用電腦,操作會影響使用它的所有 Agent。", stopAndTakeOver: "停止並接管", takeoverBusy: "請先停止任務,再接手滑鼠。", + preparingDesktop: "正在準備 Agent 的獨立桌面…", computerPreviewHint: "開啟電腦後,畫面會顯示在這裡。", bootingProgress: "啟動中…", restartComputer: "重啟", shutDownComputer: "關閉電腦", computerPower: "電腦選單", sharedPowerHint: "此為共用電腦,操作會影響使用它的所有 Agent。", stopAndTakeOver: "停止並接管", takeoverBusy: "請先停止任務,再接手滑鼠。", takeoverInFlight: "電腦上還有進行中的操作。請稍候再接手。", hudBooting: "電腦啟動中…", hudWaking: "喚醒中…", hudConnecting: "連線中…", hudHandoff: "換手中…", pasteToRemoteComputer: "貼到遠端電腦", pasteRemoteHelp: "把外面的文字貼在這裡,再送進 VNC。這個方式在區網 HTTP 也能使用。", pasteTextPlaceholder: "在此貼上文字…", pasteIntoVnc: "貼入 VNC", botNamePlaceholder: "例如:研究助理", sharedComputerHint: "與其他機器人共用環境", create: "建立", diff --git a/crates/api/src/computer.rs b/crates/api/src/computer.rs index 35893d7..a1bb8cd 100644 --- a/crates/api/src/computer.rs +++ b/crates/api/src/computer.rs @@ -1110,6 +1110,23 @@ async fn runner_agent_barrier( Ok(()) } +async fn abort_takeover_pause( + state: &AppState, + actor: &Actor, + bot_id: &str, + computer: &ComputerRow, +) { + if computer.state == "running" + && computer.provider_ref.is_some() + && let Err(error) = runner_agent_barrier(state, actor, bot_id, computer, false).await + { + tracing::warn!("takeover rollback resume_agent {bot_id}: {error}"); + } + if let Err(error) = crate::operations::set_agent_barrier(state.pool(), bot_id, false).await { + tracing::warn!("takeover rollback barrier {bot_id}: {error}"); + } +} + pub async fn takeover( state: &AppState, actor: &Actor, @@ -1173,13 +1190,47 @@ pub async fn takeover( crate::operations::set_agent_barrier(state.pool(), bot_id, true) .await .map_err(|error| error.to_string())?; - runner_agent_barrier(state, actor, bot_id, &computer, true).await?; - if crate::operations::agent_has_unresolved_effects(state.pool(), bot_id) - .await - .map_err(|error| error.to_string())? - { - return Err("Takeover incomplete: an operation is still in flight or has unknown effects; reconcile it before taking control".into()); + let admitted = async { + runner_agent_barrier(state, actor, bot_id, &computer, true).await?; + crate::operations::wait_for_accepted_operations(state.pool(), bot_id) + .await + .map_err(|error| error.to_string())?; + if crate::operations::agent_has_inflight_effects(state.pool(), bot_id) + .await + .map_err(|error| error.to_string())? + { + return Err("Takeover incomplete: an operation is still in flight or has unknown effects; reconcile it before taking control".into()); + } + if crate::operations::agent_has_unresolved_effects(state.pool(), bot_id) + .await + .map_err(|error| error.to_string())? + { + tracing::info!( + bot_id, + "takeover granted with historical unknown effects still on the ledger" + ); + } + Ok(()) } + .await; + if let Err(error) = admitted { + abort_takeover_pause(state, actor, bot_id, &computer).await; + return Err(error); + } + let granted = grant_takeover_lease(state, bot_id, &computer_id, &screen).await; + if let Err(error) = granted { + abort_takeover_pause(state, actor, bot_id, &computer).await; + return Err(error); + } + granted +} + +async fn grant_takeover_lease( + state: &AppState, + bot_id: &str, + computer_id: &str, + screen: &ScreenRow, +) -> Result<(String, String), String> { let lease_id = Uuid::new_v4().to_string(); let expires = Utc::now() + TimeDelta::minutes(15); sqlx::query( @@ -1196,7 +1247,7 @@ pub async fn takeover( "UPDATE computers SET control_holder = 'user', control_lease_id = $2, control_lease_expires_at = $3, control_bot_id = $4, updated_at = now() WHERE id = $1", ) - .bind(&computer_id) + .bind(computer_id) .bind(&lease_id) .bind(expires) .bind(bot_id) @@ -3405,3 +3456,95 @@ mod runner_event_tests { assert!(upgraded[1] > moved[1]); } } + +#[cfg(test)] +mod takeover_effect_tests { + use super::*; + + fn app(pool: sqlx::PgPool) -> AppState { + AppState { + db: crate::db::Db { pool }, + sandbox: std::sync::Arc::new(lazyboy_sandbox::FakeSandbox::new()), + data_dir: String::new(), + auth: crate::auth::AuthConfig::from_env(), + memory: crate::memory::MemoryService::from_env(), + mcp: crate::mcp::McpHub::new(), + calls: crate::state::CallRegistry::default(), + wakes: crate::state::WakeBus::default(), + } + } + + async fn seed(pool: &sqlx::PgPool) { + for statement in [ + "INSERT INTO users(id,name) VALUES ('u','test')", + "INSERT INTO spaces(id,user_id,name) VALUES ('s','u','test')", + "INSERT INTO computers(id,space_id,user_id,scope,scope_key,home_key,state,provider_ref) VALUES ('c','s','u','team','team:s','home','running','provider')", + "INSERT INTO bots(id,space_id,user_id,name,computer_id) VALUES ('b','s','u','bot','c')", + ] { + sqlx::query(statement).execute(pool).await.unwrap(); + } + } + + fn actor() -> Actor { + Actor { + user_id: "u".into(), + space_id: "s".into(), + } + } + + async fn barrier_paused(pool: &sqlx::PgPool) -> bool { + sqlx::query_scalar( + "SELECT COALESCE((SELECT paused FROM agent_mutation_barriers WHERE bot_id='b'), false)", + ) + .fetch_one(pool) + .await + .unwrap() + } + + #[sqlx::test(migrations = "../../migrations")] + async fn unknown_history_does_not_block_takeover(pool: sqlx::PgPool) { + seed(&pool).await; + sqlx::query( + "INSERT INTO computer_operations(id,bot_id,payload_hash,status) + VALUES ('b:old','b','hash','unknown')", + ) + .execute(&pool) + .await + .unwrap(); + let state = app(pool.clone()); + let granted = takeover(&state, &actor(), "b").await.unwrap(); + assert!(!granted.0.is_empty()); + assert!(barrier_paused(&pool).await); + let holder: String = + sqlx::query_scalar("SELECT control_holder FROM computers WHERE id='c'") + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(holder, "user"); + } + + #[sqlx::test(migrations = "../../migrations")] + async fn in_flight_operation_refuses_takeover_and_unpauses(pool: sqlx::PgPool) { + seed(&pool).await; + sqlx::query( + "INSERT INTO computer_operations(id,bot_id,payload_hash,status) + VALUES ('b:live','b','hash','accepted')", + ) + .execute(&pool) + .await + .unwrap(); + let state = app(pool.clone()); + let error = takeover(&state, &actor(), "b").await.unwrap_err(); + assert!(error.contains("Takeover incomplete"), "{error}"); + assert!( + !barrier_paused(&pool).await, + "a refused takeover must not leave the agent paused" + ); + let holder: String = + sqlx::query_scalar("SELECT control_holder FROM computers WHERE id='c'") + .fetch_one(&pool) + .await + .unwrap(); + assert_ne!(holder, "user"); + } +} diff --git a/crates/api/src/operations.rs b/crates/api/src/operations.rs index 354c10a..733f77e 100644 --- a/crates/api/src/operations.rs +++ b/crates/api/src/operations.rs @@ -179,6 +179,74 @@ pub async fn agent_has_unresolved_effects( .bind(bot_id).fetch_one(pool).await } +/// Live work that must finish (or be interrupted) before a human lease is +/// granted. Historical `unknown` rows stay on the ledger so they are not +/// retried; they do not mean the agent is still mutating. +pub async fn agent_has_inflight_effects(pool: &PgPool, bot_id: &str) -> Result { + sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM computer_operations WHERE bot_id=$1 AND status='accepted') + OR EXISTS(SELECT 1 FROM package_transitions WHERE bot_id=$1 AND status IN ('pending','cancelling'))", + ) + .bind(bot_id) + .fetch_one(pool) + .await +} + +pub async fn wait_for_accepted_operations(pool: &PgPool, bot_id: &str) -> Result<(), sqlx::Error> { + let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(1); + loop { + let pending: bool = sqlx::query_scalar( + "SELECT EXISTS(SELECT 1 FROM computer_operations WHERE bot_id=$1 AND status='accepted')", + ) + .bind(bot_id) + .fetch_one(pool) + .await?; + if !pending || tokio::time::Instant::now() >= deadline { + return Ok(()); + } + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + } +} + +/// Record a cancelled or timed-out attempt. A row that already finished is +/// left as-is: takeover/halt must not turn a durable success into unknown. +pub async fn close_if_accepted( + pool: &PgPool, + bot_id: &str, + run_id: &str, + operation_id: &str, + error_code: &str, +) -> Result<(), sqlx::Error> { + if operation_id.is_empty() { + return Ok(()); + } + let key = ledger_key(bot_id, operation_id); + let stored = json!({ + "text": "operation interrupted before a result was recorded", + "errorCode": error_code, + "pause": false, + }) + .to_string(); + let mut payload = json!({ + "runId": if run_id.is_empty() { None } else { Some(run_id) }, + "botId": bot_id, + "operationId": operation_id, + "errorCode": error_code, + "status": "unknown", + }); + redact_json(&mut payload); + persist_completion( + pool, + &key, + bot_id, + "unknown", + &stored, + &format!("tool:{key}"), + &payload, + ) + .await +} + fn replay_from_stored(stored: &str) -> Begin { if let Ok(value) = serde_json::from_str::(stored) { Begin::Replay { @@ -298,7 +366,19 @@ async fn persist_completion( if recovered { return transaction.commit().await; } - return Err(sqlx::Error::RowNotFound); + let existing: Option = + sqlx::query_scalar("SELECT status FROM computer_operations WHERE id=$1 AND bot_id=$2") + .bind(key) + .bind(bot_id) + .fetch_optional(&mut *transaction) + .await?; + match existing.as_deref() { + // Interrupted after a durable result, or no intent was recorded. + None | Some("succeeded" | "failed" | "unknown") => { + return transaction.commit().await; + } + _ => return Err(sqlx::Error::RowNotFound), + } } sqlx::query( "INSERT INTO operation_outbox (id, event_key, payload) VALUES ($1,$2,$3) @@ -740,6 +820,11 @@ mod barrier_tests { sqlx::query("INSERT INTO computer_operations(id,bot_id,payload_hash,status) VALUES ('unknown','a','hash','unknown')") .execute(&pool).await.unwrap(); assert!(agent_has_unresolved_effects(&pool, "a").await.unwrap()); + assert_eq!( + agent_has_inflight_effects(&pool, "a").await.unwrap(), + admitted, + "unknown history is not live work" + ); set_agent_barrier(&pool, "a", false).await.unwrap(); assert_eq!( persist_begin(&pool, "after", None, "a", "run", "hash") @@ -754,4 +839,57 @@ mod barrier_tests { Begin::InProgress ); } + + #[sqlx::test(migrations = "../../migrations")] + async fn interrupted_attempt_is_recorded_unknown_and_does_not_clobber_success(pool: PgPool) { + for sql in [ + "INSERT INTO users(id,name) VALUES ('u','test')", + "INSERT INTO spaces(id,user_id,name) VALUES ('s','u','test')", + ] { + sqlx::query(sql).execute(&pool).await.unwrap(); + } + close_if_accepted(&pool, "a", "r", "missing", "UNKNOWN_EFFECT") + .await + .unwrap(); + assert_eq!( + persist_begin(&pool, &ledger_key("a", "live"), None, "a", "r", "hash") + .await + .unwrap(), + Begin::Proceed + ); + assert!(agent_has_inflight_effects(&pool, "a").await.unwrap()); + close_if_accepted(&pool, "a", "r", "live", "TIMEOUT_EFFECT_UNKNOWN") + .await + .unwrap(); + let status: String = + sqlx::query_scalar("SELECT status FROM computer_operations WHERE id=$1") + .bind(ledger_key("a", "live")) + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(status, "unknown"); + assert!(!agent_has_inflight_effects(&pool, "a").await.unwrap()); + assert_eq!( + persist_begin(&pool, &ledger_key("a", "done"), None, "a", "r", "hash") + .await + .unwrap(), + Begin::Proceed + ); + sqlx::query( + "UPDATE computer_operations SET status='succeeded', result='{\"text\":\"done\"}' WHERE id=$1", + ) + .bind(ledger_key("a", "done")) + .execute(&pool) + .await + .unwrap(); + close_if_accepted(&pool, "a", "r", "done", "UNKNOWN_EFFECT") + .await + .unwrap(); + let kept: String = sqlx::query_scalar("SELECT status FROM computer_operations WHERE id=$1") + .bind(ledger_key("a", "done")) + .fetch_one(&pool) + .await + .unwrap(); + assert_eq!(kept, "succeeded"); + } } diff --git a/crates/api/src/runs.rs b/crates/api/src/runs.rs index 5842040..f23429d 100644 --- a/crates/api/src/runs.rs +++ b/crates/api/src/runs.rs @@ -1533,12 +1533,13 @@ async fn execute_run( // effects are unknown and let the model re-observe. Err(_) => { tool_timed_out = true; + close_interrupted_operation(&ctx, "TIMEOUT_EFFECT_UNKNOWN").await; ToolOutcome { text: format!("tool {name} timed out after 150 seconds. Its effects are unknown: observe the current screen or files before anything else, and never repeat a step that already worked."), image: None, pause: false, blocks: Vec::new(), - error_code: None, + error_code: Some("TIMEOUT_EFFECT_UNKNOWN".into()), } } } @@ -2581,6 +2582,20 @@ async fn renew_or_halt( } } +async fn close_interrupted_operation(ctx: &crate::tools::ToolCtx, error_code: &str) { + let Some(operation_id) = ctx.last_operation_id.lock().unwrap().clone() else { + return; + }; + let _ = crate::operations::close_if_accepted( + &ctx.pool, + &ctx.bot_id, + &ctx.run_id, + &operation_id, + error_code, + ) + .await; +} + async fn finish_halt( state: &AppState, thread_id: &str, @@ -2592,6 +2607,7 @@ async fn finish_halt( ctx: &crate::tools::ToolCtx, used_gui: bool, ) -> Result<(), String> { + close_interrupted_operation(ctx, "UNKNOWN_EFFECT").await; match halt { RunHalt::Cancelled => Ok(()), RunHalt::Takeover => { diff --git a/tests/frontend.test.mjs b/tests/frontend.test.mjs index 0059bef..a48032d 100644 --- a/tests/frontend.test.mjs +++ b/tests/frontend.test.mjs @@ -833,6 +833,7 @@ test('taking the screen flips the mouse on the click, not on the reply',()=>{ assert.doesNotMatch(app,/action\(\(\)=>api\(`\/api\/computer\/[^`]*takeover/,'takeover no longer rides the global busy path'); assert.equal((app.match(/\/api\/computer\/\$\{[^}]*\}\/(takeover|release)/g)||[]).length,0,'every handoff goes through setControl'); assert.equal((app.match(/holder==="user"\?"takeover":"release"/g)||[]).length,1,'one request path, one owner of it'); + assert.match(app,/if\(message==="Takeover incomplete: an operation is still in flight or has unknown effects; reconcile it before taking control"\)return t\("takeoverInFlight"\)/); }); test('shared CUA desktops keep the mouse live and still show takeover',()=>{