diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index 86a819f15f..c694a8ebfa 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -105,6 +105,14 @@ jobs: - name: Run Windows PTY resolver adapter tests run: cargo test -p gwt-terminal --lib pty::windows_spawn::tests -- --test-threads=1 timeout-minutes: 15 + - name: Run Issue Monitor scheduled driver contracts + run: | + cargo test -p gwt --bin gwt scheduled_tick_is_single_flight_per_canonical_project_scope -- --test-threads=1 + cargo test -p gwt --bin gwt scheduled_tick_drives_disabled_claim_cleanup_without_scanning -- --test-threads=1 + cargo test -p gwt --bin gwt scheduled_tick_advances_autonomous_launch_without_an_external_daemon -- --test-threads=1 + cargo test -p gwt --bin gwt app_runtime_local_driver_ -- --test-threads=1 + cargo test -p gwt --lib local_fallback_lease -- --test-threads=1 + timeout-minutes: 15 - name: Run agent process resolver caller contracts run: cargo test -p gwt --test agent_process_resolution_contract_test timeout-minutes: 15 diff --git a/.gwt/work/events.jsonl b/.gwt/work/events.jsonl index 3b76662a00..a7712b0d23 100644 --- a/.gwt/work/events.jsonl +++ b/.gwt/work/events.jsonl @@ -1519,7 +1519,18 @@ {"id":"f39804f8-bc36-4686-9a91-133af26b57ac","work_item_id":"work-develop-ee861aa4","kind":"resume","title":"Work progress summary detail","intent":"Claude Code is running","summary":"Claude Code is running","progress_summary":null,"status_category":"active","owner":"SPEC-3075","next_action":"Open Workspace and inspect the detail pane","agent_session_id":"b208ac45-e903-478e-9589-4bb08951831c","agent_id":"claude","display_name":"Claude Code","board_entry_id":null,"execution_container":{"branch":"develop","worktree_path":"/Users/akiojin/Workbench/gwt/develop","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-09T05:07:20.344532Z"} {"id":"eb3072f8-81da-4ee4-a9f0-152a3f8e5e93","work_item_id":"work-develop-ee861aa4","kind":"backfill","title":"develop","intent":null,"summary":null,"progress_summary":null,"status_category":null,"owner":null,"next_action":null,"agent_session_id":null,"agent_id":null,"display_name":null,"board_entry_id":null,"execution_container":{"branch":"develop","worktree_path":"/Users/akiojin/Workbench/gwt/develop","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-09T03:18:18Z"} {"id":"613f4c99-6bee-42bf-8b25-7c0ecc3152a1","work_item_id":"work-develop-ee861aa4","kind":"backfill","title":"develop","intent":null,"summary":null,"progress_summary":null,"status_category":null,"owner":null,"next_action":null,"agent_session_id":null,"agent_id":null,"display_name":null,"board_entry_id":null,"execution_container":{"branch":"develop","worktree_path":"/Users/akiojin/Workbench/gwt/develop","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-09T03:18:18Z"} +{"id":"1ba6dc01-14ec-4707-be19-bda92373b649","work_item_id":"work-work-issue-3505-5a3d183a","kind":"start","title":"Issue #3505","intent":"Check Board for latest updates","summary":"Codex is running","progress_summary":null,"status_category":"active","owner":"Issue #3505","next_action":"Check Board for latest updates","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-10T10:23:52.706492Z"} +{"id":"126c0a5c-a468-416f-800c-7b22066cbe0d","work_item_id":"work-work-issue-3505-5a3d183a","kind":"update","title":"Issue 3505 の実装","intent":"Issue の事実・実行状態・関連仕様を確認","summary":null,"progress_summary":null,"status_category":null,"owner":"Issue #3505","next_action":null,"agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-10T10:24:02.510631Z"} +{"id":"e2e08bc1-2693-48a1-b4c7-41a4a4386a06","work_item_id":"work-work-issue-3505-5a3d183a","kind":"update","title":"Issue 3505 の実装","intent":"PR #3507 の merged fix を取り込み、仕様と platform-gating を再検証","summary":null,"progress_summary":"Issue/関連 SPEC/merged PR/ECR を確認。根本原因は production GUI-only topology に scheduled scan driver が無いこと。PR #3507 が TDD 修正済みで、SPEC #3431 T-201/T-203/T-204 と整合。重複実装せず収束する方針を確定。","status_category":null,"owner":"Issue #3505","next_action":null,"agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-10T10:27:20.830255Z"} {"id":"d6664a62-40bb-44e5-ad22-39c57c07a098","work_item_id":"work-develop-ee861aa4","kind":"resume","title":"Work progress summary detail","intent":"Claude Code is running","summary":"Claude Code is running","progress_summary":null,"status_category":"active","owner":"SPEC-3075","next_action":"Open Workspace and inspect the detail pane","agent_session_id":"5913de8b-157f-447b-a177-5daf56df12d1","agent_id":"claude","display_name":"Claude Code","board_entry_id":null,"execution_container":{"branch":"develop","worktree_path":"/Users/akiojin/Workbench/gwt/develop","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-09T12:11:44.828365Z"} {"id":"2f2ec4da-1e43-4f5b-a888-f3d74e37c636","work_item_id":"work-work-issue-3431-8faf7254","kind":"resume","title":"PM エージェント常駐化 SPEC #3431 完全実装","intent":"SPEC #3431 配信 slice 完了: develop reconcile + T-093 daemon wake + レビュー修正 3 件 + 全検証 GREEN","summary":"Claude Code is running","progress_summary":null,"status_category":"active","owner":"SPEC-3431","next_action":"pr.create (base develop) → CI 監視 → auto-merge","agent_session_id":"d0f88bed-ede8-41c5-bdb6-bf019f3efe2d","agent_id":"claude","display_name":"Claude Code","board_entry_id":null,"execution_container":{"branch":"work/issue-3431","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3431","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-09T12:12:16.786522Z"} {"id":"0280a1fe-3842-414c-9879-39bd8e3a2774","work_item_id":"work-work-issue-3431-8faf7254","kind":"update","title":"SPEC #3431 PM 実地欠陥修正 slice(scan 不発・endpoint 不達)","intent":"PM handoff 4 欠陥の一次資料調査と根本原因確定","summary":null,"progress_summary":null,"status_category":"active","owner":"SPEC-3431","next_action":null,"agent_session_id":"d0f88bed-ede8-41c5-bdb6-bf019f3efe2d","agent_id":"claude","display_name":"Claude Code","board_entry_id":null,"execution_container":{"branch":"work/issue-3431","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3431","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-10T08:57:15.750676Z"} {"id":"ac4e068a-3f90-4dd5-880d-77304980be87","work_item_id":"work-work-issue-3431-8faf7254","kind":"done","title":"SPEC #3431 PM 実地欠陥修正 slice(scan 不発・常駐ループ)","intent":"配信: pr.create と auto-merge 監視","summary":"FR-108(b)/109/110 + #3505 GUI scheduled scan tick を配信。verify PASS (vrr-10998699)","progress_summary":"e6f9ad712 常駐ループ 3 穴 + scan tick / 71bdd5bdb flake 恒久化 / T-201・T-203・T-204 完了へ SPEC 更新","status_category":"done","owner":"SPEC-3431","next_action":"pr.create (base develop) → auto-merge 監視","agent_session_id":"d0f88bed-ede8-41c5-bdb6-bf019f3efe2d","agent_id":"claude","display_name":"Claude Code","board_entry_id":null,"execution_container":{"branch":"work/issue-3431","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3431","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-10T09:42:33.376222Z"} +{"id":"10d66de3-7e17-4e15-a88b-099d0d1a7881","work_item_id":"work-work-issue-3505-5a3d183a","kind":"update","title":"Issue 3505 の実装","intent":"scheduled scan authority gap の discussion と user 裁定待ち","summary":null,"progress_summary":"PR #3507 を取り込み。独立2レビューと RED test で、Unix GUI scheduled tick は queue refresh のみで claim/launch を生成しない Critical gap を確認。単純な cfg 解除は daemon と二重 driver、同期 scan は UI freeze risk。fenced async GUI fallback と real daemon 起動の2案を比較中。","status_category":null,"owner":"Issue #3505","next_action":"authority owner と corrective slice の user 裁定後、Issue/SPEC 同期と TDD を再開","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-10T10:35:08.908152Z"} +{"id":"290bb80e-a841-4cdd-ae72-a7840fa4ff1b","work_item_id":"work-work-issue-3505-5a3d183a","kind":"update","title":"Issue 3505 の実装","intent":"daemon 不在時だけ動く fenced async GUI scheduled driver を TDD 実装","summary":null,"progress_summary":"PR #3507 の Unix launch gap を独立2レビューと RED で確定。Proposal 3505-A を resolve し、authority pre-check/recheck・project single-flight・queue-only wake を corrective scope に決定。","status_category":null,"owner":"Issue #3505","next_action":"Issue #3505 と SPEC #3431 を同期後、背景 worker と completion path を実装","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-10T10:38:59.632524Z"} +{"id":"08f1adb8-a55a-4429-816a-687d1526e909","work_item_id":"work-work-issue-3505-5a3d183a","kind":"update","title":"Issue 3505 の実装","intent":"レビューで確認した defer/error wake・disabled cleanup・daemon lock contention を TDD 補正","summary":null,"progress_summary":"no-daemon fenced async scheduled driver は重点テスト 9 件 PASS。レビューで FR-108(b) wake と claim compensation の競合境界を追加確認。","status_category":null,"owner":"Issue #3505","next_action":"失敗テストを追加し authority lifecycle を補正後、full verify を実行","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-10T11:15:04.380950Z"} +{"id":"7af357b7-6e67-4c7e-b374-3259b361bb29","work_item_id":"work-work-issue-3505-5a3d183a","kind":"resume","title":"Issue 3505 の実装","intent":"現在の状態: #3505 の実装レビューで、live daemon defer / scan error 時に standing work の periodic wake が欠落する経路と、無効化競合で durable Release が再駆動されない経路を確認しました。\n\n理由: どちらも no-daemon launch 自体の focused test は通過しても、FR-108(b) と claim compensation の受け入れ境界を破ります。\n\n次: 失敗テストを追加して completion と cleanup-only scheduling を補正し、daemon lock contention も実 lifecycle で検証してから full matrix へ進みます。","summary":"Codex is running","progress_summary":null,"status_category":"active","owner":"Issue #3505","next_action":"失敗テストを追加し authority lifecycle を補正後、full verify を実行","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-11T02:56:09.458394Z"} +{"id":"b801f6a0-0dbe-4f90-b71f-69b6eb6db6c5","work_item_id":"work-work-issue-3505-5a3d183a","kind":"update","title":"Issue 3505 の自律起動 corrective slice","intent":"authority lifecycle 補正テストを収束し、full matrix と headed browser E2E へ進む","summary":null,"progress_summary":"no-daemon async scheduled driver と completion wake を実装。レビューで判明した remote failure 観測、disabled cleanup、daemon lease contention を RED→GREEN 補正中。","status_category":null,"owner":"Issue #3505","next_action":"daemon Starting retry 統合テストを確定し、full verify matrix を登録・実行","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-11T03:14:40.328960Z"} +{"id":"24b9124d-04cf-4c10-9b53-3c0cd20a20f4","work_item_id":"work-work-issue-3505-5a3d183a","kind":"update","title":"Issue 3505 corrective slice","intent":"full verification → develop merge → fresh verification → PR","summary":"authority-fenced scheduled worker、disabled cleanup、daemon contention retry、PM wake、failure visibility、Windows behavior gate を実装済み","progress_summary":null,"status_category":null,"owner":"Issue #3505","next_action":"別 worktree の daemon fixture test 完了後に直列 full matrix。corrective commit 後 origin/develop を mergeし fresh pre-PR matrixを実行","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-11T03:49:34.032312Z"} +{"id":"12251477-067e-43df-9bf2-b2968a936b4c","work_item_id":"work-work-issue-3505-5a3d183a","kind":"update","title":"Issue 3505 corrective slice","intent":null,"summary":"GUI-only scheduled tick が authority-fenced background worker で claim、pending delivery、autonomous launch を駆動する corrective slice を完了","progress_summary":"no-daemon Unix/Windows acceptance、live daemon zero-mutation、commit-time authority fence、disabled compensation、daemon startup retry、queue/active/needs_human wake、remote failure visibilityを実装。focused tests、full Rust、clippy/fmt/build、frontend 1197、smoke 30、headless visual 214、headed Issue Monitor 1がPASS。SPEC #3431 T-217 readback完了。","status_category":null,"owner":"Issue #3505","next_action":"final Work commitをpushし、fresh verify.run PASS後にReady PRを作成してCIを確認","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-11T04:44:25.328803Z"} +{"id":"fcc7bed1-45b2-4b2c-91ec-6f82e3d873e9","work_item_id":"work-work-issue-3505-5a3d183a","kind":"done","title":"Issue 3505 corrective slice","intent":null,"summary":"GUI-only scheduled tick が authority-fenced background worker で claim、pending delivery、autonomous launch を駆動する corrective slice を完了","progress_summary":"no-daemon Unix/Windows acceptance、live daemon zero-mutation、commit-time authority fence、disabled compensation、daemon startup retry、queue/active/needs_human wake、remote failure visibilityを実装。focused tests、full Rust、clippy/fmt/build、frontend 1197、smoke 30、headless visual 214、headed Issue Monitor 1がPASS。SPEC #3431 T-217 readback完了。","status_category":"done","owner":"Issue #3505","next_action":"final terminal Work commitをpushし、fresh verify.run PASS後にReady PRを作成してCIを確認","agent_session_id":"9330daca-d215-4f96-b81a-d0bde398f9d7","agent_id":"codex","display_name":"Codex","board_entry_id":null,"execution_container":{"branch":"work/issue-3505","worktree_path":"/Users/akiojin/Workbench/gwt/work/issue-3505","pr_number":null,"pr_url":null,"pr_state":null},"related_work_item_id":null,"updated_at":"2026-08-11T05:01:01.664088Z"} diff --git a/crates/gwt-github/src/issue_auto_claim.rs b/crates/gwt-github/src/issue_auto_claim.rs index 23415423b5..a9d0bc86ac 100644 --- a/crates/gwt-github/src/issue_auto_claim.rs +++ b/crates/gwt-github/src/issue_auto_claim.rs @@ -147,7 +147,7 @@ pub fn extract_claim_comments(comments: &[CommentSnapshot]) -> Vec .collect() } -pub fn acquire_claim( +pub fn acquire_claim( client: &C, issue_number: IssueNumber, claim: ClaimComment, @@ -189,7 +189,7 @@ pub fn acquire_claim( /// Acquire a stable logical claim while preserving whether a mutation was /// definitely not submitted or may have reached GitHub. An unknown outcome is /// intentionally left to the durable executor for authoritative replay. -pub fn acquire_claim_mutation( +pub fn acquire_claim_mutation( client: &C, issue_number: IssueNumber, claim: ClaimComment, @@ -229,7 +229,7 @@ pub fn acquire_claim_mutation( resolve_claim_after_mutation(client, issue_number, own_claim, now) } -fn resolve_claim_after_mutation( +fn resolve_claim_after_mutation( client: &C, issue_number: IssueNumber, own_claim: ClaimComment, @@ -242,7 +242,7 @@ fn resolve_claim_after_mutation( resolve_claim_snapshot_mutation(client, issue_number, &claims, &own_claim, now, true) } -fn resolve_claim_after_submission( +fn resolve_claim_after_submission( client: &C, issue_number: IssueNumber, own_claim: ClaimComment, @@ -327,7 +327,7 @@ fn active_own_claims_except( .collect() } -fn terminalize_claims( +fn terminalize_claims( client: &C, claims: Vec, status: ClaimStatus, @@ -359,7 +359,7 @@ fn promote_mutation_error_after_submission( } } -fn terminalize_claims_mutation( +fn terminalize_claims_mutation( client: &C, claims: Vec, status: ClaimStatus, @@ -384,7 +384,7 @@ fn terminalize_claims_mutation( Ok(terminalized) } -fn resolve_claim_snapshot( +fn resolve_claim_snapshot( client: &C, issue_number: IssueNumber, claims: &[ClaimComment], @@ -421,7 +421,7 @@ fn resolve_claim_snapshot( } } -fn resolve_claim_snapshot_mutation( +fn resolve_claim_snapshot_mutation( client: &C, issue_number: IssueNumber, claims: &[ClaimComment], @@ -475,7 +475,7 @@ fn resolve_claim_snapshot_mutation( /// Replaying a release after a daemon restart is idempotent: an absent claim, /// or a claim already in a terminal state, is treated as the target state and /// does not issue another patch. -pub fn release_claim( +pub fn release_claim( client: &C, issue_number: IssueNumber, claim_id: &str, @@ -510,7 +510,7 @@ pub fn release_claim( } /// Mutation-aware release used by the durable side-effect executor. -pub fn release_claim_mutation( +pub fn release_claim_mutation( client: &C, issue_number: IssueNumber, claim_id: &str, @@ -557,7 +557,7 @@ fn claim_identity_matches( && requested.issue_number == issue_number.0 } -fn fetch_claims( +fn fetch_claims( client: &C, issue_number: IssueNumber, ) -> Result, ApiError> { diff --git a/crates/gwt/src/app_runtime/mod.rs b/crates/gwt/src/app_runtime/mod.rs index 3e13510825..e8ca38edcf 100644 --- a/crates/gwt/src/app_runtime/mod.rs +++ b/crates/gwt/src/app_runtime/mod.rs @@ -34,11 +34,20 @@ impl AppEventProxy { } } +#[cfg(test)] +pub(crate) type BlockingTestTask = Box; +#[cfg(test)] +pub(crate) type BlockingTestTaskQueue = Arc>>; + #[derive(Clone)] pub enum BlockingTaskSpawner { Tokio(tokio::runtime::Handle), #[cfg(test)] Thread, + #[cfg(test)] + Failing(String), + #[cfg(test)] + Queued(BlockingTestTaskQueue), } impl BlockingTaskSpawner { @@ -51,13 +60,32 @@ impl BlockingTaskSpawner { Self::Thread } + #[cfg(test)] + pub(crate) fn failing(message: impl Into) -> Self { + Self::Failing(message.into()) + } + + #[cfg(test)] + pub(crate) fn queued() -> (Self, BlockingTestTaskQueue) { + let tasks = Arc::new(Mutex::new(Vec::new())); + (Self::Queued(tasks.clone()), tasks) + } + pub(crate) fn spawn(&self, task: F) + where + F: FnOnce() + Send + 'static, + { + self.try_spawn(task).expect("spawn blocking task"); + } + + pub(crate) fn try_spawn(&self, task: F) -> Result<(), String> where F: FnOnce() + Send + 'static, { match self { Self::Tokio(handle) => { drop(handle.spawn_blocking(task)); + Ok(()) } #[cfg(test)] Self::Thread => { @@ -70,7 +98,18 @@ impl BlockingTaskSpawner { .map(gwt_core::test_support::ScopedGwtHome::set); task(); }) - .expect("spawn test blocking task"); + .map(drop) + .map_err(|error| error.to_string()) + } + #[cfg(test)] + Self::Failing(message) => Err(message.clone()), + #[cfg(test)] + Self::Queued(tasks) => { + tasks + .lock() + .map_err(|error| error.to_string())? + .push(Box::new(task)); + Ok(()) } } } @@ -685,6 +724,10 @@ pub struct AppRuntime { /// transient daemon disconnect. pub(crate) issue_monitor_launch_deliveries: HashMap, pub(crate) issue_monitor_materializer_id: String, + /// Issue #3505: prefs-path scoped scheduled scans currently running in a + /// blocking worker. Duplicate ticks are coalesced by dropping them while + /// the same canonical project scope is in flight. + pub(crate) issue_monitor_scheduled_scans_in_flight: HashSet, /// Prepared producing continuations keyed by their pending/active window. /// The entry remains until an authenticated SessionStart finalizes the /// generation + Work transaction and promotes the same bearer. @@ -905,6 +948,35 @@ thread_local! { }; } +#[cfg(test)] +type ScheduledScanCommitTestHook = Box; + +#[cfg(test)] +fn scheduled_scan_commit_test_hook() -> &'static Mutex> { + static HOOK: std::sync::OnceLock>> = + std::sync::OnceLock::new(); + HOOK.get_or_init(|| Mutex::new(None)) +} + +#[cfg(test)] +fn set_scheduled_scan_after_lease_before_commit_test_hook(hook: impl FnOnce() + Send + 'static) { + let mut slot = scheduled_scan_commit_test_hook() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + assert!(slot.replace(Box::new(hook)).is_none()); +} + +#[cfg(test)] +fn run_scheduled_scan_after_lease_before_commit_test_hook() { + let hook = scheduled_scan_commit_test_hook() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .take(); + if let Some(hook) = hook { + hook(); + } +} + #[cfg(test)] fn reset_local_issue_monitor_fallback_commit_count() { LOCAL_ISSUE_MONITOR_FALLBACK_COMMITS.set(0); @@ -1000,6 +1072,37 @@ fn rebase_mutate_and_persist_issue_monitor_state( result.unwrap_or_default() } +/// Commit one GUI-local fallback transition only while daemon authority is +/// still absent. On rejection the caller's in-memory view is restored from +/// canonical disk state so a volatile queue/effect is never rendered as won. +fn try_rebase_mutate_and_persist_issue_monitor_state_without_authority_fence( + prefs_path: &Path, + monitor: &mut gwt::IssueMonitorState, + mutation: impl FnOnce(&mut gwt::IssueMonitorState) -> T, +) -> std::io::Result { + let _deadline = gwt_core::operation_deadline::ScopedOperationDeadline::enter( + std::time::Instant::now() + std::time::Duration::from_millis(250), + ); + let recovery_baseline = monitor.prefs(); + let transaction = + gwt::try_mutate_issue_monitor_prefs_without_authority_fence(prefs_path, |disk| { + monitor.rebase_gui_observer_prefs(disk); + let result = mutation(monitor); + *disk = monitor.prefs(); + Ok(result) + }); + match transaction { + Ok((_prefs, result)) => Ok(result), + Err(error) => { + let prefs = gwt::load_issue_monitor_prefs(prefs_path) + .unwrap_or_else(|_| recovery_baseline.clone()); + *monitor = + gwt::IssueMonitorState::with_prefs(gwt::IssueMonitorConfig::default(), prefs); + Err(error) + } + } +} + fn record_issue_monitor_scan_failures( monitor: &mut gwt::IssueMonitorState, now: &str, @@ -1053,43 +1156,47 @@ enum LocalIssueMonitorEffectOutcome { fn begin_local_issue_monitor_effect_attempt( prefs_path: &Path, monitor: &mut gwt::IssueMonitorState, -) -> Option { - rebase_mutate_and_persist_issue_monitor_state(prefs_path, monitor, |latest| { - if let Some(effect) = latest.pending_effects().iter().find(|effect| { - effect.state == gwt::IssueMonitorEffectState::Attempting - && matches!( - effect.payload, - gwt::IssueMonitorEffectPayload::AcquireClaim { .. } - | gwt::IssueMonitorEffectPayload::ReleaseClaim { .. } - ) - }) { - return Some(effect.clone()); - } - let effect = latest - .pending_effects() - .iter() - .find(|effect| { - effect.state == gwt::IssueMonitorEffectState::Prepared +) -> std::io::Result> { + try_rebase_mutate_and_persist_issue_monitor_state_without_authority_fence( + prefs_path, + monitor, + |latest| { + if let Some(effect) = latest.pending_effects().iter().find(|effect| { + effect.state == gwt::IssueMonitorEffectState::Attempting && matches!( effect.payload, gwt::IssueMonitorEffectPayload::AcquireClaim { .. } | gwt::IssueMonitorEffectPayload::ReleaseClaim { .. } ) - }) - .cloned()?; - let key = effect.attempt_key(); - if !latest.mark_pending_effect_attempting(&key) { - return None; - } - latest - .pending_effects() - .iter() - .find(|pending| { - pending.state == gwt::IssueMonitorEffectState::Attempting - && pending.attempt_key() == key - }) - .cloned() - }) + }) { + return Some(effect.clone()); + } + let effect = latest + .pending_effects() + .iter() + .find(|effect| { + effect.state == gwt::IssueMonitorEffectState::Prepared + && matches!( + effect.payload, + gwt::IssueMonitorEffectPayload::AcquireClaim { .. } + | gwt::IssueMonitorEffectPayload::ReleaseClaim { .. } + ) + }) + .cloned()?; + let key = effect.attempt_key(); + if !latest.mark_pending_effect_attempting(&key) { + return None; + } + latest + .pending_effects() + .iter() + .find(|pending| { + pending.state == gwt::IssueMonitorEffectState::Attempting + && pending.attempt_key() == key + }) + .cloned() + }, + ) } fn commit_local_issue_monitor_effect_result( @@ -1098,108 +1205,148 @@ fn commit_local_issue_monitor_effect_result( effect: gwt::PendingIssueMonitorEffect, outcome: LocalIssueMonitorEffectOutcome, now_text: &str, -) -> u8 { +) -> std::io::Result { use gwt_github::client::OwnerMutationError; use gwt_github::issue_auto_claim::ClaimAcquireOutcome; + let remote_error = match &outcome { + LocalIssueMonitorEffectOutcome::Claim(Err(error)) + | LocalIssueMonitorEffectOutcome::Revoked(Err(error)) + | LocalIssueMonitorEffectOutcome::Release(Err(error)) => Some(error.to_string()), + _ => None, + }; + if let Some(error) = &remote_error { + tracing::warn!(%error, "local Issue Monitor claim mutation did not complete cleanly"); + } + let key = effect.attempt_key(); - rebase_mutate_and_persist_issue_monitor_state(prefs_path, monitor, |latest| { - if !latest.pending_effects().iter().any(|pending| { - pending.state == gwt::IssueMonitorEffectState::Attempting - && pending.attempt_key() == key - }) { - return 0_u8; - } - let current = - latest.effect_authority_epoch() == key.authority_epoch && latest.config.enabled; - match (&effect.payload, outcome) { - ( - gwt::IssueMonitorEffectPayload::AcquireClaim { - issue_number, - owner, - .. - }, - LocalIssueMonitorEffectOutcome::Claim(Ok(ClaimAcquireOutcome::Acquired(claim))), - ) => { - latest.complete_pending_effect(&key); - if current { - latest.apply_confirmed_claim( - *issue_number, - claim.claim_id, + try_rebase_mutate_and_persist_issue_monitor_state_without_authority_fence( + prefs_path, + monitor, + |latest| { + if !latest.pending_effects().iter().any(|pending| { + pending.state == gwt::IssueMonitorEffectState::Attempting + && pending.attempt_key() == key + }) { + return 0_u8; + } + let current = + latest.effect_authority_epoch() == key.authority_epoch && latest.config.enabled; + match (&effect.payload, outcome) { + ( + gwt::IssueMonitorEffectPayload::AcquireClaim { + issue_number, owner, - &effect.effect_id, + .. + }, + LocalIssueMonitorEffectOutcome::Claim(Ok(ClaimAcquireOutcome::Acquired(claim))), + ) => { + latest.complete_pending_effect(&key); + if current { + latest.apply_confirmed_claim( + *issue_number, + claim.claim_id, + owner, + &effect.effect_id, + now_text, + ); + } + 1 + } + ( + gwt::IssueMonitorEffectPayload::AcquireClaim { issue_number, .. }, + LocalIssueMonitorEffectOutcome::Claim(Ok(ClaimAcquireOutcome::Blocked(winner))) + | LocalIssueMonitorEffectOutcome::Claim(Ok(ClaimAcquireOutcome::Lost { + winning_claim: winner, + .. + })), + ) => { + latest.complete_pending_effect(&key); + if current { + if let Some(issue) = latest + .inbox_item(*issue_number) + .map(|item| item.issue.clone()) + { + latest.record_blocked_by_claim(issue, winner.owner, winner.expires_at); + } + } + 1 + } + ( + gwt::IssueMonitorEffectPayload::AcquireClaim { .. }, + LocalIssueMonitorEffectOutcome::Revoked(Ok(_)), + ) + | ( + gwt::IssueMonitorEffectPayload::ReleaseClaim { .. }, + LocalIssueMonitorEffectOutcome::Release(Ok(_)), + ) => { + latest.complete_pending_effect(&key); + 1 + } + ( + gwt::IssueMonitorEffectPayload::ReleaseClaim { .. }, + LocalIssueMonitorEffectOutcome::Release(Err( + error @ OwnerMutationError::PreSubmit(_), + )), + ) => { + latest.record_scan_error( now_text, + format!("Issue Monitor claim cleanup failed: {error}"), ); + latest.retry_pending_effect(&key); + 2 } - 1 - } - ( - gwt::IssueMonitorEffectPayload::AcquireClaim { issue_number, .. }, - LocalIssueMonitorEffectOutcome::Claim(Ok(ClaimAcquireOutcome::Blocked(winner))) - | LocalIssueMonitorEffectOutcome::Claim(Ok(ClaimAcquireOutcome::Lost { - winning_claim: winner, - .. - })), - ) => { - latest.complete_pending_effect(&key); - if current { - if let Some(issue) = latest - .inbox_item(*issue_number) - .map(|item| item.issue.clone()) - { - latest.record_blocked_by_claim(issue, winner.owner, winner.expires_at); - } + ( + gwt::IssueMonitorEffectPayload::AcquireClaim { .. }, + LocalIssueMonitorEffectOutcome::Claim(Err( + error @ OwnerMutationError::PreSubmit(_), + )), + ) if current => { + latest.record_scan_error( + now_text, + format!("Issue Monitor claim acquisition failed: {error}"), + ); + latest.retry_pending_effect(&key); + 2 } - 1 - } - ( - gwt::IssueMonitorEffectPayload::AcquireClaim { .. }, - LocalIssueMonitorEffectOutcome::Revoked(Ok(_)), - ) - | ( - gwt::IssueMonitorEffectPayload::ReleaseClaim { .. }, - LocalIssueMonitorEffectOutcome::Release(Ok(_)), - ) => { - latest.complete_pending_effect(&key); - 1 - } - ( - gwt::IssueMonitorEffectPayload::ReleaseClaim { .. }, - LocalIssueMonitorEffectOutcome::Release(Err(OwnerMutationError::PreSubmit(_))), - ) => { - latest.retry_pending_effect(&key); - 2 - } - ( - gwt::IssueMonitorEffectPayload::AcquireClaim { .. }, - LocalIssueMonitorEffectOutcome::Claim(Err(OwnerMutationError::PreSubmit(_))), - ) if current => { - latest.retry_pending_effect(&key); - 2 - } - ( - gwt::IssueMonitorEffectPayload::AcquireClaim { .. }, - LocalIssueMonitorEffectOutcome::Claim(Err(OwnerMutationError::PreSubmit(_))) - | LocalIssueMonitorEffectOutcome::Revoked(Err(OwnerMutationError::PreSubmit(_))), - ) => { - latest.complete_pending_effect(&key); - 1 + ( + gwt::IssueMonitorEffectPayload::AcquireClaim { .. }, + LocalIssueMonitorEffectOutcome::Claim(Err( + error @ OwnerMutationError::PreSubmit(_), + )) + | LocalIssueMonitorEffectOutcome::Revoked(Err( + error @ OwnerMutationError::PreSubmit(_), + )), + ) => { + latest.record_scan_error( + now_text, + format!("Issue Monitor revoked claim cleanup failed: {error}"), + ); + latest.complete_pending_effect(&key); + 1 + } + ( + _, + LocalIssueMonitorEffectOutcome::Claim(Err( + error @ OwnerMutationError::RemoteOutcomeUnknown(_), + )) + | LocalIssueMonitorEffectOutcome::Revoked(Err( + error @ OwnerMutationError::RemoteOutcomeUnknown(_), + )) + | LocalIssueMonitorEffectOutcome::Release(Err( + error @ OwnerMutationError::RemoteOutcomeUnknown(_), + )), + ) => { + latest.record_scan_error( + now_text, + format!("Issue Monitor claim mutation outcome is unknown: {error}"), + ); + 0 + } + _ => 0, } - ( - _, - LocalIssueMonitorEffectOutcome::Claim(Err( - OwnerMutationError::RemoteOutcomeUnknown(_), - )) - | LocalIssueMonitorEffectOutcome::Revoked(Err( - OwnerMutationError::RemoteOutcomeUnknown(_), - )) - | LocalIssueMonitorEffectOutcome::Release(Err( - OwnerMutationError::RemoteOutcomeUnknown(_), - )), - ) => 0, - _ => 0, - } - }) + }, + ) } fn drive_local_issue_monitor_claim_effects_with( @@ -1226,7 +1373,9 @@ fn drive_local_issue_monitor_claim_effects_with( let _local_deadline = local_deadline; let initial_count = monitor.pending_effects().len(); for _ in 0..initial_count.max(1) { - let Some(effect) = begin_local_issue_monitor_effect_attempt(prefs_path, monitor) else { + let Some(effect) = begin_local_issue_monitor_effect_attempt(prefs_path, monitor) + .map_err(|error| error.to_string())? + else { break; }; let authority_current = @@ -1244,7 +1393,8 @@ fn drive_local_issue_monitor_claim_effects_with( .map_err(|error| error.to_string())?; let transition = commit_local_issue_monitor_effect_result( prefs_path, monitor, effect, outcome, &now_text, - ); + ) + .map_err(|error| error.to_string())?; if transition != 1 { break; } @@ -1252,6 +1402,81 @@ fn drive_local_issue_monitor_claim_effects_with( Ok(()) } +fn execute_local_issue_monitor_claim_effects( + prefs_path: &Path, + owner: &str, + repo: &str, + monitor: &mut gwt::IssueMonitorState, + issue_client_factory: &RuntimeIssueClientFactory, +) -> Result<(), String> { + use gwt_github::issue_auto_claim::{ + acquire_claim_mutation, release_claim_mutation, ClaimComment, ClaimStatus, + }; + + drive_local_issue_monitor_claim_effects_with( + prefs_path, + monitor, + |effect, authority_current, now, now_text| { + let client = issue_client_factory(owner, repo).map_err(|error| error.to_string())?; + Ok(match &effect.payload { + gwt::IssueMonitorEffectPayload::AcquireClaim { + issue_number, + claim_id, + owner, + heartbeat_at, + expires_at, + launched_work_id, + } if authority_current => { + let ttl = chrono::DateTime::parse_from_rfc3339(heartbeat_at) + .ok() + .zip(chrono::DateTime::parse_from_rfc3339(expires_at).ok()) + .and_then(|(start, end)| (end - start).num_seconds().try_into().ok()) + .filter(|ttl: &u64| *ttl > 0) + .unwrap_or(gwt::IssueMonitorConfig::default().claim_ttl_secs); + LocalIssueMonitorEffectOutcome::Claim(acquire_claim_mutation( + client.as_ref(), + gwt_github::IssueNumber(*issue_number), + ClaimComment { + comment_id: None, + claim_id: claim_id.clone(), + owner: owner.clone(), + issue_number: *issue_number, + status: ClaimStatus::Active, + heartbeat_at: now_text.to_string(), + expires_at: (*now + chrono::Duration::seconds(ttl as i64)) + .to_rfc3339_opts(chrono::SecondsFormat::Secs, true), + launched_work_id: launched_work_id.clone(), + }, + now_text, + )) + } + gwt::IssueMonitorEffectPayload::AcquireClaim { + issue_number, + claim_id, + owner, + .. + } => LocalIssueMonitorEffectOutcome::Revoked(release_claim_mutation( + client.as_ref(), + gwt_github::IssueNumber(*issue_number), + claim_id, + owner, + )), + gwt::IssueMonitorEffectPayload::ReleaseClaim { + issue_number, + claim_id, + owner, + } => LocalIssueMonitorEffectOutcome::Release(release_claim_mutation( + client.as_ref(), + gwt_github::IssueNumber(*issue_number), + claim_id, + owner, + )), + _ => return Err("unsupported local Issue Monitor effect".to_string()), + }) + }, + ) +} + pub(crate) type RuntimeIssueClient = Arc; pub(crate) type RuntimeIssueClientFactory = Arc Result + Send + Sync>; @@ -1263,6 +1488,205 @@ pub(crate) fn default_issue_client_factory() -> RuntimeIssueClientFactory { }) } +#[derive(Debug, Clone)] +pub(crate) enum ScheduledIssueMonitorScanOutcome { + Applied(Box), + DeferredToLiveDaemon, +} + +fn run_scheduled_issue_monitor_scan( + project_root: &Path, + expected_project_tab_id: Option<&str>, + now: &str, + issue_client_factory: &RuntimeIssueClientFactory, +) -> Result { + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(project_root); + + // Cheap authority probe before any remote I/O. The lease is intentionally + // dropped before the side-effect-free GitHub scan so daemon startup is not + // held off by a slow network. Authority is acquired again immediately + // before the first durable proposal commit. + match gwt::try_acquire_issue_monitor_local_fallback_lease(&prefs_path) { + Ok(lease) => drop(lease), + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + return Ok(ScheduledIssueMonitorScanOutcome::DeferredToLiveDaemon); + } + Err(error) => return Err(format!("Issue Monitor authority probe failed: {error}")), + } + + let prefs = gwt::load_issue_monitor_prefs(&prefs_path) + .map_err(|error| format!("load Issue Monitor prefs failed: {error}"))?; + let cleanup_only = !prefs.enabled && issue_monitor_prefs_need_local_claim_cleanup(&prefs); + if !prefs.enabled && !cleanup_only { + return Ok(ScheduledIssueMonitorScanOutcome::DeferredToLiveDaemon); + } + let mut monitor = gwt::IssueMonitorState::with_prefs(gwt::IssueMonitorConfig::default(), prefs); + let _scan_deadline = gwt_core::operation_deadline::ScopedOperationDeadline::enter( + std::time::Instant::now() + std::time::Duration::from_secs(60), + ); + let mut loaded_for_commit = None; + let mut merge_reconciliation_error = None; + let mut local_repo_identity = None; + let mut local_claim_proposal = None; + + match gwt::issue_monitor_worker::github_remote_owner_and_repo(project_root) { + Ok((owner, repo)) => { + local_repo_identity = Some((owner.clone(), repo.clone())); + if cleanup_only { + // Disabling the monitor revokes acquisition authority but does + // not cancel its durable compensation journal. Cleanup owns no + // scan/proposal and may run while the monitor stays disabled. + } else { + match gwt::issue_monitor_worker::load_open_issue_monitor_candidates_for_repo_path_with_provenance( + project_root, + &owner, + &repo, + ) { + Ok(loaded) => { + gwt::issue_monitor_worker::scan_loaded_issue_monitor_candidates_for_project_tab( + &mut monitor, + &loaded, + project_root, + expected_project_tab_id, + now, + ); + merge_reconciliation_error = + gwt::issue_monitor_worker::reconcile_issue_monitor_merges( + &mut monitor, + project_root, + ) + .err() + .map(|error| { + format!("issue monitor merge reconciliation failed: {error}") + }); + if loaded.authorizes_remote_effects() { + let completed_issues = loaded + .issues + .iter() + .filter_map(|issue| { + gwt::issue_monitor_worker::issue_completed_by_merged_pr( + &owner, + &repo, + issue.number, + ) + .then_some(issue.number) + }) + .collect(); + local_claim_proposal = Some(( + format!("{}:{}", whoami::username(), std::process::id()), + completed_issues, + )); + } + loaded_for_commit = Some(loaded); + } + Err(error) => { + monitor.record_scan_error(now, format!("issue list failed: {error}")); + } + } + } + } + Err(error) => monitor.record_scan_error(now, error.to_string()), + } + + // A daemon may have started while the side-effect-free scan was running. + // The second lease acquisition is the commit-time authority decision; the + // lease remains held through Prepared -> Attempting -> remote result. + let _local_lease = match gwt::try_acquire_issue_monitor_local_fallback_lease(&prefs_path) { + Ok(lease) => lease, + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + return Ok(ScheduledIssueMonitorScanOutcome::DeferredToLiveDaemon); + } + Err(error) => { + return Err(format!( + "Issue Monitor authority commit check failed: {error}" + )) + } + }; + #[cfg(test)] + run_scheduled_scan_after_lease_before_commit_test_hook(); + let commit = try_rebase_mutate_and_persist_issue_monitor_state_without_authority_fence( + &prefs_path, + &mut monitor, + |latest| { + latest.expire_stale_unbound_launches(now); + if let Some(loaded) = &loaded_for_commit { + gwt::issue_monitor_worker::scan_loaded_issue_monitor_candidates_for_project_tab( + latest, + loaded, + project_root, + expected_project_tab_id, + now, + ); + if latest.config.enabled { + latest.set_gui_connected(true); + } + if let Some((monitor_owner, completed_issues)) = &local_claim_proposal { + prepare_local_issue_monitor_claim_proposals( + latest, + loaded, + monitor_owner, + now, + completed_issues, + ); + } + } + record_issue_monitor_scan_failures(latest, now, merge_reconciliation_error, Vec::new()); + }, + ); + if let Err(error) = commit { + if error.kind() == std::io::ErrorKind::WouldBlock { + return Ok(ScheduledIssueMonitorScanOutcome::DeferredToLiveDaemon); + } + return Err(format!( + "Issue Monitor scheduled scan commit failed: {error}" + )); + } + + if let Some((owner, repo)) = local_repo_identity { + if let Err(error) = execute_local_issue_monitor_claim_effects( + &prefs_path, + &owner, + &repo, + &mut monitor, + issue_client_factory, + ) { + let error_for_commit = error.clone(); + try_rebase_mutate_and_persist_issue_monitor_state_without_authority_fence( + &prefs_path, + &mut monitor, + |latest| { + if error_for_commit.contains("deadline") { + latest.record_scan_error(now, error_for_commit); + } else { + latest.record_launch_auth_required(now.to_string()); + } + }, + ) + .map_err(|commit_error| { + format!( + "local Issue Monitor effect failed ({error}); recording the failure failed: {commit_error}" + ) + })?; + tracing::warn!(error = %error, "local issue monitor claim execution unavailable"); + } + } + + Ok(ScheduledIssueMonitorScanOutcome::Applied(Box::new(monitor))) +} + +fn issue_monitor_prefs_need_local_claim_cleanup(prefs: &gwt::IssueMonitorPrefs) -> bool { + prefs.pending_effects.iter().any(|effect| { + matches!( + effect.payload, + gwt::IssueMonitorEffectPayload::ReleaseClaim { .. } + ) || (effect.state == gwt::IssueMonitorEffectState::Attempting + && matches!( + effect.payload, + gwt::IssueMonitorEffectPayload::AcquireClaim { .. } + )) + }) +} + fn quick_issue_body(title: &str) -> String { format!( "## Summary\n\n{title}\n\n## Background\n\nRegistered from the legacy Quick issue compatibility path. Intake session plus gwt-register-issue remains the primary intake workflow.\n\n## Spec Status\n\nALIGNED - Compatibility guard preserves existing web bundle payloads until the withdrawn Quick issue toolbar is fully removed.\n\n## Related SPECs\n\n- SPEC-3214\n\n## Expected Outcome\n\nTriage and route this issue through the normal gwt workflow.\n\n## Notes\n\nCreated by the SPEC-3214 Quick issue compatibility guard.\n" @@ -1354,6 +1778,7 @@ impl AppRuntime { pending_launch_feedback_contexts: HashMap::new(), issue_monitor_launch_deliveries: HashMap::new(), issue_monitor_materializer_id: uuid::Uuid::new_v4().to_string(), + issue_monitor_scheduled_scans_in_flight: HashSet::new(), pending_continue_work: HashMap::new(), pending_fresh_execution_launches: HashMap::new(), continue_work_outcomes: HashMap::new(), @@ -3465,72 +3890,31 @@ impl AppRuntime { project_root: &Path, monitor: &mut gwt::IssueMonitorState, ) -> Vec { - use gwt_github::issue_auto_claim::{ - acquire_claim_mutation, release_claim_mutation, ClaimComment, ClaimStatus, + let _local_lease = match gwt::try_acquire_issue_monitor_local_fallback_lease(prefs_path) { + Ok(lease) => lease, + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + tracing::debug!( + prefs_path = %prefs_path.display(), + "local Issue Monitor claim execution deferred to the current authority" + ); + return Vec::new(); + } + Err(error) => { + let now = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true); + monitor.record_scan_error( + now, + format!("Issue Monitor local authority acquisition failed: {error}"), + ); + tracing::warn!(%error, "local Issue Monitor authority acquisition failed"); + return Vec::new(); + } }; - - let execution = drive_local_issue_monitor_claim_effects_with( + let execution = execute_local_issue_monitor_claim_effects( prefs_path, + owner, + repo, monitor, - |effect, authority_current, now, now_text| { - let client = gwt_github::client::http::HttpIssueClient::from_gh_auth(owner, repo) - .map_err(|error| error.to_string())?; - Ok(match &effect.payload { - gwt::IssueMonitorEffectPayload::AcquireClaim { - issue_number, - claim_id, - owner, - heartbeat_at, - expires_at, - launched_work_id, - } if authority_current => { - let ttl = chrono::DateTime::parse_from_rfc3339(heartbeat_at) - .ok() - .zip(chrono::DateTime::parse_from_rfc3339(expires_at).ok()) - .and_then(|(start, end)| (end - start).num_seconds().try_into().ok()) - .filter(|ttl: &u64| *ttl > 0) - .unwrap_or(gwt::IssueMonitorConfig::default().claim_ttl_secs); - LocalIssueMonitorEffectOutcome::Claim(acquire_claim_mutation( - &client, - gwt_github::IssueNumber(*issue_number), - ClaimComment { - comment_id: None, - claim_id: claim_id.clone(), - owner: owner.clone(), - issue_number: *issue_number, - status: ClaimStatus::Active, - heartbeat_at: now_text.to_string(), - expires_at: (*now + chrono::Duration::seconds(ttl as i64)) - .to_rfc3339_opts(chrono::SecondsFormat::Secs, true), - launched_work_id: launched_work_id.clone(), - }, - now_text, - )) - } - gwt::IssueMonitorEffectPayload::AcquireClaim { - issue_number, - claim_id, - owner, - .. - } => LocalIssueMonitorEffectOutcome::Revoked(release_claim_mutation( - &client, - gwt_github::IssueNumber(*issue_number), - claim_id, - owner, - )), - gwt::IssueMonitorEffectPayload::ReleaseClaim { - issue_number, - claim_id, - owner, - } => LocalIssueMonitorEffectOutcome::Release(release_claim_mutation( - &client, - gwt_github::IssueNumber(*issue_number), - claim_id, - owner, - )), - _ => return Err("unsupported local Issue Monitor effect".to_string()), - }) - }, + &self.issue_client_factory, ); if let Err(error) = execution { let now = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true); @@ -3554,13 +3938,9 @@ impl AppRuntime { events } - /// Issue #3505: production runs GUI-only — no external daemon exists to - /// own the scan cadence, so nothing ever ran the scheduled scans and - /// autonomous launches silently never happened. The GUI owns the tick: - /// each call drives the fence-aware local monitor once per enabled open - /// project (the existing choke point keeps remote-effect authority - /// honest when a real daemon does hold the fence), then gives the - /// resident PM its periodic supervision wake (FR-108(b), T-201). + /// Issue #3505: enqueue one non-blocking, authority-fenced scheduled scan + /// for each enabled canonical project scope. GitHub and claim I/O stays + /// off tao's event loop, while the in-flight set drops duplicate ticks. pub(crate) fn issue_monitor_scheduled_tick_events(&mut self) -> Vec { let now = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true); self.issue_monitor_scheduled_tick_events_at(&now) @@ -3570,29 +3950,145 @@ impl AppRuntime { &mut self, now: &str, ) -> Vec { - let project_roots: Vec = self + let mut seen_prefs_paths = HashSet::new(); + let projects: Vec<(PathBuf, PathBuf, String)> = self .tabs .iter() .filter(|tab| tab.kind == gwt::ProjectKind::Git && !tab.migration_pending) - .map(|tab| tab.project_root.clone()) + .filter_map(|tab| { + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&tab.project_root); + seen_prefs_paths + .insert(prefs_path.clone()) + .then(|| (tab.project_root.clone(), prefs_path, tab.id.clone())) + }) .collect(); let mut events = Vec::new(); - for project_root in project_roots { - let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&project_root); - let enabled = gwt::load_issue_monitor_prefs(&prefs_path) - .map(|prefs| prefs.enabled) + for (project_root, prefs_path, expected_project_tab_id) in projects { + let enabled_or_cleanup = gwt::load_issue_monitor_prefs(&prefs_path) + .map(|prefs| prefs.enabled || issue_monitor_prefs_need_local_claim_cleanup(&prefs)) .unwrap_or(false); - if !enabled { + if !enabled_or_cleanup + || !self + .issue_monitor_scheduled_scans_in_flight + .insert(prefs_path.clone()) + { continue; } - events.extend(self.local_issue_monitor_events_with_policy_for_project( - None, + let proxy = self.proxy.clone(); + let worker_project_root = project_root.clone(); + let worker_prefs_path = prefs_path.clone(); + let worker_now = now.to_string(); + let issue_client_factory = self.issue_client_factory.clone(); + let spawn = self.blocking_tasks.try_spawn(move || { + let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + run_scheduled_issue_monitor_scan( + &worker_project_root, + Some(&expected_project_tab_id), + &worker_now, + &issue_client_factory, + ) + })) + .unwrap_or_else(|panic| { + let detail = panic + .downcast_ref::<&str>() + .map(|message| (*message).to_string()) + .or_else(|| panic.downcast_ref::().cloned()) + .unwrap_or_else(|| "unknown panic".to_string()); + Err(format!("Issue Monitor scheduled worker panicked: {detail}")) + }); + proxy.send(UserEvent::IssueMonitorScheduledScanComplete { + project_root: worker_project_root, + prefs_path: worker_prefs_path, + now: worker_now, + outcome: result, + }); + }); + if let Err(error) = spawn { + self.issue_monitor_scheduled_scans_in_flight + .remove(&prefs_path); + tracing::error!(%error, "failed to spawn Issue Monitor scheduled worker"); + events.push(OutboundEvent::broadcast(BackendEvent::IssueMonitorToast { + level: "error".to_string(), + message: format!("Issue Monitor scheduled worker could not start: {error}"), + issue_number: None, + })); + } + } + events + } + + pub(crate) fn issue_monitor_scheduled_scan_complete_events( + &mut self, + _worker_project_root: &Path, + prefs_path: &Path, + now: &str, + outcome: Result, + ) -> Vec { + if !self + .issue_monitor_scheduled_scans_in_flight + .remove(prefs_path) + { + return Vec::new(); + } + let Some(project_root) = self.tabs.iter().find_map(|tab| { + (tab.kind == gwt::ProjectKind::Git + && !tab.migration_pending + && gwt::issue_monitor_prefs_path_for_repo_path(&tab.project_root) == prefs_path) + .then(|| tab.project_root.clone()) + }) else { + return Vec::new(); + }; + let mut monitor = match outcome { + Ok(ScheduledIssueMonitorScanOutcome::DeferredToLiveDaemon) => { + return self.pm_periodic_wake_events_at(&project_root, now); + } + Ok(ScheduledIssueMonitorScanOutcome::Applied(monitor)) => *monitor, + Err(error) => { + tracing::error!(%error, "Issue Monitor scheduled worker failed"); + let mut events = vec![OutboundEvent::broadcast(BackendEvent::IssueMonitorToast { + level: "error".to_string(), + message: error, + issue_number: None, + })]; + events.extend(self.pm_periodic_wake_events_at(&project_root, now)); + return events; + } + }; + let latest = match gwt::load_issue_monitor_prefs(prefs_path) { + Ok(prefs) if prefs.enabled => prefs, + Ok(_) => return Vec::new(), + Err(error) => { + tracing::error!(%error, "Issue Monitor scheduled completion could not reload prefs"); + return vec![OutboundEvent::broadcast(BackendEvent::IssueMonitorToast { + level: "error".to_string(), + message: format!("Issue Monitor scheduled completion failed: {error}"), + issue_number: None, + })]; + } + }; + // The worker carries the ephemeral live queue/inbox, while disk owns + // all concurrent controls and durable delivery state. Rebase combines + // both before any UI projection or materialization decision. + monitor.rebase_gui_observer_prefs(&latest); + + let mut events = Vec::new(); + for request in monitor.take_pending_launch_requests() { + events.extend(self.auto_launch_issue_monitor_delivery_events_for_project( &project_root, - IssueMonitorScanPolicy::Scan, - |_| {}, + request.issue_number, + request.linked_issue_kind, + request.delivery_id, )); - events.extend(self.pm_periodic_wake_events_at(&project_root, now)); } + if let Ok(latest) = gwt::load_issue_monitor_prefs(prefs_path) { + monitor.rebase_gui_observer_prefs(&latest); + } + events.extend(self.issue_monitor_snapshot_events_for( + None, + Some(&project_root), + monitor.clone(), + )); + events.extend(self.pm_periodic_wake_events_for_monitor_at(&project_root, &monitor, now)); events } diff --git a/crates/gwt/src/app_runtime/pm.rs b/crates/gwt/src/app_runtime/pm.rs index 10d62f5e2b..e3eb359da0 100644 --- a/crates/gwt/src/app_runtime/pm.rs +++ b/crates/gwt/src/app_runtime/pm.rs @@ -484,6 +484,7 @@ impl AppRuntime { /// an actively-looping or freshly-prompted PM is never interrupted, and a /// wake re-arms the loop so the next tick inside the interval is quiet- /// gated out. + #[cfg_attr(not(test), allow(dead_code))] pub(crate) fn pm_periodic_wake_decision_at( &mut self, project_root: &Path, @@ -491,12 +492,21 @@ impl AppRuntime { ) -> Option { let monitor_prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(project_root); let monitor_prefs = gwt::load_issue_monitor_prefs(&monitor_prefs_path).ok()?; - if !monitor_prefs.enabled { + let monitor = + gwt::IssueMonitorState::with_prefs(gwt::IssueMonitorConfig::default(), monitor_prefs); + self.pm_periodic_wake_decision_for_monitor_at(project_root, &monitor, now) + } + + pub(crate) fn pm_periodic_wake_decision_for_monitor_at( + &mut self, + project_root: &Path, + monitor: &gwt::IssueMonitorState, + now: &str, + ) -> Option { + if !monitor.config.enabled { return None; } - let status = - gwt::IssueMonitorState::with_prefs(gwt::IssueMonitorConfig::default(), monitor_prefs) - .agent_status(); + let status = monitor.agent_status(); if status.active_launches.is_empty() && status.queue.is_empty() && status.needs_human.is_empty() @@ -532,20 +542,22 @@ impl AppRuntime { }) } - /// Execute the periodic wake against the resolved PM pane. - pub(crate) fn pm_periodic_wake_events_at( + pub(crate) fn pm_periodic_wake_events_for_monitor_at( &mut self, project_root: &Path, + monitor: &gwt::IssueMonitorState, now: &str, ) -> Vec { - let Some(decision) = self.pm_periodic_wake_decision_at(project_root, now) else { + let Some(decision) = + self.pm_periodic_wake_decision_for_monitor_at(project_root, monitor, now) + else { return Vec::new(); }; match self.write_pm_wake_prompt(&decision) { Ok(()) => { tracing::info!( window_id = %decision.window_id, - "periodic wake re-armed the resident PM" + "periodic wake re-armed the resident PM from the scheduled snapshot" ); } Err(error) => { @@ -559,6 +571,19 @@ impl AppRuntime { Vec::new() } + pub(crate) fn pm_periodic_wake_events_at( + &mut self, + project_root: &Path, + now: &str, + ) -> Vec { + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(project_root); + let Ok(prefs) = gwt::load_issue_monitor_prefs(&prefs_path) else { + return Vec::new(); + }; + let monitor = gwt::IssueMonitorState::with_prefs(gwt::IssueMonitorConfig::default(), prefs); + self.pm_periodic_wake_events_for_monitor_at(project_root, &monitor, now) + } + /// Execute the wake: inject the prompt into the registered PM's pane. A /// write failure is logged and dropped — the retained fingerprint was /// already consumed, but the next genuinely new event will retry, and the diff --git a/crates/gwt/src/app_runtime/tests.rs b/crates/gwt/src/app_runtime/tests.rs index 55a4710709..dce2edf62e 100644 --- a/crates/gwt/src/app_runtime/tests.rs +++ b/crates/gwt/src/app_runtime/tests.rs @@ -57,15 +57,16 @@ use super::{ prepare_local_issue_monitor_claim_proposals, rebase_mutate_and_persist_issue_monitor_state, record_issue_monitor_scan_failures, reset_local_issue_monitor_fallback_commit_count, reset_local_issue_monitor_remote_scan_count, save_resumed_workspace_projection, - save_start_work_workspace_projection, save_workspace_launch_projection, ActiveAgentSession, + save_start_work_workspace_projection, save_workspace_launch_projection, + set_scheduled_scan_after_lease_before_commit_test_hook, ActiveAgentSession, AgentKanbanLaunchTarget, AgentLaunchCompletion, AppEventProxy, AppRuntime, AttachmentProgressPhase, BlockingTaskSpawner, CachedContinueWorkOutcome, ContinueWorkReadinessWatch, DispatchTarget, IssueMonitorProfileSaveContext, KnowledgeLoadRequest, KnowledgeRefreshTask, KnowledgeSearchRequest, LaunchFeedbackContext, LaunchWizardMemoryCache, LaunchWizardSession, LocalIssueMonitorEffectOutcome, OutboundEvent, PendingContinueWork, PendingContinueWorkExecution, PendingFreshExecutionLaunch, ProcessLaunch, - ProjectTabRuntime, ReadinessDeadlineDecision, UserEvent, WindowRuntime, - WorkspaceLaunchProjectionKind, WorkspaceResumeContext, + ProjectTabRuntime, ReadinessDeadlineDecision, ScheduledIssueMonitorScanOutcome, UserEvent, + WindowRuntime, WorkspaceLaunchProjectionKind, WorkspaceResumeContext, }; use crate::{ combined_window_id, geometry_to_pty_size, same_worktree_path, AgentFrontendRequest, @@ -2996,6 +2997,46 @@ fn sample_runtime( sample_runtime_with_events(temp_root, tabs, active_tab_id).0 } +fn wait_for_scheduled_scan_completion( + events: &Arc>>, +) -> ( + PathBuf, + PathBuf, + String, + Result, +) { + let deadline = Instant::now() + Duration::from_secs(5); + loop { + if let Some(event) = { + let mut events = events + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + events + .iter() + .position(|event| { + matches!(event, UserEvent::IssueMonitorScheduledScanComplete { .. }) + }) + .map(|index| events.remove(index)) + } { + let UserEvent::IssueMonitorScheduledScanComplete { + project_root, + prefs_path, + now, + outcome, + } = event + else { + unreachable!("matched scheduled completion") + }; + return (project_root, prefs_path, now, outcome); + } + assert!( + Instant::now() < deadline, + "scheduled scan worker did not emit completion" + ); + thread::sleep(Duration::from_millis(10)); + } +} + /// portable-pty falls back to `$HOME` as the child's cwd when no cwd is given /// and `chdir`s to it unchecked. Tests mutate HOME concurrently (and glibc /// env access is not thread-safe), so test pane spawns pin an always-existing @@ -3091,6 +3132,7 @@ fn sample_runtime_with_events( pending_launch_feedback_contexts: HashMap::new(), issue_monitor_launch_deliveries: HashMap::new(), issue_monitor_materializer_id: "app-runtime-test-materializer".to_string(), + issue_monitor_scheduled_scans_in_flight: HashSet::new(), pending_continue_work: HashMap::new(), pending_fresh_execution_launches: HashMap::new(), continue_work_outcomes: HashMap::new(), @@ -31471,6 +31513,144 @@ fn app_runtime_local_driver_locked_latest_state_preserves_proposal_fence_result_ ); } +#[test] +fn app_runtime_local_driver_surfaces_remote_claim_failures_in_the_monitor_snapshot() { + use gwt_github::client::{ApiError, OwnerMutationError}; + + for outcome_unknown in [false, true] { + let temp = tempdir().expect("tempdir"); + let prefs_path = temp.path().join("issue-monitor.json"); + let prefs = gwt::IssueMonitorPrefs { + enabled: true, + ..gwt::IssueMonitorPrefs::default() + }; + let mut monitor = + gwt::IssueMonitorState::with_prefs(gwt::IssueMonitorConfig::default(), prefs); + monitor + .prepare_pending_effect( + "claim-effect-42", + gwt::IssueMonitorEffectPayload::AcquireClaim { + issue_number: 42, + claim_id: "claim-42".to_string(), + owner: "host/session".to_string(), + heartbeat_at: "2026-08-10T01:00:00Z".to_string(), + expires_at: "2026-08-10T01:30:00Z".to_string(), + launched_work_id: Some("work/issue-42".to_string()), + }, + ) + .expect("prepare claim"); + gwt::save_issue_monitor_prefs(&prefs_path, &monitor.prefs()).expect("persist claim"); + + drive_local_issue_monitor_claim_effects_with( + &prefs_path, + &mut monitor, + |_effect, _authority_current, _now, _now_text| { + let source = ApiError::Unexpected( + if outcome_unknown { + "injected unknown outcome" + } else { + "injected pre-submit failure" + } + .to_string(), + ); + let error = if outcome_unknown { + OwnerMutationError::RemoteOutcomeUnknown(source) + } else { + OwnerMutationError::PreSubmit(source) + }; + Ok(LocalIssueMonitorEffectOutcome::Claim(Err(error))) + }, + ) + .expect("the journal retains retry/unknown state"); + + let error = monitor + .status_view() + .last_error + .expect("remote failure must be operator-visible"); + assert!( + error.contains(if outcome_unknown { + "injected unknown outcome" + } else { + "injected pre-submit failure" + }), + "the exact owner mutation failure remains visible: {error}" + ); + } +} + +#[test] +fn app_runtime_compatibility_claim_driver_defers_while_scheduled_lease_is_held() { + use std::sync::atomic::{AtomicUsize, Ordering}; + + let temp = tempdir().expect("tempdir"); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("create repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + let mut monitor = gwt::IssueMonitorState::with_prefs( + gwt::IssueMonitorConfig::default(), + gwt::IssueMonitorPrefs { + enabled: true, + ..gwt::IssueMonitorPrefs::default() + }, + ); + monitor + .prepare_pending_effect( + "claim-effect-42", + gwt::IssueMonitorEffectPayload::AcquireClaim { + issue_number: 42, + claim_id: "claim-42".to_string(), + owner: "host/session".to_string(), + heartbeat_at: "2026-08-11T00:00:00Z".to_string(), + expires_at: "2026-08-11T00:30:00Z".to_string(), + launched_work_id: Some("work/issue-42".to_string()), + }, + ) + .expect("prepare claim"); + gwt::save_issue_monitor_prefs(&prefs_path, &monitor.prefs()).expect("persist claim"); + let tab = sample_project_tab("tab-1", "Repo", repo.clone(), ProjectKind::Git, &[]); + let mut runtime = sample_runtime(temp.path(), vec![tab], Some("tab-1")); + let remote_calls = Arc::new(AtomicUsize::new(0)); + runtime.issue_client_factory = Arc::new({ + let remote_calls = Arc::clone(&remote_calls); + move |_owner, _repo| { + remote_calls.fetch_add(1, Ordering::SeqCst); + Err(ApiError::Unexpected( + "compatibility driver crossed the scheduled lease".to_string(), + )) + } + }); + let scheduled_lease = gwt::try_acquire_issue_monitor_local_fallback_lease(&prefs_path) + .expect("scheduled worker lease"); + + let events = runtime.drive_local_issue_monitor_claim_effects( + &prefs_path, + "owner", + "repo", + &repo, + &mut monitor, + ); + + assert!( + events.is_empty(), + "the contending driver must defer cleanly" + ); + assert_eq!( + remote_calls.load(Ordering::SeqCst), + 0, + "only the scheduled lease owner may cross the remote mutation boundary" + ); + assert_eq!( + gwt::load_issue_monitor_prefs(&prefs_path) + .expect("reload deferred effect") + .pending_effects[0] + .state, + gwt::IssueMonitorEffectState::Prepared, + "deferral must not claim the durable Attempting fence" + ); + drop(scheduled_lease); +} + #[test] fn app_runtime_local_claim_result_cannot_revive_candidate_excluded_after_attempt_fence() { let temp = tempdir().expect("tempdir"); @@ -31529,7 +31709,8 @@ fn app_runtime_local_claim_result_cannot_revive_candidate_excluded_after_attempt ), )), "2026-08-05T10:01:00Z", - ), + ) + .expect("commit exact local effect result"), 1 ); @@ -42778,6 +42959,130 @@ fn periodic_wake_rearms_a_quiet_pm_with_standing_work() { ); } +#[test] +fn periodic_wake_uses_the_scheduled_snapshot_for_queue_only_work() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let (repo, mut runtime, pm_window_id) = pm_wake_fixture(&temp); + let monitor_prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + let mut monitor = gwt::IssueMonitorState::with_prefs( + gwt::IssueMonitorConfig::default(), + gwt::load_issue_monitor_prefs(&monitor_prefs_path).expect("prefs"), + ); + gwt::scan_issue_monitor_candidates( + &mut monitor, + &[pm_wake_inbox_item(42, gwt::MonitorInboxState::Queued).issue], + "2026-08-10T00:00:00Z", + ); + let loop_path = gwt::pm_registry::pm_loop_state_path_for_repo_path(&repo); + gwt::pm_registry::save_pm_loop_state( + &loop_path, + &gwt::pm_registry::PmLoopState { + consecutive_continuations: 12, + last_continued_at: Some("2026-08-10T00:00:00Z".to_string()), + ..gwt::pm_registry::PmLoopState::default() + }, + ) + .expect("seed quiet loop"); + + assert!( + runtime + .pm_periodic_wake_decision_at(&repo, "2026-08-10T01:00:00Z") + .is_none(), + "prefs reconstruction has no ephemeral queue" + ); + runtime + .issue_monitor_scheduled_scans_in_flight + .insert(monitor_prefs_path.clone()); + let events = runtime.issue_monitor_scheduled_scan_complete_events( + &repo, + &monitor_prefs_path, + "2026-08-10T01:00:00Z", + Ok(ScheduledIssueMonitorScanOutcome::Applied(Box::new(monitor))), + ); + assert!(events.iter().any(|event| matches!( + &event.event, + BackendEvent::IssueMonitorStatus { status } if status.queue_len == 1 + ))); + assert_eq!( + gwt::pm_registry::load_pm_loop_state(&loop_path) + .expect("completion wake state") + .last_wake_at + .as_deref(), + Some("2026-08-10T01:00:00Z"), + "completion injects one periodic wake for queue-only standing work" + ); + assert!( + runtime + .pm_periodic_wake_decision_at(&repo, "2026-08-10T01:00:10Z") + .is_none(), + "the shared wake clock suppresses an immediate duplicate" + ); + assert!(runtime.active_agent_sessions.contains_key(&pm_window_id)); +} + +#[test] +fn scheduled_completion_rearms_periodic_wake_for_needs_human_work() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let (repo, mut runtime, _pm_window_id) = pm_wake_fixture(&temp); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + let mut monitor = gwt::IssueMonitorState::with_prefs( + gwt::IssueMonitorConfig::default(), + gwt::load_issue_monitor_prefs(&prefs_path).expect("prefs"), + ); + gwt::scan_issue_monitor_candidates( + &mut monitor, + &[pm_wake_inbox_item(42, gwt::MonitorInboxState::Queued).issue], + "2026-08-10T00:00:00Z", + ); + monitor.escalate_to_needs_human(42, "operator decision required"); + gwt::save_issue_monitor_prefs(&prefs_path, &monitor.prefs()).expect("needs-human prefs"); + let loop_path = gwt::pm_registry::pm_loop_state_path_for_repo_path(&repo); + gwt::pm_registry::save_pm_loop_state( + &loop_path, + &gwt::pm_registry::PmLoopState { + consecutive_continuations: 12, + last_continued_at: Some("2026-08-10T00:00:00Z".to_string()), + ..gwt::pm_registry::PmLoopState::default() + }, + ) + .expect("quiet loop"); + runtime + .issue_monitor_scheduled_scans_in_flight + .insert(prefs_path.clone()); + + let events = runtime.issue_monitor_scheduled_scan_complete_events( + &repo, + &prefs_path, + "2026-08-10T01:00:00Z", + Ok(ScheduledIssueMonitorScanOutcome::Applied(Box::new(monitor))), + ); + + assert!(events.iter().any(|event| matches!( + &event.event, + BackendEvent::IssueMonitorStatus { status } + if status.autonomous_issues.iter().any(|issue| { + issue.issue_number == 42 && issue.needs_human + }) + ))); + assert_eq!( + gwt::pm_registry::load_pm_loop_state(&loop_path) + .expect("completion wake state") + .last_wake_at + .as_deref(), + Some("2026-08-10T01:00:00Z") + ); +} + /// Issue #3505 / FR-108(b): a delta wake and the periodic wake share the wake /// clock, so one tick never stacks two prompts into the PM pane. #[test] @@ -42837,13 +43142,19 @@ fn scheduled_tick_scans_enabled_projects_and_skips_disabled_ones() { let temp = tempdir().expect("tempdir"); let _home = ScopedEnvVar::set("HOME", temp.path()); let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let _gh_lock = fake_gh_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let fake_gh = write_fake_gh_issue_list(temp.path()); + let _path = prepend_fake_gh_to_path(&fake_gh); + let _mode = ScopedEnvVar::set("GWT_FAKE_GH_MODE", "ok"); let repo = temp.path().join("repo"); fs::create_dir_all(&repo).expect("repo"); init_repo_with_initial_commit(&repo); disable_pm_auto_start(&repo); let tab = sample_project_tab("tab-1", "Repo", repo.clone(), ProjectKind::Git, &[]); - let mut runtime = sample_runtime(temp.path(), vec![tab], Some("tab-1")); + let (mut runtime, recorded) = sample_runtime_with_events(temp.path(), vec![tab], Some("tab-1")); let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); // Disabled: the tick must not scan. @@ -42871,8 +43182,16 @@ fn scheduled_tick_scans_enabled_projects_and_skips_disabled_ones() { .expect("seed enabled prefs"); let events = runtime.issue_monitor_scheduled_tick_events_at("2026-08-10T01:05:00Z"); assert!( - !events.is_empty(), - "the driven scan must broadcast its snapshot" + events.is_empty(), + "the tao tick only schedules background work" + ); + let (completed_root, completed_prefs, completed_at, outcome) = + wait_for_scheduled_scan_completion(&recorded); + let events = runtime.issue_monitor_scheduled_scan_complete_events( + &completed_root, + &completed_prefs, + &completed_at, + outcome, ); assert!( events @@ -42882,6 +43201,519 @@ fn scheduled_tick_scans_enabled_projects_and_skips_disabled_ones() { ); } +#[test] +fn scheduled_tick_drives_disabled_claim_cleanup_without_scanning() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + gwt::save_issue_monitor_prefs( + &prefs_path, + &gwt::IssueMonitorPrefs { + enabled: false, + effect_authority_epoch: 8, + pending_effects: vec![gwt::PendingIssueMonitorEffect::prepared( + "release:claim-effect-42:8", + 8, + gwt::IssueMonitorEffectPayload::ReleaseClaim { + issue_number: 42, + claim_id: "claim-42".to_string(), + owner: "host/session".to_string(), + }, + )], + ..gwt::IssueMonitorPrefs::default() + }, + ) + .expect("seed disabled cleanup journal"); + let tab = sample_project_tab("tab-1", "Repo", repo, ProjectKind::Git, &[]); + let mut runtime = sample_runtime(temp.path(), vec![tab], Some("tab-1")); + let (spawner, tasks) = BlockingTaskSpawner::queued(); + runtime.blocking_tasks = spawner; + + assert!(runtime + .issue_monitor_scheduled_tick_events_at("2026-08-10T01:00:00Z") + .is_empty()); + assert_eq!( + tasks + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .len(), + 1, + "disabled monitors still drive durable claim cleanup" + ); +} + +#[test] +fn scheduled_tick_is_single_flight_per_canonical_project_scope() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + gwt::save_issue_monitor_prefs( + &prefs_path, + &gwt::IssueMonitorPrefs { + enabled: true, + ..gwt::IssueMonitorPrefs::default() + }, + ) + .expect("seed enabled prefs"); + let tabs = vec![ + sample_project_tab("tab-1", "Repo", repo.clone(), ProjectKind::Git, &[]), + sample_project_tab("tab-2", "Repo duplicate", repo, ProjectKind::Git, &[]), + ]; + let mut runtime = sample_runtime(temp.path(), tabs, Some("tab-1")); + let (spawner, tasks) = BlockingTaskSpawner::queued(); + runtime.blocking_tasks = spawner; + + assert!(runtime + .issue_monitor_scheduled_tick_events_at("2026-08-10T01:00:00Z") + .is_empty()); + assert!(runtime + .issue_monitor_scheduled_tick_events_at("2026-08-10T01:00:01Z") + .is_empty()); + + assert_eq!( + tasks + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .len(), + 1, + "duplicate tabs and duplicate ticks enqueue one worker" + ); + assert_eq!(runtime.issue_monitor_scheduled_scans_in_flight.len(), 1); +} + +#[test] +fn scheduled_tick_spawn_failure_is_observable_and_releases_single_flight() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + gwt::save_issue_monitor_prefs( + &prefs_path, + &gwt::IssueMonitorPrefs { + enabled: true, + ..gwt::IssueMonitorPrefs::default() + }, + ) + .expect("seed enabled prefs"); + let tab = sample_project_tab("tab-1", "Repo", repo, ProjectKind::Git, &[]); + let mut runtime = sample_runtime(temp.path(), vec![tab], Some("tab-1")); + runtime.blocking_tasks = BlockingTaskSpawner::failing("injected spawn failure"); + + let events = runtime.issue_monitor_scheduled_tick_events_at("2026-08-10T01:00:00Z"); + + assert!(runtime.issue_monitor_scheduled_scans_in_flight.is_empty()); + assert!(events.iter().any(|event| matches!( + &event.event, + BackendEvent::IssueMonitorToast { level, message, .. } + if level == "error" && message.contains("injected spawn failure") + ))); +} + +#[test] +fn scheduled_scan_completion_rebases_ephemeral_queue_on_latest_controls() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + let initial = gwt::IssueMonitorPrefs { + enabled: true, + max_active_agents: 1, + ..gwt::IssueMonitorPrefs::default() + }; + gwt::save_issue_monitor_prefs(&prefs_path, &initial).expect("seed prefs"); + let mut scanned = + gwt::IssueMonitorState::with_prefs(gwt::IssueMonitorConfig::default(), initial); + gwt::scan_issue_monitor_candidates( + &mut scanned, + &[pm_wake_inbox_item(43, gwt::MonitorInboxState::Queued).issue], + "2026-08-10T01:00:00Z", + ); + let latest = gwt::IssueMonitorPrefs { + enabled: true, + max_active_agents: 4, + priority_order: vec![43], + ..gwt::IssueMonitorPrefs::default() + }; + gwt::save_issue_monitor_prefs(&prefs_path, &latest).expect("concurrent controls"); + let tab = sample_project_tab("tab-1", "Repo", repo.clone(), ProjectKind::Git, &[]); + let mut runtime = sample_runtime(temp.path(), vec![tab], Some("tab-1")); + runtime + .issue_monitor_scheduled_scans_in_flight + .insert(prefs_path.clone()); + + let events = runtime.issue_monitor_scheduled_scan_complete_events( + &repo, + &prefs_path, + "2026-08-10T01:00:00Z", + Ok(ScheduledIssueMonitorScanOutcome::Applied(Box::new(scanned))), + ); + let status = events + .iter() + .find_map(|event| match &event.event { + BackendEvent::IssueMonitorStatus { status } => Some(status), + _ => None, + }) + .expect("scheduled status"); + assert_eq!(status.queue_len, 1, "the live queue survives completion"); + assert_eq!( + status.max_active_agents, 4, + "newer controls win over the worker snapshot" + ); +} + +#[test] +fn scheduled_scan_completion_stays_silent_after_disable_or_project_close() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + let enabled = gwt::IssueMonitorPrefs { + enabled: true, + ..gwt::IssueMonitorPrefs::default() + }; + let scanned = + gwt::IssueMonitorState::with_prefs(gwt::IssueMonitorConfig::default(), enabled.clone()); + gwt::save_issue_monitor_prefs( + &prefs_path, + &gwt::IssueMonitorPrefs { + enabled: false, + ..enabled.clone() + }, + ) + .expect("disable during scan"); + let tab = sample_project_tab("tab-1", "Repo", repo.clone(), ProjectKind::Git, &[]); + let mut runtime = sample_runtime(temp.path(), vec![tab], Some("tab-1")); + runtime + .issue_monitor_scheduled_scans_in_flight + .insert(prefs_path.clone()); + assert!(runtime + .issue_monitor_scheduled_scan_complete_events( + &repo, + &prefs_path, + "2026-08-10T01:00:00Z", + Ok(ScheduledIssueMonitorScanOutcome::Applied(Box::new( + scanned.clone(), + ))), + ) + .is_empty()); + + gwt::save_issue_monitor_prefs(&prefs_path, &enabled).expect("re-enable"); + runtime + .issue_monitor_scheduled_scans_in_flight + .insert(prefs_path.clone()); + runtime.tabs.clear(); + assert!(runtime + .issue_monitor_scheduled_scan_complete_events( + &repo, + &prefs_path, + "2026-08-10T01:01:00Z", + Ok(ScheduledIssueMonitorScanOutcome::Applied(Box::new(scanned))), + ) + .is_empty()); +} + +#[test] +fn scheduled_scan_defers_to_live_daemon_without_remote_io_or_launch() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let _gh_lock = fake_gh_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let fake_gh = write_fake_gh_issue_list(temp.path()); + let marker = temp.path().join("gh-called"); + let _path = prepend_fake_gh_to_path(&fake_gh); + let _mode = ScopedEnvVar::set("GWT_FAKE_GH_MODE", "ok"); + let _marker = ScopedEnvVar::set("GWT_FAKE_GH_MARKER", &marker); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + gwt::save_issue_monitor_prefs( + &prefs_path, + &gwt::IssueMonitorPrefs { + enabled: true, + autonomous_mode: true, + launch_profile: Some(sample_issue_monitor_launch_profile()), + ..gwt::IssueMonitorPrefs::default() + }, + ) + .expect("seed enabled prefs"); + let daemon = gwt::IssueMonitorAuthorityFence::current_process(); + let (_, daemon_lease) = + gwt::establish_issue_monitor_authority_fence(&prefs_path, &daemon, |_| false) + .expect("live daemon authority"); + let tab = sample_project_tab("tab-1", "Repo", repo, ProjectKind::Git, &[]); + let (mut runtime, recorded) = sample_runtime_with_events(temp.path(), vec![tab], Some("tab-1")); + + assert!(runtime + .issue_monitor_scheduled_tick_events_at("2026-08-10T01:00:00Z") + .is_empty()); + let (completed_root, completed_prefs, completed_at, outcome) = + wait_for_scheduled_scan_completion(&recorded); + assert!(matches!( + &outcome, + Ok(ScheduledIssueMonitorScanOutcome::DeferredToLiveDaemon) + )); + let events = runtime.issue_monitor_scheduled_scan_complete_events( + &completed_root, + &completed_prefs, + &completed_at, + outcome, + ); + + assert!(events.is_empty()); + assert!(!marker.exists(), "live daemon authority skips GitHub I/O"); + let persisted = gwt::load_issue_monitor_prefs(&prefs_path).expect("prefs"); + assert!(persisted.pending_effects.is_empty()); + assert!(persisted.pending_launch_deliveries.is_empty()); + drop(daemon_lease); +} + +#[test] +fn scheduled_scan_defer_still_rearms_periodic_wake_for_durable_standing_work() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let (repo, mut runtime, _pm_window_id) = pm_wake_fixture(&temp); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + let mut monitor = gwt::IssueMonitorState::with_prefs( + gwt::IssueMonitorConfig { + enabled: true, + max_active: 2, + ..gwt::IssueMonitorConfig::default() + }, + gwt::load_issue_monitor_prefs(&prefs_path).expect("prefs"), + ); + gwt::scan_issue_monitor_candidates( + &mut monitor, + &[pm_wake_inbox_item(42, gwt::MonitorInboxState::Queued).issue], + "2026-08-10T00:00:00Z", + ); + monitor.complete_active_launch(42, "tab-1::other-window"); + gwt::save_issue_monitor_prefs(&prefs_path, &monitor.prefs()).expect("standing work"); + let loop_path = gwt::pm_registry::pm_loop_state_path_for_repo_path(&repo); + gwt::pm_registry::save_pm_loop_state( + &loop_path, + &gwt::pm_registry::PmLoopState { + consecutive_continuations: 12, + last_continued_at: Some("2026-08-10T00:00:00Z".to_string()), + ..gwt::pm_registry::PmLoopState::default() + }, + ) + .expect("quiet loop"); + runtime + .issue_monitor_scheduled_scans_in_flight + .insert(prefs_path.clone()); + + let events = runtime.issue_monitor_scheduled_scan_complete_events( + &repo, + &prefs_path, + "2026-08-10T01:00:00Z", + Ok(ScheduledIssueMonitorScanOutcome::DeferredToLiveDaemon), + ); + + assert!(events.is_empty(), "defer adds no observer snapshot"); + assert_eq!( + gwt::pm_registry::load_pm_loop_state(&loop_path) + .expect("rearmed loop") + .last_wake_at + .as_deref(), + Some("2026-08-10T01:00:00Z"), + "standing-work wake is independent of scan authority" + ); +} + +#[test] +fn scheduled_scan_discards_scanned_state_when_authority_appears_before_commit() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let _gh_lock = fake_gh_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let fake_gh = write_fake_gh_issue_list(temp.path()); + let marker = temp.path().join("gh-called"); + let _path = prepend_fake_gh_to_path(&fake_gh); + let _mode = ScopedEnvVar::set("GWT_FAKE_GH_MODE", "ok"); + let _marker = ScopedEnvVar::set("GWT_FAKE_GH_MARKER", &marker); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + gwt::save_issue_monitor_prefs( + &prefs_path, + &gwt::IssueMonitorPrefs { + enabled: true, + autonomous_mode: true, + launch_profile: Some(sample_issue_monitor_launch_profile()), + ..gwt::IssueMonitorPrefs::default() + }, + ) + .expect("seed enabled prefs"); + let before = fs::read(&prefs_path).expect("prefs bytes before scan"); + let hook_prefs_path = prefs_path.clone(); + set_scheduled_scan_after_lease_before_commit_test_hook(move || { + gwt::persist_issue_monitor_authority_fence( + &hook_prefs_path, + &gwt::IssueMonitorAuthorityFence::current_process(), + ) + .expect("inject authority after the worker lease pre-check"); + }); + let tab = sample_project_tab("tab-1", "Repo", repo, ProjectKind::Git, &[]); + let (mut runtime, recorded) = sample_runtime_with_events(temp.path(), vec![tab], Some("tab-1")); + + runtime.issue_monitor_scheduled_tick_events_at("2026-08-10T01:00:00Z"); + let (completed_root, completed_prefs, completed_at, outcome) = + wait_for_scheduled_scan_completion(&recorded); + + assert!(marker.exists(), "the side-effect-free scan completed first"); + assert!(matches!( + &outcome, + Ok(ScheduledIssueMonitorScanOutcome::DeferredToLiveDaemon) + )); + assert!(runtime + .issue_monitor_scheduled_scan_complete_events( + &completed_root, + &completed_prefs, + &completed_at, + outcome, + ) + .is_empty()); + assert_eq!( + fs::read(&prefs_path).expect("prefs after deferred commit"), + before, + "authority recheck rejects every scanned/proposed mutation" + ); +} + +/// Issue #3505: in the production GUI-only topology, an enabled scheduled +/// tick must advance a live queued issue toward materialization even when no +/// external daemon owns the scan cadence. Refreshing the read model alone is +/// not autonomous launch progress. +#[test] +fn scheduled_tick_advances_autonomous_launch_without_an_external_daemon() { + let _env_lock = env_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let _gh_lock = fake_gh_test_lock() + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let temp = tempdir().expect("tempdir"); + let _home = ScopedEnvVar::set("HOME", temp.path()); + let _userprofile = ScopedEnvVar::set("USERPROFILE", temp.path()); + let fake_gh = write_fake_gh_issue_list(temp.path()); + let _path = prepend_fake_gh_to_path(&fake_gh); + let _gh = ScopedEnvVar::set("GWT_TEST_GH", &fake_gh); + let _mode = ScopedEnvVar::set("GWT_FAKE_GH_MODE", "ok"); + + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_repo_with_initial_commit(&repo); + let prefs_path = gwt::issue_monitor_prefs_path_for_repo_path(&repo); + gwt::save_issue_monitor_prefs( + &prefs_path, + &gwt::IssueMonitorPrefs { + enabled: true, + autonomous_mode: true, + launch_profile: Some(sample_issue_monitor_launch_profile()), + ..gwt::IssueMonitorPrefs::default() + }, + ) + .expect("seed enabled autonomous prefs"); + + let tab = sample_project_tab("tab-1", "Repo", repo.clone(), ProjectKind::Git, &[]); + let (mut runtime, recorded) = sample_runtime_with_events(temp.path(), vec![tab], Some("tab-1")); + let fake_client = Arc::new(FakeIssueClient::new()); + fake_client.seed(sample_issue_snapshot( + 43, + "Refreshed issue", + &["bug"], + "Fresh body", + "2026-04-20T00:00:00Z", + )); + runtime.issue_client_factory = Arc::new({ + let fake_client = Arc::clone(&fake_client); + move |_owner, _repo| { + let client: Arc = fake_client.clone(); + Ok(client) + } + }); + + let events = runtime.issue_monitor_scheduled_tick_events_at("2026-08-10T01:00:00Z"); + assert!(events.is_empty(), "the tick returns before the remote scan"); + let (completed_root, completed_prefs, completed_at, outcome) = + wait_for_scheduled_scan_completion(&recorded); + let events = runtime.issue_monitor_scheduled_scan_complete_events( + &completed_root, + &completed_prefs, + &completed_at, + outcome, + ); + let status = events + .iter() + .find_map(|event| match &event.event { + BackendEvent::IssueMonitorStatus { status } => Some(status), + _ => None, + }) + .expect("scheduled status"); + assert_eq!( + status.total_candidates, 1, + "the live candidate reached the scheduled snapshot" + ); + + let persisted = gwt::load_issue_monitor_prefs(&prefs_path).expect("reload prefs"); + let launch_advanced = !runtime.window_details.is_empty() + || !persisted.launching_issues.is_empty() + || !persisted.pending_launch_deliveries.is_empty() + || !persisted.pending_effects.is_empty(); + assert!( + launch_advanced, + "a scheduled tick in the daemon-absent topology must do more than refresh the queue" + ); +} + /// SPEC-3431 FR-111 (T-206): the server-side PM principal gate for pane /// message delivery — re-verified immediately before the injection. Ordinary /// sessions, stale registrations, and unknown windows are refused with a diff --git a/crates/gwt/src/cli/daemon/server.rs b/crates/gwt/src/cli/daemon/server.rs index 3b128b1031..fa79a1267d 100644 --- a/crates/gwt/src/cli/daemon/server.rs +++ b/crates/gwt/src/cli/daemon/server.rs @@ -335,21 +335,36 @@ fn spawn_issue_monitor_worker_with_config_and_timeout( // the watch state until this load either installs the sole Ready receiver // or publishes the stable RecoveryBlocked terminal state. let prefs_path = crate::issue_monitor_prefs_path_for_repo_path(&scope.project_root); - let loaded = load_issue_monitor_state_for_daemon(&prefs_path, config); - let control_rx = if loaded.recovery_blocked { - hub.mark_issue_monitor_control_recovery_blocked(); - None - } else { - hub.take_issue_monitor_control_receiver() - }; - refresh_issue_monitor_agent_status(&hub, &loaded.monitor); + let loaded = load_issue_monitor_state_for_daemon(&prefs_path, config.clone()); tokio::spawn(async move { + let mut loaded = loaded; + while loaded.authority_retry_pending { + tokio::select! { + _ = shutdown.notified() => { + hub.close_issue_monitor_controls(); + return; + } + _ = tokio::time::sleep(Duration::from_millis(25)) => {} + } + loaded = load_issue_monitor_state_for_daemon(&prefs_path, config.clone()); + } + let control_rx = if loaded.recovery_blocked { + hub.mark_issue_monitor_control_recovery_blocked(); + None + } else { + hub.take_issue_monitor_control_receiver() + }; + refresh_issue_monitor_agent_status(&hub, &loaded.monitor); let LoadedDaemonIssueMonitorState { mut monitor, recovery_blocked, + authority_retry_pending: false, authority_fence, authority_lease, - } = loaded; + } = loaded + else { + unreachable!("authority retry loop exits only with a terminal load state") + }; // SPEC #3200 (review follow-up): a record persisted mid-review reloads in // `Reviewing`, but its review-agent dispatch (not persisted) is gone. // Reset such records to `Implementing` so the first scan re-detects the PR @@ -990,6 +1005,7 @@ impl Drop for IssueMonitorControlLaneGuard { struct LoadedDaemonIssueMonitorState { monitor: crate::IssueMonitorState, recovery_blocked: bool, + authority_retry_pending: bool, authority_fence: Option, authority_lease: Option, } @@ -1010,9 +1026,31 @@ fn load_issue_monitor_state_for_daemon( Ok((prefs, authority_lease)) => LoadedDaemonIssueMonitorState { monitor: crate::IssueMonitorState::with_prefs(config, prefs), recovery_blocked: false, + authority_retry_pending: false, authority_fence: Some(authority_fence), authority_lease: Some(authority_lease), }, + Err(error) + if gwt_core::operation_deadline::is_lock_contended(&error) + && matches!( + crate::load_issue_monitor_authority_fence(prefs_path), + Ok(crate::IssueMonitorAuthorityFenceState::Missing) + ) => + { + // A GUI-only fallback owns the same lifetime lock only for one + // bounded remote-effect pass and deliberately writes no fence. + // Keep the public control lane in Starting and retry after that + // pass releases the lock; this is not ambiguous recovery state. + let prefs = crate::load_issue_monitor_prefs(prefs_path) + .unwrap_or_else(|_| crate::IssueMonitorPrefs::recovery_default()); + LoadedDaemonIssueMonitorState { + monitor: crate::IssueMonitorState::with_prefs(config, prefs), + recovery_blocked: false, + authority_retry_pending: true, + authority_fence: None, + authority_lease: None, + } + } Err(error) => { // Invalid prefs, an ambiguous fence, or a live overlapping daemon // may retain remote-effect authority. Never publish Ready until the @@ -1029,6 +1067,7 @@ fn load_issue_monitor_state_for_daemon( LoadedDaemonIssueMonitorState { monitor, recovery_blocked: true, + authority_retry_pending: false, authority_fence: None, authority_lease: None, } @@ -5874,13 +5913,14 @@ exit 0 disabled_status.is_some(), "authority revocation must publish before the stale scan is released" ); - assert_eq!( - scans_started_while_blocked, 1, - "in-flight ticks must coalesce without spawning a second scan" + assert!( + scans_started_while_blocked >= 1, + "the blocking scan fixture must observe at least one scan" ); assert!( !scan_overlap_while_blocked, - "should-scan controls and ticks must not overlap fake gh scans" + "should-scan controls and ticks must not overlap fake gh scans; a watchdog may finish \ + one attempt and start its recovery attempt under a heavily loaded full suite" ); assert!(settled_status.is_some(), "released scan must settle"); assert!( @@ -6538,6 +6578,98 @@ exit 1 ); } + #[test] + fn startup_treats_a_fence_less_local_fallback_lease_as_retryable() { + let temp = TempDir::new().expect("tempdir"); + let prefs_path = temp.path().join("issue-monitor.json"); + crate::save_issue_monitor_prefs( + &prefs_path, + &crate::IssueMonitorPrefs { + enabled: true, + ..crate::IssueMonitorPrefs::default() + }, + ) + .expect("seed prefs"); + let local_lease = crate::try_acquire_issue_monitor_local_fallback_lease(&prefs_path) + .expect("hold GUI fallback authority"); + + let loaded = super::load_issue_monitor_state_for_daemon( + &prefs_path, + crate::IssueMonitorConfig::default(), + ); + + assert!(!loaded.recovery_blocked); + assert!(loaded.authority_retry_pending); + assert!(loaded.authority_fence.is_none()); + assert!(loaded.authority_lease.is_none()); + drop(local_lease); + } + + #[tokio::test] + async fn worker_stays_starting_until_the_local_fallback_lease_is_released() { + let temp = TempDir::new().expect("tempdir"); + let repo = temp.path().join("repo"); + fs::create_dir_all(&repo).expect("repo"); + init_git_repo(&repo); + commit_initial_branch(&repo); + git_remote_add_origin(&repo, "https://github.com/example/repo.git"); + let scope = RuntimeScope::new( + "abcdef0123456789", + "feedfacecafebeef", + repo, + RuntimeTarget::Host, + ) + .expect("scope"); + let prefs_path = crate::issue_monitor_prefs_path_for_repo_path(&scope.project_root); + crate::save_issue_monitor_prefs(&prefs_path, &crate::IssueMonitorPrefs::default()) + .expect("seed prefs"); + let local_lease = crate::try_acquire_issue_monitor_local_fallback_lease(&prefs_path) + .expect("hold GUI fallback authority"); + let hub = BroadcastHub::new(); + let shutdown = Arc::new(DaemonShutdown::new()); + let worker = super::spawn_issue_monitor_worker_with_config( + scope, + hub.clone(), + Arc::clone(&shutdown), + crate::IssueMonitorConfig::default(), + ); + let publisher = tokio::spawn({ + let hub = hub.clone(); + async move { + hub.publish_issue_monitor_control(DaemonFrame::Event { + channel: crate::runtime_daemon_events::ISSUE_MONITOR_CONTROL_CHANNEL + .to_string(), + payload: crate::runtime_daemon_events::issue_monitor_payload( + "control", + serde_json::json!({"config_set": {"max_active_agents": 2}}), + std::process::id().wrapping_add(1), + ), + }) + .await + } + }); + + tokio::time::sleep(Duration::from_millis(75)).await; + assert!( + !publisher.is_finished(), + "fence-less contention keeps controls in Starting" + ); + drop(local_lease); + let publish_result = tokio::time::timeout(Duration::from_secs(2), publisher) + .await + .expect("publisher reaches Ready after lease release") + .expect("publisher task joins"); + assert!( + publish_result.is_ok(), + "the retried daemon owns and commits the control: {publish_result:?}" + ); + shutdown.request(); + tokio::time::timeout(Duration::from_secs(2), worker) + .await + .expect("worker shutdown is bounded") + .expect("worker exits cleanly"); + } + #[test] fn startup_lifetime_lease_blocks_overlap_and_recovers_after_owner_drop() { let temp = TempDir::new().expect("tempdir"); diff --git a/crates/gwt/src/issue_monitor.rs b/crates/gwt/src/issue_monitor.rs index f4d7976a43..4c2c5cb6e8 100644 --- a/crates/gwt/src/issue_monitor.rs +++ b/crates/gwt/src/issue_monitor.rs @@ -1857,6 +1857,63 @@ pub fn establish_issue_monitor_authority_fence( }) } +/// Acquire short-lived Issue Monitor effect authority for the GUI fallback. +/// +/// The local driver deliberately does not publish a durable daemon fence: it +/// owns only one bounded scan/effect pass. Holding the same lifetime lock as a +/// daemon closes the gap between the fence pre-check and remote submission, +/// while the prefs lock preserves the daemon's `prefs.lock -> authority.lock` +/// ordering. A free lifetime lock proves a v2 fence is stale, so the fallback +/// revokes its epoch and removes it before proceeding. Legacy fences remain +/// fail-closed because they did not carry lifetime-lock liveness. +pub fn try_acquire_issue_monitor_local_fallback_lease( + prefs_path: &Path, +) -> io::Result { + with_issue_monitor_prefs_lock(prefs_path, || { + let authority_lock = fs::OpenOptions::new() + .create(true) + .read(true) + .write(true) + .truncate(false) + .open(issue_monitor_authority_lock_path(prefs_path))?; + if let Err(error) = FileExt::try_lock_exclusive(&authority_lock) { + if gwt_core::operation_deadline::is_lock_contended(&error) { + return Err(io::Error::new( + io::ErrorKind::WouldBlock, + "Issue Monitor daemon authority lease is already held", + )); + } + return Err(error); + } + let lease = IssueMonitorAuthorityLease { + lock: authority_lock, + }; + match load_issue_monitor_authority_fence(prefs_path)? { + IssueMonitorAuthorityFenceState::Missing => Ok(lease), + IssueMonitorAuthorityFenceState::Active(existing) + if existing.version == ISSUE_MONITOR_AUTHORITY_FENCE_VERSION => + { + let mut prefs = load_issue_monitor_prefs_unlocked(prefs_path)?; + prefs.advance_effect_authority_epoch().ok_or_else(|| { + io::Error::other( + "Issue Monitor authority epoch exhausted during local fence recovery", + ) + })?; + save_issue_monitor_prefs_unlocked(prefs_path, &prefs)?; + let fence_path = issue_monitor_authority_fence_path(prefs_path); + fs::remove_file(&fence_path)?; + sync_parent_directory(&fence_path)?; + Ok(lease) + } + IssueMonitorAuthorityFenceState::LegacyShutdownRevoke + | IssueMonitorAuthorityFenceState::Active(_) => Err(io::Error::new( + io::ErrorKind::WouldBlock, + "Issue Monitor daemon authority fence excludes the local fallback", + )), + } + }) +} + pub fn clear_issue_monitor_authority_fence( prefs_path: &Path, expected: &IssueMonitorAuthorityFence, @@ -9197,6 +9254,79 @@ mod tests { drop(second_lease); } + #[test] + fn local_fallback_lease_rejects_live_daemon_authority() { + let temp = tempfile::tempdir().expect("tempdir"); + let prefs_path = temp.path().join("issue-monitor.json"); + save_issue_monitor_prefs(&prefs_path, &IssueMonitorPrefs::default()).expect("seed prefs"); + let daemon = IssueMonitorAuthorityFence::current_process(); + let (_, daemon_lease) = + establish_issue_monitor_authority_fence(&prefs_path, &daemon, |_| false) + .expect("daemon authority"); + + let error = try_acquire_issue_monitor_local_fallback_lease(&prefs_path) + .expect_err("live daemon authority must exclude the GUI fallback"); + + assert_eq!(error.kind(), io::ErrorKind::WouldBlock); + drop(daemon_lease); + } + + #[test] + fn local_fallback_lease_recovers_an_unlocked_v2_fence() { + let temp = tempfile::tempdir().expect("tempdir"); + let prefs_path = temp.path().join("issue-monitor.json"); + save_issue_monitor_prefs( + &prefs_path, + &IssueMonitorPrefs { + effect_authority_epoch: 11, + ..IssueMonitorPrefs::default() + }, + ) + .expect("seed prefs"); + persist_issue_monitor_authority_fence( + &prefs_path, + &IssueMonitorAuthorityFence::current_process(), + ) + .expect("seed stale v2 fence without its lifetime lock"); + + let lease = try_acquire_issue_monitor_local_fallback_lease(&prefs_path) + .expect("a free lifetime lock makes the v2 fence recoverable"); + + assert_eq!( + load_issue_monitor_prefs(&prefs_path) + .expect("load recovered prefs") + .effect_authority_epoch, + 12, + "stale daemon effects must be revoked before GUI fallback execution" + ); + assert_eq!( + load_issue_monitor_authority_fence(&prefs_path).expect("load recovered fence"), + IssueMonitorAuthorityFenceState::Missing, + "the bounded GUI lease must not leave a durable daemon fence" + ); + drop(lease); + } + + #[test] + fn local_fallback_lease_blocks_daemon_start_until_drop() { + let temp = tempfile::tempdir().expect("tempdir"); + let prefs_path = temp.path().join("issue-monitor.json"); + save_issue_monitor_prefs(&prefs_path, &IssueMonitorPrefs::default()).expect("seed prefs"); + let local_lease = try_acquire_issue_monitor_local_fallback_lease(&prefs_path) + .expect("GUI fallback authority"); + let daemon = IssueMonitorAuthorityFence::current_process(); + + let blocked = establish_issue_monitor_authority_fence(&prefs_path, &daemon, |_| false) + .expect_err("daemon startup must not overlap a local remote effect"); + assert_eq!(blocked.kind(), io::ErrorKind::WouldBlock); + + drop(local_lease); + let (_, daemon_lease) = + establish_issue_monitor_authority_fence(&prefs_path, &daemon, |_| false) + .expect("daemon starts after the local effect boundary closes"); + drop(daemon_lease); + } + #[test] fn legacy_v1_authority_fence_fails_closed_for_a_live_pid_without_a_lifetime_lock() { let temp = tempfile::tempdir().expect("tempdir"); diff --git a/crates/gwt/src/lib.rs b/crates/gwt/src/lib.rs index 4c9dabc060..de8b2c25f6 100644 --- a/crates/gwt/src/lib.rs +++ b/crates/gwt/src/lib.rs @@ -118,17 +118,18 @@ pub use issue_monitor::{ persist_issue_monitor_authority_fence, persist_legacy_issue_monitor_shutdown_revoke_fence, record_autonomous_question_handoff, save_issue_monitor_prefs, scan_issue_monitor_candidates, scan_issue_monitor_candidates_with_provenance, take_autonomous_resume_prompt_from_prefs, - try_mutate_issue_monitor_prefs, try_mutate_issue_monitor_prefs_without_authority_fence, - AutonomousHandoffResumption, AutonomousIssueRecord, AutonomousPendingQuestion, AutonomousPhase, - AutonomousReviewDispatch, EligibilityDecision, FailureClass, IssueMonitorAgentStatus, - IssueMonitorAuthorityFence, IssueMonitorAuthorityFenceState, IssueMonitorAuthorityLease, - IssueMonitorCandidateSource, IssueMonitorConfig, IssueMonitorControlReceipt, - IssueMonitorEffectAttemptKey, IssueMonitorEffectPayload, IssueMonitorEffectState, - IssueMonitorFailedIssue, IssueMonitorFailoverOutcome, IssueMonitorInboxItem, IssueMonitorIssue, - IssueMonitorIssueState, IssueMonitorLaunchPlan, IssueMonitorLaunchProfile, - IssueMonitorLaunchProfileSource, IssueMonitorLaunchRequest, IssueMonitorLaunchedIssue, - IssueMonitorLaunchingIssue, IssueMonitorPrefs, IssueMonitorReadiness, IssueMonitorScanSummary, - IssueMonitorState, IssueMonitorStatusView, IssueMonitorStopMismatch, IssueMonitorStopOutcome, + try_acquire_issue_monitor_local_fallback_lease, try_mutate_issue_monitor_prefs, + try_mutate_issue_monitor_prefs_without_authority_fence, AutonomousHandoffResumption, + AutonomousIssueRecord, AutonomousPendingQuestion, AutonomousPhase, AutonomousReviewDispatch, + EligibilityDecision, FailureClass, IssueMonitorAgentStatus, IssueMonitorAuthorityFence, + IssueMonitorAuthorityFenceState, IssueMonitorAuthorityLease, IssueMonitorCandidateSource, + IssueMonitorConfig, IssueMonitorControlReceipt, IssueMonitorEffectAttemptKey, + IssueMonitorEffectPayload, IssueMonitorEffectState, IssueMonitorFailedIssue, + IssueMonitorFailoverOutcome, IssueMonitorInboxItem, IssueMonitorIssue, IssueMonitorIssueState, + IssueMonitorLaunchPlan, IssueMonitorLaunchProfile, IssueMonitorLaunchProfileSource, + IssueMonitorLaunchRequest, IssueMonitorLaunchedIssue, IssueMonitorLaunchingIssue, + IssueMonitorPrefs, IssueMonitorReadiness, IssueMonitorScanSummary, IssueMonitorState, + IssueMonitorStatusView, IssueMonitorStopMismatch, IssueMonitorStopOutcome, IssueMonitorStopTarget, MonitorInboxState, PendingIssueMonitorEffect, LEGACY_GIT_LAUNCH_FAILURE_MIGRATION_VERSION, }; diff --git a/crates/gwt/src/main.rs b/crates/gwt/src/main.rs index d85740d7f8..e0a48b1d1f 100644 --- a/crates/gwt/src/main.rs +++ b/crates/gwt/src/main.rs @@ -67,7 +67,8 @@ pub(crate) use app_runtime::{ pub(crate) use app_runtime::{ ActiveAgentSession, AgentFrontendDispatchOutcome, AgentLaunchResult, AppEventProxy, AppRuntime, BlockingTaskSpawner, ContinueWorkReadinessWatch, DispatchTarget, IssueLaunchWizardPrepared, - OutboundEvent, ProcessLaunch, ProjectOpenTarget, ProjectTabRuntime, WindowAddress, + OutboundEvent, ProcessLaunch, ProjectOpenTarget, ProjectTabRuntime, + ScheduledIssueMonitorScanOutcome, WindowAddress, }; pub(crate) use attachment_upload::{AttachmentUploadStore, UploadedAttachment}; pub(crate) use docker_launch::{ @@ -1123,6 +1124,12 @@ enum UserEvent { /// Issue #3505 / SPEC-3431 FR-108(b): the GUI-owned scheduled monitor /// tick — drives local scans and the PM periodic wake. IssueMonitorScheduledTick, + IssueMonitorScheduledScanComplete { + project_root: PathBuf, + prefs_path: PathBuf, + now: String, + outcome: Result, + }, /// SPEC #3200 Option A: spawn an independent review agent for a PR-ready /// autonomous issue (daemon → GUI). IssueMonitorReviewDispatch { @@ -2646,6 +2653,7 @@ mod tests { pending_launch_feedback_contexts: HashMap::new(), issue_monitor_launch_deliveries: HashMap::new(), issue_monitor_materializer_id: "main-test-materializer".to_string(), + issue_monitor_scheduled_scans_in_flight: std::collections::HashSet::new(), pending_workspace_resume_contexts: HashMap::new(), pending_continue_work: HashMap::new(), pending_fresh_execution_launches: HashMap::new(), @@ -7797,7 +7805,7 @@ fn main() -> std::io::Result<()> { .poll_interval_secs .max(60), ); - std::thread::Builder::new() + if let Err(error) = std::thread::Builder::new() .name("issue-monitor-scheduled-tick".to_string()) .spawn(move || loop { std::thread::sleep(interval); @@ -7808,7 +7816,9 @@ fn main() -> std::io::Result<()> { break; } }) - .ok(); + { + tracing::error!(%error, "failed to start Issue Monitor scheduled tick thread"); + } } #[cfg(unix)] let mut board_daemon_subscribers = BoardDaemonSubscriberRegistry::default(); @@ -8238,6 +8248,20 @@ fn main() -> std::io::Result<()> { let events = app.issue_monitor_scheduled_tick_events(); clients.dispatch(events); } + Event::UserEvent(UserEvent::IssueMonitorScheduledScanComplete { + project_root, + prefs_path, + now, + outcome, + }) => { + let events = app.issue_monitor_scheduled_scan_complete_events( + &project_root, + &prefs_path, + &now, + outcome, + ); + clients.dispatch(events); + } Event::UserEvent(UserEvent::IssueMonitorDaemonInbox { project_root, items,