diff --git a/.yoi/tickets/00001KX6BPY7M/item.md b/.yoi/tickets/00001KX6BPY7M/item.md index f19f5bd9..5805fff5 100644 --- a/.yoi/tickets/00001KX6BPY7M/item.md +++ b/.yoi/tickets/00001KX6BPY7M/item.md @@ -1,8 +1,8 @@ --- title: 'Add Backend Worker/Workdir registry and link model' -state: 'inprogress' +state: 'closed' created_at: '2026-07-10T15:53:02Z' -updated_at: '2026-07-10T16:40:17Z' +updated_at: '2026-07-10T17:38:20Z' assignee: null queued_by: 'workspace-panel' queued_at: '2026-07-10T16:10:57Z' diff --git a/.yoi/tickets/00001KX6BPY7M/resolution.md b/.yoi/tickets/00001KX6BPY7M/resolution.md new file mode 100644 index 00000000..523132ae --- /dev/null +++ b/.yoi/tickets/00001KX6BPY7M/resolution.md @@ -0,0 +1,35 @@ +Backend Worker/Workdir registry and link model を実装・レビュー・merge・検証した。 + +実装内容: +- Backend SQLite schema に `worker_registry`, `workdir_registry`, `worker_workdir_links` を追加。 +- Worker registry / Workdir registry / Worker-Workdir link の typed store API と projection を追加。 +- Backend 経由の Workdir 作成で Runtime materialization 前に Backend-managed pending row を作成。 +- Backend / Runtime scoped Worker creation で Worker registry row、Workdir registry row、Worker-Workdir link を同期。 +- Worker list は Runtime observation を同期したうえで Backend `worker_registry` / link authority から projection。 +- Worker stop は canonical registry row を削除せず、Worker lifecycle と linked Workdir status を同期。 +- Workdir list / launch options は Backend-managed rows を基準にし、`runtime_unmanaged` Workdirs は managed Browser list/options に混ぜず runtime/diagnostic projection で区別可能にした。 +- `pinned` retention が通常 Runtime sync で `normal` に clobber されないよう保護。 +- Browser-facing summaries に safe `management_kind` を追加し、raw Runtime materialized path を Backend registry / Browser response に保存・表示しない方針を維持。 + +Review: +- 初回 review は 5 blocker で `request_changes`。 +- `8695089f fix: sync backend worker registry projections` で主要 blocker を修正。 +- 2回目 review は unmanaged Workdir が managed list/options に混ざる blocker で `request_changes`。 +- `239f93b0 fix: keep unmanaged workdirs out of managed lists` で残 blocker を修正。 +- 最終 focused re-review は `approve`。 + +Merge / validation: +- Merge commit: `03652f82 merge: worker workdir registry`。 +- Final validation passed: + - `git diff --check` + - `cargo test -p yoi-workspace-server --lib` + - `cargo test -p worker-runtime --features ws-server,fs-store` + - `cargo check -p yoi` + - `cd web/workspace && deno task check` + - `cd web/workspace && deno task test` + - `yoi ticket doctor` + - `nix build .#yoi --no-link` +- Validation log: `/run/user/1000/yoi/yoi-orchestrator/bash-output/worker-workdir-final-validation-1783704956.txt` + +Non-blocking follow-up noted by reviewer: +- Some runtime detail/create response paths still return direct Runtime summaries with `management_kind: None`; list/options surfaces required by this Ticket are safe, but a future cleanup can project detail responses through Backend records for more uniform diagnostics。 \ No newline at end of file diff --git a/.yoi/tickets/00001KX6BPY7M/thread.md b/.yoi/tickets/00001KX6BPY7M/thread.md index 45699162..602c4cf9 100644 --- a/.yoi/tickets/00001KX6BPY7M/thread.md +++ b/.yoi/tickets/00001KX6BPY7M/thread.md @@ -197,3 +197,264 @@ Implementation progress report: - External review will be requested via a read-only sibling Reviewer Pod. `StopPod` は使わない。 --- + + + +## Review: request changes + +External review result: request_changes + +Blockers: + +1. Browser Worker list が Backend registry based になっていない。 + - Ticket は Worker/Workdir lists を Backend-managed registry + Runtime observation merge に基づかせることを要求している。 + - `workers_response` はまだ `api.runtime.list_workers(limit)` を呼び、live Runtime items を opportunistic に記録したうえで `items: runtime_workers.items` を返している(`crates/workspace-server/src/server.rs:2529-2546`)。 + - `worker_registry` / `worker_workdir_links` を read していないため、Backend-only / stopped / archive rows を list できず、relation authority が Browser Worker list に使われていない。 + +2. Runtime Worker creation path が Worker registry / Workdir registry / links を書いていない。 + - `create_runtime_worker` は `working_directory_request` を解決して `api.runtime.spawn_worker` に forward し、Runtime result を返すだけ(`server.rs:2269-2285`)。 + - Runtime materialization 前の pending Backend Workdir row、`record_worker_summary`、Worker-Workdir link がない。 + - Browser `/api/workers` path のみ existing Workdir を record/link-select している(`server.rs:2141-2170`)。 + - Backend-created/internal Runtime Workers with materialized workdirs が untracked になり、「Backend worker creation writes/updates Worker registry and Worker-Workdir link」 acceptance に違反する。 + +3. Stopped Worker -> removed/missing Workdir relationship が Backend registry から安定して答えられない。 + - `stop_runtime_worker` は Runtime stop を forward して返すだけ(`server.rs:2330-2339`)。 + - Runtime cleanup 後に Worker lifecycle、Workdir materialization status、link state を Backend SQLite に更新しない。 + - `working_directory_summaries` は `list_managed_workdir_registry` だけを読む(`server.rs:3242-3249`)ため、Runtime stop が bound workdir を remove しても Backend record は stale `present` のままになり得る。 + +4. Runtime Workdir observation / unmanaged distinction が partial。 + - runtime-scoped workdir list は raw Runtime summaries を直接返す(`server.rs:1184-1200`)。 + - workspace-scoped list は Backend-managed rows のみ返す。 + - `runtime_unmanaged` rows は detail / worker-create paths で opportunistic に作られるだけで、list 時に sync/typed されない。 + - Browser summaries も management kind を露出していないため、unmanaged/external Runtime Workdirs を Backend-managed Workdirs と区別して visible/diagnosable にする requirement を満たしていない。 + +5. pinned retention が sync で clobber され得る。 + - schema は `pinned` を保存できるが、`record_worker_summary` は常に `retention_state: "normal"` を書く(`server.rs:3290-3318`)。 + - `upsert_worker_registry` は conflict 時に `retention_state = excluded.retention_state` で上書きする(`store.rs:223-242`)。 + - pinned Worker への後続 Runtime observation / list sync が pinned を normal に downgrade できるため、`pinned` が durable retention metadata として保護されていない。 + +Non-blocking follow-ups: +- schema 方向性自体は妥当。`worker_registry`, `workdir_registry`, `worker_workdir_links` は typed checks / FK relations 付きで存在する(`store.rs:488-545`)。 +- Backend workdir creation が Runtime materialization 前に pending Backend row を作る点は良い(`server.rs:1276-1327`)。 +- inspected paths では raw materialized Runtime paths が新 Backend registry records に保存されている箇所は見当たらなかった。 +- tests は schema roundtrip / pinned storage / removed Workdir / link basics を cover しているが、上記 blocker paths(runtime worker spawn with workdir, stop cleanup syncing, Backend-registry Worker list projection, unmanaged list diagnostics, pinned preservation under sync)は不足。 + +Validation performed: +- Ticket item/thread including routing decision, IntentPacket, implementation report を確認。 +- implementation commit `3e5546ea` の `git show --name-only` / `git show --check` を確認。 +- implementation worktree 内だけで branch/status を確認。 +- `store.rs`, `server.rs`, `catalog.rs`, `worker_backend.rs`, `working_directory.rs` を review。 +- acceptance blockers が inspection で明確だったため、full test/build matrix は再実行していない。 + +--- + + + +## Implementation report + +Follow-up implementation progress report after request_changes: + +- Coder Pod `yoi-coder-00001KX6BPY7M-worker-workdir` completed the focused fix pass and committed follow-up commit `8695089fba346b8f10d55016b8328337e1c37445 fix: sync backend worker registry projections` on branch `work/00001KX6BPY7M-worker-workdir-registry`。 +- Implementation worktree status was clean at Orchestrator inspection。 +- Orchestrator inspected branch log / commit stats / `git diff --check 86f652a7..HEAD` in the implementation worktree。 + +Coder-reported fix summary: +- Browser Worker list now projects from Backend `worker_registry` rows after Runtime observation sync, preserving Backend-only stopped/archive rows。 +- Worker-Workdir relationships are projected from `worker_workdir_links` as authority。 +- Runtime-scoped Worker spawn now writes pending Backend-managed Workdir rows, Worker registry rows, and Worker-Workdir links。 +- Worker stop now updates Backend Worker lifecycle and syncs linked Workdir status to present/removed/missing without deleting canonical rows。 +- Runtime Workdir observation/listing now syncs Runtime observations into Backend registry and projects safe `management_kind` (`backend_managed` / `runtime_unmanaged`)。 +- Existing `pinned` Worker retention is protected from ordinary Runtime sync/upsert clobbering。 +- Tests were added/strengthened for archive/link projection, unmanaged Workdir projection, raw path redaction, removed Workdir preservation, and pinned retention preservation。 +- Frontend type exposure for Workdir `management_kind` was added。 + +Coder-reported validation passed: +- `git diff --check` +- `cargo test -p yoi-workspace-server --lib` +- `cargo test -p worker-runtime --features ws-server,fs-store` +- `cargo check -p yoi` +- `cd web/workspace && deno task check && deno task test` +- `yoi ticket doctor` +- `nix build .#yoi --no-link` + +Next action: +- Request focused external re-review against commits `3e5546ea..8695089f` and the whole Ticket acceptance criteria before merge/close decisions。 + +--- + + + +## Review: request changes + +Focused re-review result: request_changes + +Evidence reviewed: +- Implementation worktree is on `work/00001KX6BPY7M-worker-workdir-registry` at `8695089fba346b8f10d55016b8328337e1c37445`, clean。 +- Follow-up fixes most prior blockers: + - Worker list now syncs Runtime observations, then projects from `worker_registry` / `worker_workdir_links` (`workers_response`, `merge_worker_registry_projection`)。 + - Runtime Worker create now writes pending Backend Workdir rows, records Worker rows, and links selected/materialized Workdirs。 + - Stop path now updates Worker lifecycle and attempts linked Workdir status sync。 + - Store upsert preserves existing `pinned` when a Runtime sync supplies `normal`。 + - New registry records still do not store raw materialized Runtime paths。 + +Blocker: +- Unmanaged Runtime Workdirs are now mixed into normal Backend-managed Browser lists / launch options。 + - Ticket requires Runtime-direct/unmanaged Workdirs to be distinguishable and “not mixed into managed lists” (`item.md:61`)。 + - Follow-up syncs unknown Runtime Workdirs as `runtime_unmanaged`, but `working_directory_summaries()` returns all `workdir_registry` rows via `list_workdir_registry()` (`server.rs:3300-3308`)。 + - This feeds Browser launch options (`server.rs:3278`) and the workspace-scoped Workdirs list。 + - The worker-new UI select iterates `options.working_directories` without filtering or labeling `management_kind`, so a `runtime_unmanaged` Workdir can appear/select like a normal managed Workdir。 + - This is not a safe typed UI distinction and does not satisfy the Backend-managed vs unmanaged boundary。 + +Non-blocking notes: +- Runtime-scoped Workdir lists being typed from registry is a good direction。 +- Consider ensuring detail/cleanup/create responses also project through `workdir_summary_from_record`; several paths still return Runtime summaries with `management_kind: None`, even after storing a typed Backend record。 + +--- + + + +## Implementation report + +Second follow-up implementation progress report: + +- Coder Pod `yoi-coder-00001KX6BPY7M-worker-workdir` completed the focused unmanaged-workdir fix and committed `239f93b084fdcc2201e17fcdbc23eb0140531f16 fix: keep unmanaged workdirs out of managed lists` on branch `work/00001KX6BPY7M-worker-workdir-registry`。 +- Implementation worktree status was clean at Orchestrator inspection。 +- Orchestrator inspected branch log / commit stats / `git diff --check 86f652a7..HEAD` in the implementation worktree。 + +Coder-reported fix summary: +- Workspace-scoped Browser Workdir summaries / launch options now include only `backend_managed` Workdir registry rows。 +- Runtime-direct/unmanaged Workdirs remain visible through runtime-scoped projection with explicit `management_kind: runtime_unmanaged`。 +- Tests prove: + - `runtime_unmanaged` rows are excluded from managed Browser Workdir summaries / launch options。 + - unmanaged rows remain distinguishable in runtime projection。 + - store-level managed listing excludes unmanaged rows while all-row listing still includes them。 + +Coder-reported validation passed: +- `git diff --check` +- `cargo test -p yoi-workspace-server --lib` +- `cargo test -p worker-runtime --features ws-server,fs-store` +- `cargo check -p yoi` +- `cd web/workspace && deno task check && deno task test` +- `yoi ticket doctor` +- `nix build .#yoi --no-link` + +Next action: +- Request focused re-review on remaining unmanaged-list blocker and full acceptance status before merge/close decisions。 + +--- + + + +## Review: approve + +Focused re-review result: approve + +Evidence: +- Follow-up `239f93b084fdcc2201e17fcdbc23eb0140531f16` is present on the requested branch/worktree; worktree was clean during inspection。 +- Remaining blocker is fixed: + - Workspace-scoped Browser launch options still use `working_directory_summaries()`, and that now lists only `backend_managed` rows after Runtime observation sync (`server.rs:3279-3308`)。 + - Store filtering is explicit in `list_managed_workdir_registry()` with `management_kind = 'backend_managed'` (`store.rs:376-392`)。 + - Runtime-scoped Workdir projection still lists registry rows for the Runtime and preserves explicit `management_kind`, so `runtime_unmanaged` remains diagnosable (`server.rs:3310-3317`, `workdir_summary_from_record`)。 +- Previous blockers sanity-checked: + - Worker list syncs Runtime observations then projects from Backend `worker_registry` / link authority (`server.rs:2564-2583`, `merge_worker_registry_projection`)。 + - Runtime Worker create prepares pending Workdir rows, syncs Worker observations, and links Workdirs (`server.rs:2295-2352`)。 + - Stop path preserves canonical records and updates Worker/linked Workdir statuses without deleting registry rows (`server.rs:2399-2443`)。 + - Pinned retention is protected from normal Runtime sync by store upsert `CASE` logic (`store.rs:220-243`)。 + - New Backend registry structs/schema store ids/status/metadata only; inspected Browser summaries do not expose raw Runtime materialized paths。 +- Focused validation run: + - `cargo test -p yoi-workspace-server --lib workdir -- --nocapture` passed 3/3 relevant tests, including managed-list exclusion, unmanaged runtime projection, and pinned/registry roundtrip coverage。 + +Non-blocking note: +- Some runtime detail/create response paths still return direct Runtime summaries with `management_kind: None`; list/options surfaces required by this Ticket are now safe, but projecting details through Backend records would make diagnostics more uniformly typed later。 + +--- + + + +## Implementation report + +Merge / final validation report: + +- Reviewer approved the implementation after the final focused fix。 +- Merged branch `work/00001KX6BPY7M-worker-workdir-registry` into orchestration branch with merge commit `03652f82 merge: worker workdir registry`。 +- Ticket record conflicts from child branch were resolved in favor of Orchestrator-authoritative Ticket records; implementation code changes were merged。 +- Final validation was rerun in `/home/hare/Projects/yoi/.worktree/orchestration` after merge。 + +Validation passed: +- `git diff --check` +- `cargo test -p yoi-workspace-server --lib` +- `cargo test -p worker-runtime --features ws-server,fs-store` +- `cargo check -p yoi` +- `cd web/workspace && deno task check` +- `cd web/workspace && deno task test` +- `yoi ticket doctor` +- `nix build .#yoi --no-link` + +Validation log: +- `/run/user/1000/yoi/yoi-orchestrator/bash-output/worker-workdir-final-validation-1783704956.txt` + +Reviewed/merged implementation commits: +- `3e5546ea feat: add backend worker workdir registry` +- `8695089f fix: sync backend worker registry projections` +- `239f93b0 fix: keep unmanaged workdirs out of managed lists` + +--- + + + +## State changed + +Implementation was reviewed, merged into orchestration branch, and final validation passed. Moving to done before closure. + +--- + + + +## State changed + +Ticket を closed にしました。 + + +--- + + + +## 完了 + +Backend Worker/Workdir registry and link model を実装・レビュー・merge・検証した。 + +実装内容: +- Backend SQLite schema に `worker_registry`, `workdir_registry`, `worker_workdir_links` を追加。 +- Worker registry / Workdir registry / Worker-Workdir link の typed store API と projection を追加。 +- Backend 経由の Workdir 作成で Runtime materialization 前に Backend-managed pending row を作成。 +- Backend / Runtime scoped Worker creation で Worker registry row、Workdir registry row、Worker-Workdir link を同期。 +- Worker list は Runtime observation を同期したうえで Backend `worker_registry` / link authority から projection。 +- Worker stop は canonical registry row を削除せず、Worker lifecycle と linked Workdir status を同期。 +- Workdir list / launch options は Backend-managed rows を基準にし、`runtime_unmanaged` Workdirs は managed Browser list/options に混ぜず runtime/diagnostic projection で区別可能にした。 +- `pinned` retention が通常 Runtime sync で `normal` に clobber されないよう保護。 +- Browser-facing summaries に safe `management_kind` を追加し、raw Runtime materialized path を Backend registry / Browser response に保存・表示しない方針を維持。 + +Review: +- 初回 review は 5 blocker で `request_changes`。 +- `8695089f fix: sync backend worker registry projections` で主要 blocker を修正。 +- 2回目 review は unmanaged Workdir が managed list/options に混ざる blocker で `request_changes`。 +- `239f93b0 fix: keep unmanaged workdirs out of managed lists` で残 blocker を修正。 +- 最終 focused re-review は `approve`。 + +Merge / validation: +- Merge commit: `03652f82 merge: worker workdir registry`。 +- Final validation passed: + - `git diff --check` + - `cargo test -p yoi-workspace-server --lib` + - `cargo test -p worker-runtime --features ws-server,fs-store` + - `cargo check -p yoi` + - `cd web/workspace && deno task check` + - `cd web/workspace && deno task test` + - `yoi ticket doctor` + - `nix build .#yoi --no-link` +- Validation log: `/run/user/1000/yoi/yoi-orchestrator/bash-output/worker-workdir-final-validation-1783704956.txt` + +Non-blocking follow-up noted by reviewer: +- Some runtime detail/create response paths still return direct Runtime summaries with `management_kind: None`; list/options surfaces required by this Ticket are safe, but a future cleanup can project detail responses through Backend records for more uniform diagnostics。 + +--- diff --git a/.yoi/tickets/00001KX6CRVBE/artifacts/orchestration-plan.jsonl b/.yoi/tickets/00001KX6CRVBE/artifacts/orchestration-plan.jsonl new file mode 100644 index 00000000..195571fe --- /dev/null +++ b/.yoi/tickets/00001KX6CRVBE/artifacts/orchestration-plan.jsonl @@ -0,0 +1,3 @@ +{"id":"orch-plan-20260710-164538-1","ticket_id":"00001KX6CRVBE","kind":"after","related_ticket":"00001KX6BPY7M","note":"Ticket body explicitly states this is a follow-up to `00001KX6BPY7M`. The prerequisite registry/link model is currently inprogress on branch `work/00001KX6BPY7M-worker-workdir-registry` and under external review. Start this cleanup/delete Ticket only after that branch is approved/merged or an explicit combined-work decision is made.","author":"orchestrator","at":"2026-07-10T16:45:38Z"} +{"id":"orch-plan-20260710-164549-2","ticket_id":"00001KX6CRVBE","kind":"waiting_capacity_note","note":"Queue review after Dashboard authorization. `00001KX6CRVBE` is otherwise well specified, but it depends on the Backend Worker/Workdir registry/link model from `00001KX6BPY7M`, which is currently inprogress/review. Starting now would duplicate/conflict with store schema, server APIs, and Worker/Workdir projections. Leave queued until `00001KX6BPY7M` is approved/merged; then re-route/start from updated orchestration branch.","author":"orchestrator","at":"2026-07-10T16:45:49Z"} +{"id":"orch-plan-20260710-173850-3","ticket_id":"00001KX6CRVBE","kind":"accepted_plan","accepted_plan":{"summary":"Implement manual Worker/Workdir delete and cleanup operations on top of the merged Backend Worker/Workdir registry/link model. Keep cleanup manual-first with stale revision/digest checks, pinned protection, running-linked Workdir blocks, dirty discard confirmation, raw path redaction, and Browser plan preview/execution surfaces.","branch":"work/00001KX6CRVBE-manual-cleanup","worktree":"/home/hare/Projects/yoi/.worktree/00001KX6CRVBE-manual-cleanup","role_plan":"Orchestrator accepts queued Ticket after prerequisite `00001KX6BPY7M` was approved, merged, validated, closed, and its implementation worktree/branch cleaned. Use sibling Coder for implementation in a dedicated worktree, then sibling Reviewer for read-only review. Orchestrator retains merge/validation/close authority."},"author":"orchestrator","at":"2026-07-10T17:38:50Z"} diff --git a/.yoi/tickets/00001KX6CRVBE/item.md b/.yoi/tickets/00001KX6CRVBE/item.md index 63f6fa64..a1ec6d71 100644 --- a/.yoi/tickets/00001KX6CRVBE/item.md +++ b/.yoi/tickets/00001KX6CRVBE/item.md @@ -1,8 +1,8 @@ --- title: 'Add manual delete and cleanup operations for Workers and workdirs' -state: 'queued' +state: 'closed' created_at: '2026-07-10T16:11:33Z' -updated_at: '2026-07-10T16:45:14Z' +updated_at: '2026-07-10T18:37:40Z' assignee: null queued_by: 'workspace-panel' queued_at: '2026-07-10T16:45:14Z' diff --git a/.yoi/tickets/00001KX6CRVBE/resolution.md b/.yoi/tickets/00001KX6CRVBE/resolution.md new file mode 100644 index 00000000..427f38a7 --- /dev/null +++ b/.yoi/tickets/00001KX6CRVBE/resolution.md @@ -0,0 +1,31 @@ +Manual Worker/Workdir delete and cleanup operations を実装・レビュー・merge・検証した。 + +実装内容: +- Backend Worker registry row の pin/unpin API を追加し、Worker list/detail に `pinned` / `retention_state` を返すようにした。 +- Runtime ごとの manual cleanup plan API を追加し、Backend Worker/Workdir/link registry と Runtime observation に基づいて Worker delete / Workdir cleanup/delete candidates を生成するようにした。 +- Plan candidate は action kind、reason/blocking reason、linked Worker/Workdir ids、pinned state、cleanliness/file status、running-link status、safe な reclaim bytes placeholder を含む。 +- Manual cleanup execution API は expected plan revision/digest を要求し、stale plan を拒否する。 +- Pinned Worker/history delete、running-linked Workdir cleanup、dirty/unknown Workdir の confirmation なし discard を拒否する。 +- Removed/missing Workdir registry record delete は安全条件のもとで separate action として扱う。 +- Runtime-observed/materialized Workdirs は trusted clean evidence なしに `clean` とせず、`unknown` として扱い、normal clean cleanup ではなく explicit discard confirmation path に乗せる。 +- Browser UI に Runtime Workdirs cleanup preview/execution と Worker pin/unpin controls を追加した。 +- UI では verified-clean cleanup と dirty/unknown-state discard を区別し、raw Runtime materialized path を表示しない。 + +Review: +- 初回 review は dirty Workdir safety の実データ経路が不十分として `request_changes`。 +- `361569a6 fix: require discard confirmation for unknown workdirs` で observed Workdir を `unknown` 扱いにし、unknown/dirty cleanup に explicit confirmation を要求するよう修正。 +- focused re-review は `approve`。 +- `estimated_reclaim_bytes: None` は safe size source が未実装のため許容と判断された。 + +Merge / validation: +- Merge commit: `4970a58c merge: manual worker workdir cleanup`。 +- Final validation passed: + - `git diff --check` + - `cargo test -p yoi-workspace-server --lib` + - `cargo test -p worker-runtime --features ws-server,fs-store` + - `cargo check -p yoi` + - `cd web/workspace && deno task check` + - `cd web/workspace && deno task test` + - `yoi ticket doctor` + - `nix build .#yoi --no-link` +- Validation log: `/run/user/1000/yoi/yoi-orchestrator/bash-output/manual-cleanup-final-validation-1783708592.txt` \ No newline at end of file diff --git a/.yoi/tickets/00001KX6CRVBE/thread.md b/.yoi/tickets/00001KX6CRVBE/thread.md index 5c3631a2..a45bbfdd 100644 --- a/.yoi/tickets/00001KX6CRVBE/thread.md +++ b/.yoi/tickets/00001KX6CRVBE/thread.md @@ -78,4 +78,282 @@ Marked ready by `yoi ticket state`. Ticket を `workspace-panel` が queued にしました。 +--- + + + +## Decision + +Routing decision: blocked_by_dependency_or_missing_authority + +Reason: +- Ticket is a manual delete / cleanup follow-up that depends on the Backend Worker/Workdir registry and link authority from `00001KX6BPY7M`。 +- `00001KX6BPY7M` is currently `inprogress` on branch `work/00001KX6BPY7M-worker-workdir-registry` and under external review。 +- This Ticket touches the same store schema, server APIs, Worker/Workdir projections, retention/managed-unmanaged semantics, and cleanup safety rules。 +- Starting this Ticket now would create a high-conflict parallel branch and could implement delete/cleanup behavior against an unreviewed registry model。 +- Therefore this routing pass leaves the Ticket queued and does not record `queued -> inprogress`, create a worktree, or spawn role Pods。 + +Evidence checked: +- Ticket body / thread / artifacts。 +- `TicketRelationQuery(00001KX6CRVBE)`: 0 typed relations, but Ticket body/order context makes `00001KX6BPY7M` a practical implementation prerequisite。 +- `TicketOrchestrationPlanQuery(00001KX6CRVBE)`: prior record 0 件だったため、今回 `after 00001KX6BPY7M` と waiting-capacity note を記録。 +- `TicketList`: queued はこの Ticket 1件、inprogress は `00001KX6BPY7M` 1件。 +- Orchestrator worktree git status: clean on `orchestration`。 + +Next action: +- `00001KX6BPY7M` の implementation/review/merge outcome を待つ。 +- registry/link model が merge されたら、この Ticket を再 routing し、updated orchestration branch から start する。 + +Escalate if: +- 人間が `00001KX6BPY7M` とこの Ticket を同一 combined worktree で統合実装する方針に切り替えたい場合。 + +--- + + + +## Decision + +Routing decision: accept for implementation + +Reason: +- Previously recorded prerequisite `00001KX6BPY7M` has now been reviewed, merged into orchestration branch, final-validated, closed, and its child implementation worktree/branch cleaned up。 +- `TicketRelationQuery(00001KX6CRVBE)` shows no typed blockers。 +- Queue contains this Ticket and there are no other `inprogress` Tickets。 +- Ticket is concrete and well scoped: manual Worker delete / Workdir cleanup plan and execution on top of Backend Worker/Workdir registry/link authority。 + +Accepted plan: +- Create dedicated implementation worktree `/home/hare/Projects/yoi/.worktree/00001KX6CRVBE-manual-cleanup` on branch `work/00001KX6CRVBE-manual-cleanup` from current orchestration branch。 +- Route implementation to sibling Coder Pod。 +- Route external review to sibling Reviewer Pod after Coder commit/validation。 +- Orchestrator retains merge, final validation, Ticket close, and cleanup authority。 + +Design constraints for implementation: +- Cleanup/delete is manual-first only; no automatic prune policy。 +- Worker lifecycle remains `running -> stopped -> delete`; do not introduce `archived` lifecycle state。 +- `pinned` Worker/history is protected in plan generation and execution。 +- Workdir file cleanup blocks when linked to running Workers。 +- Dirty Workdir discard is explicit and separate from ordinary clean cleanup。 +- Execution must reject stale plan revision/digest。 +- Browser-facing responses/UI must not expose raw Runtime paths。 +- Workdir/Worker canonical records should remain Backend/registry based and safe after execution。 + +--- + + + +## State changed + +Prerequisite registry/link model is now closed and merged. Accepting queued manual cleanup Ticket for implementation before creating worktree or spawning role Pods. + +--- + + + +## Implementation report + +Implementation progress report: + +- Coder Pod `yoi-coder-00001KX6CRVBE-manual-cleanup` completed implementation in `/home/hare/Projects/yoi/.worktree/00001KX6CRVBE-manual-cleanup` on branch `work/00001KX6CRVBE-manual-cleanup`。 +- Implementation commit: `ba44391acfdb8076c06ced9324b74d6cc63910ab feat: add manual worker workdir cleanup`。 +- Orchestrator inspected the worktree status, branch log, commit stats, and `git diff --check 5f762a6c..HEAD`; worktree was clean and diff check passed。 + +Coder-reported implementation summary: +- Added Backend Worker retention APIs for pin/unpin registry rows; Worker list/detail carries `pinned` / `retention_state`; pin/unpin mutates Backend registry only and does not affect Runtime processes。 +- Added manual Runtime cleanup plan API with separate Worker delete and Workdir cleanup/delete candidates, action kind, reason/blocking reason, linked ids, pinned state, cleanliness/file status, running-link state, and safe reclaim-byte placeholder。 +- Added manual cleanup execution API requiring expected plan revision/digest; rejects stale plans, pinned Worker/history delete, running-linked Workdir cleanup, and dirty cleanup without explicit discard confirmation; allows safe removed/missing Workdir registry deletion。 +- Added Browser UI for manual cleanup preview/execution on Runtime Workdirs page and Backend-only pin/unpin controls on Worker page。 +- Added focused tests for plan generation, stale rejection, pinned protection, running-linked Workdir blocking, dirty confirmation, removed/missing record delete, and raw path redaction。 + +Coder-reported validation passed: +- `git diff --check` +- `cargo test -p yoi-workspace-server --lib` +- `cargo test -p worker-runtime --features ws-server,fs-store` +- `cargo check -p yoi` +- `cd web/workspace && deno task check && deno task test` +- `yoi ticket doctor` +- `nix build .#yoi --no-link` + +Next action: +- Route to external Reviewer Pod before merge/close decisions。 + +--- + + + +## Review: request changes + +External review result: request_changes + +Blocker: +- Dirty Workdir safety is not implemented against real Runtime/filesystem state。 +- Cleanup plan decides `WorkdirDirtyDiscard` solely from `WorkdirRegistryRecord.cleanliness` (`crates/workspace-server/src/server.rs:1689-1695`)。 +- Runtime observation sync always writes `cleanliness: "clean"` for observed Workdirs (`server.rs:4179-4199`) and there is no code outside tests that records `"dirty"` (`rg "cleanliness.*dirty|dirty.*cleanliness"` only finds plan/test use)。 +- Execution then treats such targets as `WorkdirCleanCleanup` and calls Runtime cleanup without dirty confirmation (`server.rs:1851-1852`)。 +- Runtime cleanup uses forced removal (`crates/worker-runtime/src/working_directory.rs:505-508`)。 +- This violates acceptance items requiring safe dirty/clean status representation and explicit confirmation before dirty discard。 + +Test gap: +- The dirty confirmation test seeds a synthetic registry row with `cleanliness = "dirty"` (`server.rs:5436-5449`, `5551-5580`), so it does not catch the real path where Runtime/workdir sync overwrites or fails to detect dirty files。 + +Positive evidence: +- pin/unpin API/routes and retention projection exist (`server.rs:1544-1584`, `2032-2072`, `4008-4022`)。 +- plan/execution endpoints include digest/revision stale checks and pinned/running-link blocking (`server.rs:1766-1840`)。 +- UI preview/execution and pin controls are present without raw path fields。 + +Non-blocking note: +- `estimated_reclaim_bytes` is present but always `None`; confirm whether “when safe” requires actual byte estimation before approval。 + +--- + + + +## Implementation report + +Follow-up implementation progress report after request_changes: + +- Coder Pod `yoi-coder-00001KX6CRVBE-manual-cleanup` completed the dirty/unknown Workdir safety fix and committed `361569a69efcba929b59b9c48113add5499d5068 fix: require discard confirmation for unknown workdirs`。 +- Orchestrator inspected the implementation worktree status, branch log, commit stats, and `git diff --check 5f762a6c..HEAD`; worktree was clean and diff check passed。 + +Coder-reported fix summary: +- Runtime-observed/materialized Workdirs are no longer marked `clean` without trusted clean evidence。 +- Observed Workdirs now sync as `cleanliness: "unknown"`。 +- Cleanup plan downgrades stale `clean` registry rows to `unknown` when real Runtime observation exists。 +- Present Workdirs with `cleanliness: "unknown"` or `dirty` are classified as discard-confirmation targets, not normal clean cleanup。 +- Direct Workdir cleanup rejects unknown/dirty Workdirs and requires the cleanup execution API with explicit discard confirmation。 +- Browser UI labels distinguish verified-clean cleanup from dirty/unknown-state discard。 +- `estimated_reclaim_bytes` remains intentionally `None` unless a safe estimate is available。 + +Regression tests added/updated: +- Real observed Workdir without verified clean evidence is classified as `unknown` and `WorkdirDirtyDiscard`。 +- Normal/direct cleanup is rejected for unknown Workdir state。 +- Dirty/unknown discard requires explicit confirmation。 +- Stale `clean` registry state is downgraded when Runtime observation exists。 +- Existing path-safe cleanup test uses explicit cleanup execution confirmation。 + +Coder-reported validation passed: +- `git diff --check` +- `cargo test -p yoi-workspace-server --lib` +- `cargo test -p worker-runtime --features ws-server,fs-store` +- `cargo check -p yoi` +- `cd web/workspace && deno task check && deno task test` +- `yoi ticket doctor` +- `nix build .#yoi --no-link` + +Next action: +- Request focused re-review on the dirty/unknown Workdir blocker and overall acceptance status before merge/close decisions。 + +--- + + + +## Review: approve + +Focused re-review result: approve + +Evidence: +- Dirty/unknown blocker is fixed。 +- Runtime/materialized Workdirs now sync as `cleanliness: "unknown"` (`crates/workspace-server/src/server.rs:4208`; pending create also around `1400`)。 +- Cleanup planning downgrades any live-observed non-dirty row to unknown before action selection (`server.rs:1689-1701`)。 +- Unknown/dirty present Workdirs become `WorkdirDirtyDiscard` with explicit-confirmation reasons (`server.rs:1721-1725`)。 +- Direct Workdir DELETE rejects that path before calling Runtime cleanup (`server.rs:1491-1509`)。 +- Execution still checks expected revision + digest for staleness (`server.rs:1780-1788`), blocks pinned/running-linked Workdirs (`server.rs:1703-1707`, `1834-1839`), and requires `confirm_dirty_discard_target_ids` before forced Runtime cleanup for dirty/unknown targets (`server.rs:1842-1850`)。 +- Normal clean cleanup is only reachable for `clean` candidates (`server.rs:1698-1699`, `1860-1861`)。 + +No observed regressions in reviewed acceptance points: +- Pinned Worker deletion remains blocked (`server.rs:1811-1816`)。 +- Workdir record delete remains separate for removed/missing (`server.rs:1696-1697`, `1869-1879`)。 +- Plan responses expose ids/status/reasons/links/retention but not raw paths (`server.rs:620-633`)。 +- UI labels dirty/unknown discard distinctly from verified-clean cleanup with a separate confirm checkbox (`web/workspace/.../workdirs/+page.svelte`)。 +- `estimated_reclaim_bytes: None` is acceptable here because no safe size source is implemented/exposed; the optional field is present and avoids raw-path traversal/leakage。 + +Validation performed by reviewer: +- `cargo test -p yoi-workspace-server cleanup --lib` passed。 +- Targeted regression tests `observed_workdir_without_verified_clean_evidence_requires_discard_confirmation` and `stale_clean_registry_row_is_downgraded_by_real_observation_before_cleanup` passed。 +- `git diff --check HEAD^..HEAD` passed。 + +--- + + + +## Implementation report + +Merge / final validation report: + +- Reviewer approved the implementation after the dirty/unknown Workdir safety fix。 +- Merged branch `work/00001KX6CRVBE-manual-cleanup` into orchestration branch with merge commit `4970a58c merge: manual worker workdir cleanup`。 +- Ticket record conflicts from child branch were resolved in favor of Orchestrator-authoritative Ticket records; implementation code changes were merged。 +- Final validation was rerun in `/home/hare/Projects/yoi/.worktree/orchestration` after merge。 + +Validation passed: +- `git diff --check` +- `cargo test -p yoi-workspace-server --lib` +- `cargo test -p worker-runtime --features ws-server,fs-store` +- `cargo check -p yoi` +- `cd web/workspace && deno task check` +- `cd web/workspace && deno task test` +- `yoi ticket doctor` +- `nix build .#yoi --no-link` + +Validation log: +- `/run/user/1000/yoi/yoi-orchestrator/bash-output/manual-cleanup-final-validation-1783708592.txt` + +Reviewed/merged implementation commits: +- `ba44391a feat: add manual worker workdir cleanup` +- `361569a6 fix: require discard confirmation for unknown workdirs` + +--- + + + +## State changed + +Implementation was reviewed, merged into orchestration branch, and final validation passed. Moving to done before closure. + +--- + + + +## State changed + +Ticket を closed にしました。 + + +--- + + + +## 完了 + +Manual Worker/Workdir delete and cleanup operations を実装・レビュー・merge・検証した。 + +実装内容: +- Backend Worker registry row の pin/unpin API を追加し、Worker list/detail に `pinned` / `retention_state` を返すようにした。 +- Runtime ごとの manual cleanup plan API を追加し、Backend Worker/Workdir/link registry と Runtime observation に基づいて Worker delete / Workdir cleanup/delete candidates を生成するようにした。 +- Plan candidate は action kind、reason/blocking reason、linked Worker/Workdir ids、pinned state、cleanliness/file status、running-link status、safe な reclaim bytes placeholder を含む。 +- Manual cleanup execution API は expected plan revision/digest を要求し、stale plan を拒否する。 +- Pinned Worker/history delete、running-linked Workdir cleanup、dirty/unknown Workdir の confirmation なし discard を拒否する。 +- Removed/missing Workdir registry record delete は安全条件のもとで separate action として扱う。 +- Runtime-observed/materialized Workdirs は trusted clean evidence なしに `clean` とせず、`unknown` として扱い、normal clean cleanup ではなく explicit discard confirmation path に乗せる。 +- Browser UI に Runtime Workdirs cleanup preview/execution と Worker pin/unpin controls を追加した。 +- UI では verified-clean cleanup と dirty/unknown-state discard を区別し、raw Runtime materialized path を表示しない。 + +Review: +- 初回 review は dirty Workdir safety の実データ経路が不十分として `request_changes`。 +- `361569a6 fix: require discard confirmation for unknown workdirs` で observed Workdir を `unknown` 扱いにし、unknown/dirty cleanup に explicit confirmation を要求するよう修正。 +- focused re-review は `approve`。 +- `estimated_reclaim_bytes: None` は safe size source が未実装のため許容と判断された。 + +Merge / validation: +- Merge commit: `4970a58c merge: manual worker workdir cleanup`。 +- Final validation passed: + - `git diff --check` + - `cargo test -p yoi-workspace-server --lib` + - `cargo test -p worker-runtime --features ws-server,fs-store` + - `cargo check -p yoi` + - `cd web/workspace && deno task check` + - `cd web/workspace && deno task test` + - `yoi ticket doctor` + - `nix build .#yoi --no-link` +- Validation log: `/run/user/1000/yoi/yoi-orchestrator/bash-output/manual-cleanup-final-validation-1783708592.txt` + --- diff --git a/crates/worker-runtime/src/catalog.rs b/crates/worker-runtime/src/catalog.rs index 744f8dfd..a42dbc70 100644 --- a/crates/worker-runtime/src/catalog.rs +++ b/crates/worker-runtime/src/catalog.rs @@ -117,6 +117,10 @@ pub struct WorkingDirectoryRequest { pub materializer: MaterializerKind, #[serde(default)] pub dirty_state_policy: DirtyStatePolicy, + /// Backend-assigned stable Workdir id. Runtimes use this when present so the + /// Backend can create canonical registry rows before materialization. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub backend_workdir_id: Option, } #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] @@ -158,6 +162,10 @@ pub struct WorkingDirectorySummary { #[serde(default, skip_serializing_if = "Option::is_none")] pub cleanup_policy: Option, pub status: WorkingDirectoryStatusKind, + /// Backend projection metadata. Runtimes leave this absent; Workspace Browser + /// APIs fill it with `backend_managed` or `runtime_unmanaged`. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub management_kind: Option, } #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] diff --git a/crates/worker-runtime/src/worker_backend.rs b/crates/worker-runtime/src/worker_backend.rs index 335600c8..c62d860c 100644 --- a/crates/worker-runtime/src/worker_backend.rs +++ b/crates/worker-runtime/src/worker_backend.rs @@ -1036,6 +1036,7 @@ mod tests { }, materializer: MaterializerKind::LocalGitWorktree, dirty_state_policy: DirtyStatePolicy::CleanPointOnly, + backend_workdir_id: None, } } diff --git a/crates/worker-runtime/src/working_directory.rs b/crates/worker-runtime/src/working_directory.rs index fb4593e8..76f7c007 100644 --- a/crates/worker-runtime/src/working_directory.rs +++ b/crates/worker-runtime/src/working_directory.rs @@ -51,6 +51,7 @@ impl WorkingDirectory { cleanup_target: Some(self.cleanup_target.clone()), cleanup_policy: Some(self.cleanup_policy.clone()), status: self.status.clone(), + management_kind: None, } } } @@ -395,7 +396,10 @@ impl WorkingDirectoryMaterializer for LocalGitWorktreeMaterializer { &self, request: &WorkingDirectoryRequest, ) -> Result { - let working_directory_id = next_working_directory_id(&request.repository.id); + let working_directory_id = request + .backend_workdir_id + .clone() + .unwrap_or_else(|| next_working_directory_id(&request.repository.id)); self.materialize_with_working_directory_id( working_directory_id, request, @@ -719,6 +723,7 @@ mod tests { }, materializer: MaterializerKind::LocalGitWorktree, dirty_state_policy: DirtyStatePolicy::CleanPointOnly, + backend_workdir_id: None, } } diff --git a/crates/workspace-server/src/hosts.rs b/crates/workspace-server/src/hosts.rs index 158a419f..66b9fef5 100644 --- a/crates/workspace-server/src/hosts.rs +++ b/crates/workspace-server/src/hosts.rs @@ -240,6 +240,10 @@ pub struct WorkerSummary { pub state: String, pub status: String, pub last_seen_at: Option, + #[serde(default)] + pub pinned: bool, + #[serde(default)] + pub retention_state: String, pub implementation: WorkerImplementationSummary, pub capabilities: WorkerCapabilitySummary, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -1191,6 +1195,8 @@ impl EmbeddedWorkerRuntime { status: embedded_worker_execution_status_label(summary.status, &summary.execution) .to_string(), last_seen_at: None, + pinned: false, + retention_state: "transient".to_string(), implementation: WorkerImplementationSummary { kind: "embedded_worker_runtime".to_string(), display_hint: "backend-internal worker-runtime Worker".to_string(), @@ -1226,6 +1232,8 @@ impl EmbeddedWorkerRuntime { status: embedded_worker_execution_status_label(detail.status, &detail.execution) .to_string(), last_seen_at: None, + pinned: false, + retention_state: "transient".to_string(), implementation: WorkerImplementationSummary { kind: "embedded_worker_runtime".to_string(), display_hint: "backend-internal worker-runtime Worker".to_string(), @@ -1927,6 +1935,8 @@ impl RemoteWorkerRuntime { status: embedded_worker_execution_status_label(summary.status, &summary.execution) .to_string(), last_seen_at: None, + pinned: false, + retention_state: "transient".to_string(), implementation: WorkerImplementationSummary { kind: "remote_worker_runtime".to_string(), display_hint: "Backend-proxied remote worker-runtime Worker".to_string(), @@ -1961,6 +1971,8 @@ impl RemoteWorkerRuntime { status: embedded_worker_execution_status_label(detail.status, &detail.execution) .to_string(), last_seen_at: None, + pinned: false, + retention_state: "transient".to_string(), implementation: WorkerImplementationSummary { kind: "remote_worker_runtime".to_string(), display_hint: "Backend-proxied remote worker-runtime Worker".to_string(), @@ -3094,6 +3106,8 @@ pub fn placeholder_worker(host_id: impl Into) -> WorkerSummary { state: "unsupported".to_string(), status: "Worker runtime control is not wired yet".to_string(), last_seen_at: None, + pinned: false, + retention_state: "transient".to_string(), implementation: WorkerImplementationSummary { kind: "placeholder".to_string(), display_hint: "unsupported".to_string(), @@ -3420,6 +3434,8 @@ mod tests { state: "running".to_string(), status: "available".to_string(), last_seen_at: None, + pinned: false, + retention_state: "transient".to_string(), implementation: WorkerImplementationSummary { kind: "fixture".to_string(), display_hint: "test fixture".to_string(), diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index e6041b2e..60b210f9 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -11,6 +11,8 @@ use axum::{Json, Router}; use chrono::{SecondsFormat, Utc}; use futures::StreamExt; use serde::{Deserialize, Serialize}; +use sha2::{Digest, Sha256}; +use std::collections::{HashMap, HashSet}; use tokio::net::TcpListener; use worker_runtime::resource::{BackendResourceError, BackendResourceFetchRequest}; use worker_runtime::worker_backend::{ProfileRuntimeWorkerFactory, WorkerRuntimeExecutionBackend}; @@ -24,10 +26,11 @@ use crate::config::{RemoteRuntimeConfigFile, WorkspaceBackendConfigFile, resolve use crate::hosts::{ ConfigBundleCheckResult, ConfigBundleSyncResult, DiagnosticSeverity, EmbeddedWorkerRuntime, HostSummary, RemoteRuntimeConfig, RemoteWorkerRuntime, RuntimeDiagnostic, RuntimeRegistry, - RuntimeRegistryUnregisterResult, RuntimeSummary, WorkerInputRequest, WorkerInputResult, - WorkerLifecycleRequest, WorkerLifecycleResult, WorkerOperationState, - WorkerSpawnAcceptanceRequirement, WorkerSpawnIntent, WorkerSpawnRequest, WorkerSpawnResult, - WorkerSpawnWorkingDirectoryRequest, WorkerSummary, + RuntimeRegistryUnregisterResult, RuntimeSummary, WorkerCapabilitySummary, + WorkerImplementationSummary, WorkerInputRequest, WorkerInputResult, WorkerLifecycleRequest, + WorkerLifecycleResult, WorkerOperationState, WorkerSpawnAcceptanceRequirement, + WorkerSpawnIntent, WorkerSpawnRequest, WorkerSpawnResult, WorkerSpawnWorkingDirectoryRequest, + WorkerSummary, WorkerWorkspaceSummary, }; use crate::identity::WorkspaceIdentity; use crate::observation::{ @@ -48,12 +51,16 @@ use crate::repositories::{ RepositoryRegistryReader, RepositorySummary, }; use crate::resource_broker::BackendResourceBroker; -use crate::store::{ControlPlaneStore, WorkspaceRecord}; +use crate::store::{ + ControlPlaneStore, WorkdirRegistryRecord, WorkerRegistryRecord, WorkerWorkdirLinkRecord, + WorkspaceRecord, +}; use crate::{Error, Result}; use worker_runtime::catalog::{ ConfigBundleRef, DirtyStatePolicy, MaterializerKind, ProfileSelector, RepositorySelector as RuntimeRepositorySelector, WorkingDirectoryClaim, - WorkingDirectoryRepository, WorkingDirectoryRequest, WorkingDirectorySummary, + WorkingDirectoryRepository, WorkingDirectoryRequest, WorkingDirectoryStatusKind, + WorkingDirectorySummary, }; use worker_runtime::config_bundle::ConfigBundle; use worker_runtime::http_server::{ @@ -468,6 +475,18 @@ pub fn build_router(api: WorkspaceApi) -> Router { "/api/w/{workspace_id}/runtimes/{runtime_id}/workers/{worker_id}", get(scoped_get_runtime_worker), ) + .route( + "/api/w/{workspace_id}/runtimes/{runtime_id}/workers/{worker_id}/pin", + put(scoped_pin_runtime_worker).delete(scoped_unpin_runtime_worker), + ) + .route( + "/api/w/{workspace_id}/runtimes/{runtime_id}/cleanup-plan", + get(scoped_runtime_cleanup_plan), + ) + .route( + "/api/w/{workspace_id}/runtimes/{runtime_id}/cleanup-executions", + post(scoped_execute_runtime_cleanup), + ) .route( "/api/runtimes/{runtime_id}/workers/{worker_id}/input", post(send_runtime_worker_input), @@ -562,6 +581,102 @@ pub struct RuntimeListResponse { pub diagnostics: Vec, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum CleanupTargetKind { + WorkerDelete, + WorkdirCleanCleanup, + WorkdirDirtyDiscard, + WorkdirRecordDelete, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct CleanupWorkerCandidate { + pub target_id: String, + pub action: CleanupTargetKind, + pub worker_id: String, + pub runtime_worker_id: String, + pub runtime_id: String, + pub reason: String, + pub blocking_reason: Option, + pub pinned: bool, + pub retention_state: String, + pub lifecycle_state: String, + pub linked_workdir_ids: Vec, + pub running_linked: bool, + pub estimated_reclaim_bytes: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct CleanupWorkdirCandidate { + pub target_id: String, + pub action: CleanupTargetKind, + pub workdir_id: String, + pub runtime_id: String, + pub repository_id: String, + pub reason: String, + pub blocking_reason: Option, + pub linked_worker_ids: Vec, + pub linked_running_worker_ids: Vec, + pub running_linked: bool, + pub pinned_linked: bool, + pub file_status: String, + pub cleanliness: String, + pub estimated_reclaim_bytes: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct RuntimeCleanupPlanResponse { + pub workspace_id: String, + pub runtime_id: String, + pub generated_at: String, + pub revision: String, + pub digest: String, + pub workers: Vec, + pub workdirs: Vec, + pub diagnostics: Vec, +} + +#[derive(Debug, Clone, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct ExecuteRuntimeCleanupRequest { + pub expected_plan_revision: String, + pub expected_plan_digest: String, + #[serde(default)] + pub worker_target_ids: Vec, + #[serde(default)] + pub workdir_target_ids: Vec, + #[serde(default)] + pub confirm_dirty_discard_target_ids: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct RuntimeCleanupExecutionResult { + pub target_id: String, + pub action: CleanupTargetKind, + pub status: String, + pub message: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct RuntimeCleanupExecutionResponse { + pub workspace_id: String, + pub runtime_id: String, + pub executed_at: String, + pub results: Vec, + pub plan_after: RuntimeCleanupPlanResponse, + pub diagnostics: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkerRetentionResponse { + pub workspace_id: String, + pub runtime_id: String, + pub worker_id: String, + pub pinned: bool, + pub retention_state: String, +} + #[derive(Debug, Serialize, Deserialize)] pub struct RuntimeConnectionSettingsResponse { pub workspace_id: String, @@ -1174,18 +1289,11 @@ async fn scoped_list_runtime_working_directories( AxumPath(path): AxumPath, ) -> ApiResult> { validate_workspace_scope(&api, &path.workspace_id)?; - let list = api - .runtime - .list_working_directories(&path.runtime_id) - .map_err(|err| err.into_error())?; + let (items, diagnostics) = runtime_working_directory_summaries(&api, &path.runtime_id)?; Ok(Json(BrowserWorkingDirectoryListResponse { workspace_id: api.config.workspace_id.clone(), - items: list - .items - .into_iter() - .map(|status| status.summary) - .collect(), - diagnostics: list.diagnostics, + items, + diagnostics, })) } @@ -1266,12 +1374,36 @@ fn create_working_directory_for_runtime( request: BrowserWorkingDirectoryCreateRequest, ) -> ApiResult> { let runtime_id = request.runtime_id.clone(); - let working_directory_request = working_directory_request_for_browser(&api, request)?; + let mut working_directory_request = working_directory_request_for_browser(&api, request)?; + let workdir_id = next_backend_workdir_id(&working_directory_request.repository.id); + working_directory_request.backend_workdir_id = Some(workdir_id.clone()); + let pending = WorkdirRegistryRecord { + workspace_id: api.config.workspace_id.clone(), + workdir_id: workdir_id.clone(), + runtime_id: runtime_id.clone(), + repository_id: working_directory_request.repository.id.clone(), + selector: working_directory_request + .repository + .selector + .as_ref() + .map(|selector| selector.as_ref().to_string()), + resolved_commit: None, + materialization_status: "pending".to_string(), + cleanliness: "unknown".to_string(), + management_kind: "backend_managed".to_string(), + created_at: now_registry_timestamp(), + updated_at: now_registry_timestamp(), + }; + api.store.upsert_workdir_registry(&pending)?; let result = api .runtime .create_working_directory(&runtime_id, working_directory_request) .map_err(|err| err.into_error())?; let Some(working_directory) = result.working_directory else { + let mut failed = pending; + failed.materialization_status = "failed".to_string(); + failed.updated_at = now_registry_timestamp(); + api.store.upsert_workdir_registry(&failed)?; return Err(ApiError::with_diagnostics( Error::RuntimeOperationFailed { runtime_id, @@ -1281,6 +1413,13 @@ fn create_working_directory_for_runtime( result.diagnostics, )); }; + let record = workdir_record_from_summary( + &api, + &runtime_id, + &working_directory.summary, + "backend_managed", + ); + api.store.upsert_workdir_registry(&record)?; Ok(Json(BrowserWorkingDirectoryDetailResponse { workspace_id: api.config.workspace_id.clone(), item: working_directory.summary, @@ -1297,21 +1436,43 @@ fn working_directory_detail_for_runtime( .runtime .working_directory(runtime_id, working_directory_id) .map_err(|err| err.into_error())?; - let Some(working_directory) = result.working_directory else { - return Err(ApiError::with_diagnostics( - Error::RuntimeOperationFailed { - runtime_id: runtime_id.to_string(), - code: "workspace_working_directory_lookup_failed".to_string(), - message: "Runtime did not return working directory".to_string(), - }, - result.diagnostics, - )); - }; - Ok(Json(BrowserWorkingDirectoryDetailResponse { - workspace_id: api.config.workspace_id.clone(), - item: working_directory.summary, - diagnostics: result.diagnostics, - })) + if let Some(working_directory) = result.working_directory { + let management_kind = api + .store + .get_workdir_registry(&api.config.workspace_id, working_directory_id)? + .map(|record| record.management_kind) + .unwrap_or_else(|| "runtime_unmanaged".to_string()); + let record = workdir_record_from_summary( + &api, + runtime_id, + &working_directory.summary, + management_kind.as_str(), + ); + api.store.upsert_workdir_registry(&record)?; + return Ok(Json(BrowserWorkingDirectoryDetailResponse { + workspace_id: api.config.workspace_id.clone(), + item: working_directory.summary, + diagnostics: result.diagnostics, + })); + } + if let Some(record) = api + .store + .get_workdir_registry(&api.config.workspace_id, working_directory_id)? + { + return Ok(Json(BrowserWorkingDirectoryDetailResponse { + workspace_id: api.config.workspace_id.clone(), + item: workdir_summary_from_record(&record), + diagnostics: result.diagnostics, + })); + } + Err(ApiError::with_diagnostics( + Error::RuntimeOperationFailed { + runtime_id: runtime_id.to_string(), + code: "workspace_working_directory_lookup_failed".to_string(), + message: "Runtime did not return working directory".to_string(), + }, + result.diagnostics, + )) } fn cleanup_working_directory_for_runtime( @@ -1319,6 +1480,26 @@ fn cleanup_working_directory_for_runtime( runtime_id: &str, working_directory_id: &str, ) -> ApiResult> { + if let Some(candidate) = build_runtime_cleanup_plan(&api, runtime_id)? + .workdirs + .into_iter() + .find(|candidate| candidate.workdir_id == working_directory_id) + { + if let Some(reason) = candidate.blocking_reason { + return Err(cleanup_api_error( + runtime_id, + "workspace_cleanup_workdir_blocked", + &reason, + )); + } + if candidate.action == CleanupTargetKind::WorkdirDirtyDiscard { + return Err(cleanup_api_error( + runtime_id, + "workspace_cleanup_dirty_confirmation_required", + "dirty Workdir discard requires the cleanup execution API with explicit confirmation", + )); + } + } let result = api .runtime .cleanup_working_directory(runtime_id, working_directory_id) @@ -1333,6 +1514,18 @@ fn cleanup_working_directory_for_runtime( result.diagnostics, )); }; + let management_kind = api + .store + .get_workdir_registry(&api.config.workspace_id, working_directory_id)? + .map(|record| record.management_kind) + .unwrap_or_else(|| "runtime_unmanaged".to_string()); + let record = workdir_record_from_summary( + &api, + runtime_id, + &working_directory.summary, + management_kind.as_str(), + ); + api.store.upsert_workdir_registry(&record)?; Ok(Json(BrowserWorkingDirectoryDetailResponse { workspace_id: api.config.workspace_id.clone(), item: working_directory.summary, @@ -1340,6 +1533,403 @@ fn cleanup_working_directory_for_runtime( })) } +async fn set_worker_retention( + api: WorkspaceApi, + runtime_id: String, + runtime_worker_id: String, + pinned: bool, +) -> ApiResult> { + let worker_id = backend_worker_id(runtime_id.as_str(), runtime_worker_id.as_str()); + if api + .store + .get_worker_registry(&api.config.workspace_id, worker_id.as_str())? + .is_none() + { + if let Ok(worker) = api + .runtime + .worker(runtime_id.as_str(), runtime_worker_id.as_str()) + { + let _ = sync_worker_observation(&api, &worker); + } + } + let retention_state = if pinned { "pinned" } else { "normal" }; + let changed = api.store.update_worker_retention( + &api.config.workspace_id, + worker_id.as_str(), + retention_state, + now_registry_timestamp().as_str(), + )?; + if !changed { + return Err(cleanup_api_error( + runtime_id.as_str(), + "workspace_worker_retention_unknown_worker", + "Worker is not known to the Backend registry", + )); + } + Ok(Json(WorkerRetentionResponse { + workspace_id: api.config.workspace_id, + runtime_id, + worker_id: runtime_worker_id, + pinned, + retention_state: retention_state.to_string(), + })) +} + +fn build_runtime_cleanup_plan( + api: &WorkspaceApi, + runtime_id: &str, +) -> ApiResult { + let _ = workers_response(api.clone()); + let (workdir_summaries, mut diagnostics) = + match runtime_working_directory_summaries(api, runtime_id) { + Ok(result) => result, + Err(error) => { + let mut diagnostics = error.diagnostics; + if diagnostics.is_empty() { + diagnostics.push(RuntimeDiagnostic { + code: "workspace_cleanup_runtime_observation_unavailable".to_string(), + severity: DiagnosticSeverity::Warning, + message: sanitize_backend_error(&error.error.to_string()), + }); + } + (Vec::new(), diagnostics) + } + }; + let workdir_records = api + .store + .list_workdir_registry(&api.config.workspace_id, 500)?; + let worker_records = api + .store + .list_worker_registry(&api.config.workspace_id, 500)?; + let worker_by_id: HashMap<_, _> = worker_records + .iter() + .map(|record| (record.worker_id.clone(), record.clone())) + .collect(); + let observed_workdirs: HashMap<_, _> = workdir_summaries + .into_iter() + .map(|summary| (summary.working_directory_id.clone(), summary)) + .collect(); + + let mut worker_candidates = Vec::new(); + for record in worker_records + .iter() + .filter(|record| record.runtime_id == runtime_id) + { + let links = api + .store + .list_worker_workdir_links(&api.config.workspace_id, record.worker_id.as_str())?; + let is_running = record.lifecycle_state == "running"; + let pinned = record.retention_state == "pinned"; + let blocking_reason = if pinned { + Some("worker is pinned".to_string()) + } else if is_running { + Some("worker is running".to_string()) + } else { + None + }; + worker_candidates.push(CleanupWorkerCandidate { + target_id: format!("worker:{}", record.worker_id), + action: CleanupTargetKind::WorkerDelete, + worker_id: record.worker_id.clone(), + runtime_worker_id: record.runtime_worker_id.clone(), + runtime_id: record.runtime_id.clone(), + reason: if blocking_reason.is_some() { + "Worker registry row cannot be deleted until blocking conditions are cleared" + .to_string() + } else { + "Stopped or unobserved Worker registry row can be manually deleted".to_string() + }, + blocking_reason, + pinned, + retention_state: record.retention_state.clone(), + lifecycle_state: record.lifecycle_state.clone(), + linked_workdir_ids: links.iter().map(|link| link.workdir_id.clone()).collect(), + running_linked: is_running, + estimated_reclaim_bytes: None, + }); + } + + let mut workdir_candidates = Vec::new(); + for record in workdir_records + .iter() + .filter(|record| record.runtime_id == runtime_id) + { + let links = api + .store + .list_workdir_worker_links(&api.config.workspace_id, record.workdir_id.as_str())?; + let linked_workers = links + .iter() + .filter_map(|link| worker_by_id.get(link.worker_id.as_str())) + .collect::>(); + let linked_worker_ids = links + .iter() + .map(|link| link.worker_id.clone()) + .collect::>(); + let linked_running_worker_ids = linked_workers + .iter() + .filter(|worker| worker.lifecycle_state == "running") + .map(|worker| worker.worker_id.clone()) + .collect::>(); + let pinned_linked = linked_workers + .iter() + .any(|worker| worker.retention_state == "pinned"); + let running_linked = !linked_running_worker_ids.is_empty(); + let observed_status = observed_workdirs + .get(record.workdir_id.as_str()) + .map(|summary| format!("{:?}", summary.status).to_lowercase()); + let file_status = observed_status.unwrap_or_else(|| record.materialization_status.clone()); + let observed_without_clean_evidence = + observed_workdirs.contains_key(record.workdir_id.as_str()); + let cleanliness = if observed_without_clean_evidence && record.cleanliness != "dirty" { + "unknown".to_string() + } else { + record.cleanliness.clone() + }; + let action = if matches!(file_status.as_str(), "removed" | "missing") { + CleanupTargetKind::WorkdirRecordDelete + } else if cleanliness == "clean" { + CleanupTargetKind::WorkdirCleanCleanup + } else { + CleanupTargetKind::WorkdirDirtyDiscard + }; + let blocking_reason = if running_linked { + Some("workdir is linked to a running Worker".to_string()) + } else if pinned_linked { + Some("workdir is linked to a pinned Worker/history".to_string()) + } else { + None + }; + workdir_candidates.push(CleanupWorkdirCandidate { + target_id: format!("workdir:{}", record.workdir_id), + action, + workdir_id: record.workdir_id.clone(), + runtime_id: record.runtime_id.clone(), + repository_id: record.repository_id.clone(), + reason: if blocking_reason.is_some() { + "Workdir cleanup is blocked until linked Worker state is safe".to_string() + } else if matches!(file_status.as_str(), "removed" | "missing") { + "Removed or missing Workdir record can be deleted from the Backend registry" + .to_string() + } else if cleanliness == "dirty" { + "Dirty Workdir requires explicit discard confirmation before cleanup".to_string() + } else if cleanliness == "unknown" { + "Workdir clean state is unknown; explicit discard confirmation is required" + .to_string() + } else { + "Clean Workdir can be manually cleaned up".to_string() + }, + blocking_reason, + linked_worker_ids, + linked_running_worker_ids, + running_linked, + pinned_linked, + file_status, + cleanliness, + estimated_reclaim_bytes: None, + }); + } + diagnostics.truncate(16); + let generated_at = now_registry_timestamp(); + let digest = cleanup_plan_digest(&worker_candidates, &workdir_candidates)?; + Ok(RuntimeCleanupPlanResponse { + workspace_id: api.config.workspace_id.clone(), + runtime_id: runtime_id.to_string(), + generated_at, + revision: digest.clone(), + digest, + workers: worker_candidates, + workdirs: workdir_candidates, + diagnostics, + }) +} + +fn cleanup_plan_digest( + workers: &[CleanupWorkerCandidate], + workdirs: &[CleanupWorkdirCandidate], +) -> ApiResult { + let bytes = serde_json::to_vec(&(workers, workdirs)).map_err(|error| { + cleanup_api_error( + "backend", + "workspace_cleanup_plan_digest_failed", + &format!("failed to serialize cleanup plan: {error}"), + ) + })?; + let mut hasher = Sha256::new(); + hasher.update(bytes); + let bytes = hasher.finalize(); + let digest = bytes + .iter() + .map(|byte| format!("{byte:02x}")) + .collect::(); + Ok(format!("sha256:{digest}")) +} + +fn execute_runtime_cleanup( + api: &WorkspaceApi, + runtime_id: &str, + request: ExecuteRuntimeCleanupRequest, +) -> ApiResult { + let plan = build_runtime_cleanup_plan(api, runtime_id)?; + if request.expected_plan_revision != plan.revision + || request.expected_plan_digest != plan.digest + { + return Err(cleanup_api_error( + runtime_id, + "workspace_cleanup_plan_stale", + "cleanup plan revision/digest is stale; refresh the preview before executing", + )); + } + let worker_targets: HashSet<_> = request.worker_target_ids.iter().cloned().collect(); + let workdir_targets: HashSet<_> = request.workdir_target_ids.iter().cloned().collect(); + let dirty_confirmations: HashSet<_> = request + .confirm_dirty_discard_target_ids + .iter() + .cloned() + .collect(); + let mut results = Vec::new(); + + for candidate in plan + .workers + .iter() + .filter(|candidate| worker_targets.contains(candidate.target_id.as_str())) + { + if let Some(reason) = &candidate.blocking_reason { + return Err(cleanup_api_error( + runtime_id, + "workspace_cleanup_worker_blocked", + reason, + )); + } + if candidate.pinned { + return Err(cleanup_api_error( + runtime_id, + "workspace_cleanup_worker_pinned", + "pinned Worker/history cannot be deleted", + )); + } + api.store + .delete_worker_registry(&api.config.workspace_id, candidate.worker_id.as_str())?; + results.push(RuntimeCleanupExecutionResult { + target_id: candidate.target_id.clone(), + action: candidate.action.clone(), + status: "deleted".to_string(), + message: "Worker registry row deleted; Runtime process state was not touched" + .to_string(), + }); + } + + for candidate in plan + .workdirs + .iter() + .filter(|candidate| workdir_targets.contains(candidate.target_id.as_str())) + { + if let Some(reason) = &candidate.blocking_reason { + return Err(cleanup_api_error( + runtime_id, + "workspace_cleanup_workdir_blocked", + reason, + )); + } + match candidate.action { + CleanupTargetKind::WorkdirDirtyDiscard => { + if !dirty_confirmations.contains(candidate.target_id.as_str()) { + return Err(cleanup_api_error( + runtime_id, + "workspace_cleanup_dirty_confirmation_required", + "dirty Workdir discard requires explicit confirmation", + )); + } + cleanup_runtime_workdir_for_execution(api, runtime_id, candidate)?; + results.push(RuntimeCleanupExecutionResult { + target_id: candidate.target_id.clone(), + action: candidate.action.clone(), + status: "discarded".to_string(), + message: + "Dirty/unknown Workdir cleanup/discard was executed after explicit confirmation" + .to_string(), + }); + } + CleanupTargetKind::WorkdirCleanCleanup => { + cleanup_runtime_workdir_for_execution(api, runtime_id, candidate)?; + results.push(RuntimeCleanupExecutionResult { + target_id: candidate.target_id.clone(), + action: candidate.action.clone(), + status: "cleaned".to_string(), + message: "Clean Workdir cleanup was executed".to_string(), + }); + } + CleanupTargetKind::WorkdirRecordDelete => { + api.store.delete_workdir_registry( + &api.config.workspace_id, + candidate.workdir_id.as_str(), + )?; + results.push(RuntimeCleanupExecutionResult { + target_id: candidate.target_id.clone(), + action: candidate.action.clone(), + status: "deleted".to_string(), + message: "Removed/missing Workdir registry row deleted".to_string(), + }); + } + CleanupTargetKind::WorkerDelete => { + return Err(cleanup_api_error( + runtime_id, + "workspace_cleanup_invalid_target_kind", + "worker delete action cannot be executed as a Workdir target", + )); + } + } + } + + let plan_after = build_runtime_cleanup_plan(api, runtime_id)?; + let executed_at = now_registry_timestamp(); + Ok(RuntimeCleanupExecutionResponse { + workspace_id: api.config.workspace_id.clone(), + runtime_id: runtime_id.to_string(), + executed_at, + results, + diagnostics: plan_after.diagnostics.clone(), + plan_after, + }) +} + +fn cleanup_runtime_workdir_for_execution( + api: &WorkspaceApi, + runtime_id: &str, + candidate: &CleanupWorkdirCandidate, +) -> ApiResult<()> { + let result = api + .runtime + .cleanup_working_directory(runtime_id, candidate.workdir_id.as_str()) + .map_err(|err| err.into_error())?; + let Some(working_directory) = result.working_directory else { + return Err(ApiError::with_diagnostics( + Error::RuntimeOperationFailed { + runtime_id: runtime_id.to_string(), + code: "workspace_cleanup_workdir_runtime_failed".to_string(), + message: "Runtime did not cleanup selected Workdir".to_string(), + }, + result.diagnostics, + )); + }; + let record = workdir_record_from_summary( + api, + runtime_id, + &working_directory.summary, + "backend_managed", + ); + api.store.upsert_workdir_registry(&record)?; + Ok(()) +} + +fn cleanup_api_error(runtime_id: &str, code: &str, message: &str) -> ApiError { + Error::RuntimeOperationFailed { + runtime_id: runtime_id.to_string(), + code: code.to_string(), + message: message.to_string(), + } + .into() +} + async fn scoped_get_runtime_connection_settings( State(api): State, AxumPath(path): AxumPath, @@ -1448,6 +2038,41 @@ async fn scoped_get_runtime_worker( get_runtime_worker(State(api), AxumPath((path.runtime_id, path.worker_id))).await } +async fn scoped_pin_runtime_worker( + State(api): State, + AxumPath(path): AxumPath, +) -> ApiResult> { + validate_workspace_scope(&api, &path.workspace_id)?; + set_worker_retention(api, path.runtime_id, path.worker_id, true).await +} + +async fn scoped_unpin_runtime_worker( + State(api): State, + AxumPath(path): AxumPath, +) -> ApiResult> { + validate_workspace_scope(&api, &path.workspace_id)?; + set_worker_retention(api, path.runtime_id, path.worker_id, false).await +} + +async fn scoped_runtime_cleanup_plan( + State(api): State, + AxumPath(path): AxumPath, +) -> ApiResult> { + validate_workspace_scope(&api, &path.workspace_id)?; + let plan = build_runtime_cleanup_plan(&api, path.runtime_id.as_str())?; + Ok(Json(plan)) +} + +async fn scoped_execute_runtime_cleanup( + State(api): State, + AxumPath(path): AxumPath, + Json(request): Json, +) -> ApiResult> { + validate_workspace_scope(&api, &path.workspace_id)?; + let response = execute_runtime_cleanup(&api, path.runtime_id.as_str(), request)?; + Ok(Json(response)) +} + async fn scoped_send_runtime_worker_input( State(api): State, AxumPath(path): AxumPath, @@ -1944,6 +2569,7 @@ fn working_directory_request_from_repository( }, materializer: MaterializerKind::LocalGitWorktree, dirty_state_policy: DirtyStatePolicy::CleanPointOnly, + backend_workdir_id: None, } } @@ -2004,6 +2630,10 @@ async fn create_workspace_worker( content: initial_text, }) }; + let selected_working_directory_id = request + .working_directory + .as_ref() + .map(|selection| selection.working_directory_id.clone()); let resolved_working_directory = request .working_directory @@ -2017,7 +2647,7 @@ async fn create_workspace_worker( .spawn_worker( &request.runtime_id, WorkerSpawnRequest { - requested_worker_name: Some(display_name), + requested_worker_name: Some(display_name.clone()), intent: WorkerSpawnIntent::WorkspaceCoding, acceptance: WorkerSpawnAcceptanceRequirement::RunAccepted { expected_segments: if initial_input.is_some() { 1 } else { 0 }, @@ -2042,6 +2672,37 @@ async fn create_workspace_worker( code: "workspace_worker_create_missing_summary".to_string(), message: "Runtime completed worker creation without returning a Worker summary".to_string(), })?; + let worker_record = sync_worker_observation(&api, &worker)?; + if let Some(workdir_id) = selected_working_directory_id.as_deref() { + if api + .store + .get_workdir_registry(&api.config.workspace_id, workdir_id)? + .is_none() + { + if let Ok(result) = api + .runtime + .working_directory(worker.runtime_id.as_str(), workdir_id) + .map_err(|err| err.into_error()) + { + if let Some(status) = result.working_directory { + let record = workdir_record_from_summary( + &api, + worker.runtime_id.as_str(), + &status.summary, + "runtime_unmanaged", + ); + api.store.upsert_workdir_registry(&record)?; + } + } + } + if api + .store + .get_workdir_registry(&api.config.workspace_id, workdir_id)? + .is_some() + { + link_worker_to_workdir(&api, &worker_record, workdir_id)?; + } + } let runtime_id = worker.runtime_id.clone(); let worker_id = worker.worker_id.clone(); let workspace_id = api.workspace_id().to_string(); @@ -2125,7 +2786,19 @@ async fn get_runtime_worker( .runtime .worker(&runtime_id, &worker_id) .map_err(|err| err.into_error())?; - Ok(Json(worker)) + let record = sync_worker_observation(&api, &worker)?; + let links = api + .store + .list_worker_workdir_links(&api.config.workspace_id, record.worker_id.as_str())?; + let workdirs = api + .store + .list_workdir_registry(&api.config.workspace_id, 500)?; + Ok(Json(merge_worker_registry_projection( + Some(&worker), + &record, + links, + &workdirs, + ))) } #[derive(Debug, Serialize, Deserialize)] @@ -2150,10 +2823,47 @@ async fn create_runtime_worker( configured_working_directory_request(&api.config, working_directory) }) .transpose()?; + let prepared_workdir_id = if let Some(working_directory_request) = + request.resolved_working_directory_request.as_mut() + { + Some(upsert_pending_backend_workdir( + &api, + &runtime_id, + working_directory_request, + )?) + } else { + request + .resolved_working_directory + .as_ref() + .map(|claim| claim.working_directory_id.clone()) + }; let result = api .runtime .spawn_worker(&runtime_id, request) .map_err(|err| err.into_error())?; + if let Some(worker) = result.worker.as_ref() { + let record = sync_worker_observation(&api, worker)?; + if worker.working_directory.is_none() { + if let Some(workdir_id) = prepared_workdir_id.as_deref() { + if api + .store + .get_workdir_registry(&api.config.workspace_id, workdir_id)? + .is_some() + { + link_worker_to_workdir(&api, &record, workdir_id)?; + } + } + } + } else if let Some(workdir_id) = prepared_workdir_id.as_deref() { + if let Some(mut record) = api + .store + .get_workdir_registry(&api.config.workspace_id, workdir_id)? + { + record.materialization_status = "failed".to_string(); + record.updated_at = now_registry_timestamp(); + api.store.upsert_workdir_registry(&record)?; + } + } Ok(Json(result)) } @@ -2208,6 +2918,16 @@ async fn stop_runtime_worker( .runtime .stop_worker(&runtime_id, &worker_id, request) .map_err(|err| err.into_error())?; + let backend_id = backend_worker_id(&runtime_id, &worker_id); + if let Some(mut record) = api + .store + .get_worker_registry(&api.config.workspace_id, backend_id.as_str())? + { + record.lifecycle_state = "stopped".to_string(); + record.updated_at = now_registry_timestamp(); + api.store.upsert_worker_registry(&record)?; + sync_linked_workdir_after_worker_stop(&api, &runtime_id, &record)?; + } Ok(Json(result)) } @@ -2387,11 +3107,37 @@ async fn list_host_workers( fn workers_response(api: WorkspaceApi) -> ApiResult> { let limit = api.config.max_records.min(200); let runtime_workers = api.runtime.list_workers(limit); + let mut observed = std::collections::BTreeMap::new(); + for worker in &runtime_workers.items { + let _ = sync_worker_observation(&api, worker); + observed.insert( + backend_worker_id(worker.runtime_id.as_str(), worker.worker_id.as_str()), + worker.clone(), + ); + } + let worker_records = api + .store + .list_worker_registry(&api.config.workspace_id, limit)?; + let workdir_records = api + .store + .list_workdir_registry(&api.config.workspace_id, 500)?; + let mut items = Vec::new(); + for record in worker_records { + let links = api + .store + .list_worker_workdir_links(&api.config.workspace_id, record.worker_id.as_str())?; + items.push(merge_worker_registry_projection( + observed.get(record.worker_id.as_str()), + &record, + links, + &workdir_records, + )); + } Ok(RuntimeListResponse { workspace_id: api.config.workspace_id, limit, - items: runtime_workers.items, - source: "worker_runtime_registry".to_string(), + items, + source: "backend_worker_registry".to_string(), diagnostics: runtime_workers.diagnostics, }) } @@ -3090,25 +3836,387 @@ fn working_directory_repository_options( } fn working_directory_summaries(api: &WorkspaceApi) -> ApiResult> { - let list = api - .runtime - .list_working_directories(EMBEDDED_WORKER_RUNTIME_ID) - .map_err(|err| err.into_error())?; - if !list.diagnostics.is_empty() { - return Err(ApiError::with_diagnostics( - Error::RuntimeOperationFailed { - runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), - code: "workspace_working_directory_list_failed".to_string(), - message: "Runtime did not list working directories".to_string(), - }, - list.diagnostics, - )); + let _ = sync_all_runtime_workdir_observations(api); + let records = api + .store + .list_managed_workdir_registry(&api.config.workspace_id, 200)?; + Ok(records + .iter() + .map(workdir_summary_from_record) + .collect::>()) +} + +fn runtime_working_directory_summaries( + api: &WorkspaceApi, + runtime_id: &str, +) -> ApiResult<(Vec, Vec)> { + let diagnostics = sync_runtime_workdir_observations(api, runtime_id)?; + let records = api + .store + .list_workdir_registry(&api.config.workspace_id, 200)?; + let items = records + .iter() + .filter(|record| record.runtime_id == runtime_id) + .map(workdir_summary_from_record) + .collect::>(); + Ok((items, diagnostics)) +} + +fn backend_worker_id(runtime_id: &str, runtime_worker_id: &str) -> String { + format!("{runtime_id}/{runtime_worker_id}") +} + +fn now_registry_timestamp() -> String { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|duration| duration.as_millis().to_string()) + .unwrap_or_else(|_| "0".to_string()) +} + +fn registry_safe_id_component(value: &str) -> String { + let sanitized: String = value + .chars() + .map(|ch| { + if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') { + ch + } else { + '-' + } + }) + .collect(); + sanitized.trim_matches('-').chars().take(48).collect() +} + +fn next_backend_workdir_id(repository_id: &str) -> String { + let repository = registry_safe_id_component(repository_id); + format!( + "backend-{}-{}", + now_registry_timestamp(), + if repository.is_empty() { + "workdir" + } else { + &repository + } + ) +} + +fn record_worker_summary( + api: &WorkspaceApi, + worker: &WorkerSummary, + display_name: &str, + profile: Option, +) -> ApiResult { + let timestamp = now_registry_timestamp(); + let worker_id = backend_worker_id(worker.runtime_id.as_str(), worker.worker_id.as_str()); + let existing = api + .store + .get_worker_registry(&api.config.workspace_id, worker_id.as_str())?; + let record = WorkerRegistryRecord { + workspace_id: api.config.workspace_id.clone(), + worker_id: worker_id.clone(), + runtime_id: worker.runtime_id.as_str().to_string(), + runtime_worker_id: worker.worker_id.as_str().to_string(), + display_name: display_name.to_string(), + profile, + lifecycle_state: worker.status.clone(), + retention_state: existing + .as_ref() + .map(|record| record.retention_state.clone()) + .unwrap_or_else(|| "normal".to_string()), + transcript_ref: Some(format!( + "runtime://{}/workers/{}/transcript", + worker.runtime_id.as_str(), + worker.worker_id.as_str() + )), + session_ref: None, + summary_ref: None, + diagnostics_ref: None, + created_at: timestamp.clone(), + updated_at: timestamp, + }; + api.store.upsert_worker_registry(&record)?; + Ok(api + .store + .get_worker_registry(&api.config.workspace_id, worker_id.as_str())? + .unwrap_or(record)) +} + +fn worker_summary_from_registry(record: &WorkerRegistryRecord) -> WorkerSummary { + WorkerSummary { + worker_id: record.runtime_worker_id.clone(), + runtime_id: record.runtime_id.clone(), + host_id: "backend-registry".to_string(), + role: None, + label: record.display_name.clone(), + status: record.lifecycle_state.clone(), + state: record.lifecycle_state.clone(), + last_seen_at: Some(record.updated_at.clone()), + pinned: record.retention_state == "pinned", + retention_state: record.retention_state.clone(), + capabilities: WorkerCapabilitySummary { + can_accept_input: false, + can_stop: false, + can_spawn_followup: false, + }, + workspace: WorkerWorkspaceSummary { + visibility: "backend_registry".to_string(), + identity: record.workspace_id.clone(), + }, + profile: record.profile.clone(), + implementation: WorkerImplementationSummary { + kind: "backend_worker_registry".to_string(), + display_hint: "Archived Worker".to_string(), + }, + working_directory: None, + diagnostics: vec![RuntimeDiagnostic { + code: "backend_worker_registry_only".to_string(), + severity: DiagnosticSeverity::Info, + message: + "Worker is preserved in the Backend registry without a live Runtime observation" + .to_string(), + }], } - Ok(list - .items +} + +fn merge_worker_registry_projection( + live: Option<&WorkerSummary>, + record: &WorkerRegistryRecord, + links: Vec, + workdirs: &[WorkdirRegistryRecord], +) -> WorkerSummary { + let mut summary = live + .cloned() + .unwrap_or_else(|| worker_summary_from_registry(record)); + summary.label = record.display_name.clone(); + summary.status = record.lifecycle_state.clone(); + summary.state = record.lifecycle_state.clone(); + summary.profile = record.profile.clone(); + summary.pinned = record.retention_state == "pinned"; + summary.retention_state = record.retention_state.clone(); + summary.working_directory = links.iter().find_map(|link| { + workdirs + .iter() + .find(|workdir| workdir.workdir_id == link.workdir_id) + .map(|workdir| workdir_summary_from_record(workdir)) + }); + summary +} + +fn sync_worker_observation( + api: &WorkspaceApi, + worker: &WorkerSummary, +) -> ApiResult { + let record = record_worker_summary(api, worker, worker.label.as_str(), worker.profile.clone())?; + if let Some(working_directory) = worker.working_directory.as_ref() { + let management_kind = api + .store + .get_workdir_registry( + &api.config.workspace_id, + &working_directory.working_directory_id, + )? + .map(|existing| existing.management_kind) + .unwrap_or_else(|| "runtime_unmanaged".to_string()); + let workdir_record = workdir_record_from_summary( + api, + worker.runtime_id.as_str(), + working_directory, + management_kind.as_str(), + ); + api.store.upsert_workdir_registry(&workdir_record)?; + link_worker_to_workdir(api, &record, &working_directory.working_directory_id)?; + } + Ok(record) +} + +fn upsert_pending_backend_workdir( + api: &WorkspaceApi, + runtime_id: &str, + request: &mut WorkingDirectoryRequest, +) -> ApiResult { + let workdir_id = request + .backend_workdir_id + .clone() + .unwrap_or_else(|| next_backend_workdir_id(&request.repository.id)); + request.backend_workdir_id = Some(workdir_id.clone()); + let timestamp = now_registry_timestamp(); + api.store.upsert_workdir_registry(&WorkdirRegistryRecord { + workspace_id: api.config.workspace_id.clone(), + workdir_id: workdir_id.clone(), + runtime_id: runtime_id.to_string(), + repository_id: request.repository.id.clone(), + selector: request + .repository + .selector + .as_ref() + .map(|selector| selector.as_ref().to_string()), + resolved_commit: None, + materialization_status: "pending".to_string(), + cleanliness: "unknown".to_string(), + management_kind: "backend_managed".to_string(), + created_at: timestamp.clone(), + updated_at: timestamp, + })?; + Ok(workdir_id) +} + +fn sync_runtime_workdir_observations( + api: &WorkspaceApi, + runtime_id: &str, +) -> ApiResult> { + let response = api + .runtime + .list_working_directories(runtime_id) + .map_err(|err| err.into_error())?; + let mut observed = std::collections::BTreeSet::new(); + for status in &response.items { + observed.insert(status.summary.working_directory_id.clone()); + let management_kind = api + .store + .get_workdir_registry( + &api.config.workspace_id, + &status.summary.working_directory_id, + )? + .map(|existing| existing.management_kind) + .unwrap_or_else(|| "runtime_unmanaged".to_string()); + let record = + workdir_record_from_summary(api, runtime_id, &status.summary, management_kind.as_str()); + api.store.upsert_workdir_registry(&record)?; + } + for mut record in api + .store + .list_workdir_registry(&api.config.workspace_id, 500)? .into_iter() - .map(|status| status.summary) - .collect()) + .filter(|record| record.runtime_id == runtime_id && !observed.contains(&record.workdir_id)) + { + if record.materialization_status == "present" { + record.materialization_status = "missing".to_string(); + record.updated_at = now_registry_timestamp(); + api.store.upsert_workdir_registry(&record)?; + } + } + Ok(response.diagnostics) +} + +fn sync_all_runtime_workdir_observations(api: &WorkspaceApi) -> Vec { + let mut diagnostics = Vec::new(); + let runtimes = api.runtime.list_runtimes(api.config.max_records.min(200)); + for runtime in runtimes.items { + if runtime.capabilities.supports_worktrees { + match sync_runtime_workdir_observations(api, runtime.runtime_id.as_str()) { + Ok(mut runtime_diagnostics) => diagnostics.append(&mut runtime_diagnostics), + Err(err) => diagnostics.extend(err.diagnostics), + } + } + } + diagnostics +} + +fn sync_linked_workdir_after_worker_stop( + api: &WorkspaceApi, + runtime_id: &str, + worker_record: &WorkerRegistryRecord, +) -> ApiResult<()> { + let links = api + .store + .list_worker_workdir_links(&api.config.workspace_id, worker_record.worker_id.as_str())?; + for link in links { + let result = api + .runtime + .working_directory(runtime_id, link.workdir_id.as_str()) + .map_err(|err| err.into_error())?; + if let Some(status) = result.working_directory { + let management_kind = api + .store + .get_workdir_registry(&api.config.workspace_id, link.workdir_id.as_str())? + .map(|record| record.management_kind) + .unwrap_or_else(|| "runtime_unmanaged".to_string()); + let record = workdir_record_from_summary( + api, + runtime_id, + &status.summary, + management_kind.as_str(), + ); + api.store.upsert_workdir_registry(&record)?; + } else if let Some(mut record) = api + .store + .get_workdir_registry(&api.config.workspace_id, link.workdir_id.as_str())? + { + record.materialization_status = "missing".to_string(); + record.updated_at = now_registry_timestamp(); + api.store.upsert_workdir_registry(&record)?; + } + } + Ok(()) +} + +fn workdir_record_from_summary( + api: &WorkspaceApi, + runtime_id: &str, + summary: &WorkingDirectorySummary, + management_kind: &str, +) -> WorkdirRegistryRecord { + let timestamp = now_registry_timestamp(); + WorkdirRegistryRecord { + workspace_id: api.config.workspace_id.clone(), + workdir_id: summary.working_directory_id.clone(), + runtime_id: runtime_id.to_string(), + repository_id: summary.repository_id.clone(), + selector: summary.requested_selector.clone(), + resolved_commit: summary.resolved_commit.clone(), + materialization_status: match summary.status { + WorkingDirectoryStatusKind::Active => "present", + WorkingDirectoryStatusKind::Removed => "removed", + WorkingDirectoryStatusKind::CleanupPending => "pending", + } + .to_string(), + cleanliness: "unknown".to_string(), + management_kind: management_kind.to_string(), + created_at: timestamp.clone(), + updated_at: timestamp, + } +} + +fn workdir_summary_from_record(record: &WorkdirRegistryRecord) -> WorkingDirectorySummary { + let status = match record.materialization_status.as_str() { + "present" => WorkingDirectoryStatusKind::Active, + "pending" => WorkingDirectoryStatusKind::CleanupPending, + _ => WorkingDirectoryStatusKind::Removed, + }; + WorkingDirectorySummary { + working_directory_id: record.workdir_id.clone(), + repository_id: record.repository_id.clone(), + requested_selector: record.selector.clone(), + materializer_kind: MaterializerKind::LocalGitWorktree, + dirty_state_policy: DirtyStatePolicy::CleanPointOnly, + resolved_commit: record.resolved_commit.clone(), + resolved_tree: None, + cleanup_target: Some(worker_runtime::catalog::WorkingDirectoryCleanupTarget { + kind: "local_git_worktree".to_string(), + working_directory_id: record.workdir_id.clone(), + repository_id: record.repository_id.clone(), + }), + cleanup_policy: Some("manual_or_worker_stop".to_string()), + status, + management_kind: Some(record.management_kind.clone()), + } +} + +fn link_worker_to_workdir( + api: &WorkspaceApi, + worker_record: &WorkerRegistryRecord, + workdir_id: &str, +) -> ApiResult<()> { + let timestamp = now_registry_timestamp(); + api.store + .upsert_worker_workdir_link(&WorkerWorkdirLinkRecord { + workspace_id: api.config.workspace_id.clone(), + worker_id: worker_record.worker_id.clone(), + workdir_id: workdir_id.to_string(), + role: "primary_cwd".to_string(), + linked_at: timestamp, + unlinked_at: None, + })?; + Ok(()) } fn validate_working_directory_claim_for_browser( @@ -3171,6 +4279,7 @@ fn working_directory_request_for_browser( }, materializer: MaterializerKind::LocalGitWorktree, dirty_state_policy: DirtyStatePolicy::CleanPointOnly, + backend_workdir_id: None, }) } @@ -3599,7 +4708,11 @@ impl IntoResponse for ApiError { Error::RuntimeOperationFailed { code, .. } if code == "profile_registry_revision_conflict" || code == "profile_source_revision_conflict" - || code == "workspace_metadata_revision_conflict" => + || code == "workspace_metadata_revision_conflict" + || code == "workspace_cleanup_plan_stale" + || code == "workspace_cleanup_worker_blocked" + || code == "workspace_cleanup_workdir_blocked" + || code == "workspace_cleanup_worker_pinned" => { StatusCode::CONFLICT } @@ -3621,6 +4734,7 @@ impl IntoResponse for ApiError { || code.starts_with("invalid_") || code.starts_with("unsupported_worker_profile") || code.starts_with("working_directory_") + || code.starts_with("workspace_cleanup_") || code.ends_with("_already_exists") || code.ends_with("_not_config_managed") || code.ends_with("_unsupported") => @@ -3675,6 +4789,148 @@ mod tests { const TEST_REPOSITORY_ID: &str = "main"; const TEST_CREATED_AT: &str = "2026-06-23T06:43:28Z"; + #[test] + fn backend_worker_projection_preserves_archive_rows_links_and_redacts_paths() { + let worker = WorkerRegistryRecord { + workspace_id: "workspace-1".to_string(), + worker_id: "embedded/worker-1".to_string(), + runtime_id: "embedded".to_string(), + runtime_worker_id: "worker-1".to_string(), + display_name: "Archived Worker".to_string(), + profile: Some("builtin:coder".to_string()), + lifecycle_state: "stopped".to_string(), + retention_state: "pinned".to_string(), + transcript_ref: Some("runtime://embedded/workers/worker-1/transcript".to_string()), + session_ref: None, + summary_ref: None, + diagnostics_ref: None, + created_at: "1".to_string(), + updated_at: "2".to_string(), + }; + let workdir = WorkdirRegistryRecord { + workspace_id: "workspace-1".to_string(), + workdir_id: "backend-1-repo".to_string(), + runtime_id: "embedded".to_string(), + repository_id: "repo".to_string(), + selector: Some("develop".to_string()), + resolved_commit: Some("abcdef".to_string()), + materialization_status: "missing".to_string(), + cleanliness: "clean".to_string(), + management_kind: "backend_managed".to_string(), + created_at: "1".to_string(), + updated_at: "3".to_string(), + }; + let link = WorkerWorkdirLinkRecord { + workspace_id: "workspace-1".to_string(), + worker_id: worker.worker_id.clone(), + workdir_id: workdir.workdir_id.clone(), + role: "primary_cwd".to_string(), + linked_at: "4".to_string(), + unlinked_at: None, + }; + + let projected = merge_worker_registry_projection(None, &worker, vec![link], &[workdir]); + + assert_eq!(projected.status, "stopped"); + assert_eq!( + projected.working_directory.as_ref().unwrap().status, + WorkingDirectoryStatusKind::Removed + ); + assert_eq!( + projected + .working_directory + .as_ref() + .unwrap() + .management_kind + .as_deref(), + Some("backend_managed") + ); + let serialized = serde_json::to_string(&projected).unwrap(); + assert!(!serialized.contains("/tmp/")); + assert!(!serialized.contains("materialized_path")); + } + + #[tokio::test] + async fn workspace_managed_workdir_summaries_exclude_runtime_unmanaged_rows() { + let dir = tempfile::tempdir().unwrap(); + let api = test_api(dir.path()).await; + api.store + .upsert_workdir_registry(&WorkdirRegistryRecord { + workspace_id: TEST_WORKSPACE_ID.to_string(), + workdir_id: "managed".to_string(), + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + repository_id: "repo".to_string(), + selector: None, + resolved_commit: None, + materialization_status: "present".to_string(), + cleanliness: "clean".to_string(), + management_kind: "backend_managed".to_string(), + created_at: "1".to_string(), + updated_at: "1".to_string(), + }) + .unwrap(); + api.store + .upsert_workdir_registry(&WorkdirRegistryRecord { + workspace_id: TEST_WORKSPACE_ID.to_string(), + workdir_id: "runtime-direct".to_string(), + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + repository_id: "repo".to_string(), + selector: None, + resolved_commit: None, + materialization_status: "present".to_string(), + cleanliness: "unknown".to_string(), + management_kind: "runtime_unmanaged".to_string(), + created_at: "1".to_string(), + updated_at: "2".to_string(), + }) + .unwrap(); + + let managed = working_directory_summaries(&api) + .unwrap_or_else(|err| panic!("working_directory_summaries failed: {}", err.error)); + assert_eq!(managed.len(), 1); + assert_eq!(managed[0].working_directory_id, "managed"); + assert_eq!( + managed[0].management_kind.as_deref(), + Some("backend_managed") + ); + + let (runtime_projection, _) = + runtime_working_directory_summaries(&api, EMBEDDED_WORKER_RUNTIME_ID).unwrap_or_else( + |err| panic!("runtime_working_directory_summaries failed: {}", err.error), + ); + assert!(runtime_projection.iter().any(|summary| { + summary.working_directory_id == "runtime-direct" + && summary.management_kind.as_deref() == Some("runtime_unmanaged") + })); + } + #[test] + fn unmanaged_runtime_workdir_projection_is_typed_and_diagnostic_safe() { + let workdir = WorkdirRegistryRecord { + workspace_id: "workspace-1".to_string(), + workdir_id: "runtime-direct".to_string(), + runtime_id: "embedded".to_string(), + repository_id: "repo".to_string(), + selector: None, + resolved_commit: None, + materialization_status: "present".to_string(), + cleanliness: "unknown".to_string(), + management_kind: "runtime_unmanaged".to_string(), + created_at: "1".to_string(), + updated_at: "2".to_string(), + }; + + let projected = workdir_summary_from_record(&workdir); + + assert_eq!(projected.status, WorkingDirectoryStatusKind::Active); + assert_eq!( + projected.management_kind.as_deref(), + Some("runtime_unmanaged") + ); + let serialized = serde_json::to_string(&projected).unwrap(); + assert!(!serialized.contains("/tmp/")); + assert!(!serialized.contains("materialized_path")); + } + #[test] fn worker_profile_candidates_are_backend_published_and_mapped() { let candidates = worker_profile_candidates(); @@ -4121,6 +5377,319 @@ mod tests { .unwrap() } + fn seed_cleanup_worker( + api: &WorkspaceApi, + runtime_worker_id: &str, + lifecycle_state: &str, + retention_state: &str, + ) -> String { + let worker_id = backend_worker_id("runtime-test", runtime_worker_id); + let now = now_registry_timestamp(); + api.store + .upsert_worker_registry(&WorkerRegistryRecord { + workspace_id: api.config.workspace_id.clone(), + worker_id: worker_id.clone(), + runtime_id: "runtime-test".to_string(), + runtime_worker_id: runtime_worker_id.to_string(), + display_name: runtime_worker_id.to_string(), + profile: None, + lifecycle_state: lifecycle_state.to_string(), + retention_state: retention_state.to_string(), + transcript_ref: None, + session_ref: None, + summary_ref: None, + diagnostics_ref: None, + created_at: now.clone(), + updated_at: now, + }) + .unwrap(); + worker_id + } + + fn seed_cleanup_workdir(api: &WorkspaceApi, workdir_id: &str, status: &str, cleanliness: &str) { + let now = now_registry_timestamp(); + api.store + .upsert_workdir_registry(&WorkdirRegistryRecord { + workspace_id: api.config.workspace_id.clone(), + workdir_id: workdir_id.to_string(), + runtime_id: "runtime-test".to_string(), + repository_id: "repo-test".to_string(), + management_kind: "backend_managed".to_string(), + selector: Some("HEAD".to_string()), + resolved_commit: None, + materialization_status: status.to_string(), + cleanliness: cleanliness.to_string(), + created_at: now.clone(), + updated_at: now, + }) + .unwrap(); + } + + fn seed_cleanup_link(api: &WorkspaceApi, worker_id: &str, workdir_id: &str) { + api.store + .upsert_worker_workdir_link(&WorkerWorkdirLinkRecord { + workspace_id: api.config.workspace_id.clone(), + worker_id: worker_id.to_string(), + workdir_id: workdir_id.to_string(), + role: "primary".to_string(), + linked_at: now_registry_timestamp(), + unlinked_at: None, + }) + .unwrap(); + } + + fn create_observed_workdir(api: &WorkspaceApi) -> String { + let Json(response) = create_working_directory_for_runtime( + api.clone(), + BrowserWorkingDirectoryCreateRequest { + runtime_id: EMBEDDED_WORKER_RUNTIME_ID.to_string(), + repository_id: TEST_REPOSITORY_ID.to_string(), + selector: Some("HEAD".to_string()), + policy: BrowserWorkingDirectoryCreatePolicy::default(), + }, + ) + .unwrap_or_else(|err| panic!("create observed workdir: {}", err.error)); + response.item.working_directory_id + } + + #[tokio::test] + async fn observed_workdir_without_verified_clean_evidence_requires_discard_confirmation() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + let workdir_id = create_observed_workdir(&api); + sync_runtime_workdir_observations(&api, EMBEDDED_WORKER_RUNTIME_ID) + .unwrap_or_else(|err| panic!("sync runtime workdir observations: {}", err.error)); + + let stored = api + .store + .get_workdir_registry(&api.config.workspace_id, workdir_id.as_str()) + .unwrap() + .expect("workdir registry row"); + assert_eq!(stored.materialization_status, "present"); + assert_eq!(stored.cleanliness, "unknown"); + + let plan = build_runtime_cleanup_plan(&api, EMBEDDED_WORKER_RUNTIME_ID) + .unwrap_or_else(|err| panic!("cleanup plan: {}", err.error)); + let candidate = plan + .workdirs + .iter() + .find(|candidate| candidate.workdir_id == workdir_id) + .expect("cleanup candidate"); + assert_eq!(candidate.cleanliness, "unknown"); + assert_eq!(candidate.action, CleanupTargetKind::WorkdirDirtyDiscard); + assert!(candidate.reason.contains("unknown")); + + let direct_cleanup = cleanup_working_directory_for_runtime( + api.clone(), + EMBEDDED_WORKER_RUNTIME_ID, + workdir_id.as_str(), + ); + assert!(direct_cleanup.is_err()); + + let target_id = candidate.target_id.clone(); + let missing_confirmation = ExecuteRuntimeCleanupRequest { + expected_plan_revision: plan.revision.clone(), + expected_plan_digest: plan.digest.clone(), + worker_target_ids: Vec::new(), + workdir_target_ids: vec![target_id.clone()], + confirm_dirty_discard_target_ids: Vec::new(), + }; + assert!( + execute_runtime_cleanup(&api, EMBEDDED_WORKER_RUNTIME_ID, missing_confirmation) + .is_err() + ); + + let confirmed = ExecuteRuntimeCleanupRequest { + expected_plan_revision: plan.revision, + expected_plan_digest: plan.digest, + worker_target_ids: Vec::new(), + workdir_target_ids: vec![target_id.clone()], + confirm_dirty_discard_target_ids: vec![target_id], + }; + let response = execute_runtime_cleanup(&api, EMBEDDED_WORKER_RUNTIME_ID, confirmed) + .unwrap_or_else(|err| panic!("cleanup execution: {}", err.error)); + assert_eq!( + response.results[0].action, + CleanupTargetKind::WorkdirDirtyDiscard + ); + assert_eq!(response.results[0].status, "discarded"); + } + + #[tokio::test] + async fn stale_clean_registry_row_is_downgraded_by_real_observation_before_cleanup() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + let workdir_id = create_observed_workdir(&api); + let mut stale_clean = api + .store + .get_workdir_registry(&api.config.workspace_id, workdir_id.as_str()) + .unwrap() + .expect("workdir registry row"); + stale_clean.cleanliness = "clean".to_string(); + api.store.upsert_workdir_registry(&stale_clean).unwrap(); + + let plan = build_runtime_cleanup_plan(&api, EMBEDDED_WORKER_RUNTIME_ID) + .unwrap_or_else(|err| panic!("cleanup plan: {}", err.error)); + let candidate = plan + .workdirs + .iter() + .find(|candidate| candidate.workdir_id == workdir_id) + .expect("cleanup candidate"); + assert_eq!(candidate.cleanliness, "unknown"); + assert_ne!(candidate.action, CleanupTargetKind::WorkdirCleanCleanup); + } + + #[tokio::test] + async fn synthetic_verified_clean_workdir_can_still_use_clean_cleanup_path() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + seed_cleanup_workdir(&api, "verified-clean", "present", "clean"); + + let plan = build_runtime_cleanup_plan(&api, "runtime-test") + .unwrap_or_else(|err| panic!("cleanup plan: {}", err.error)); + let candidate = plan + .workdirs + .iter() + .find(|candidate| candidate.workdir_id == "verified-clean") + .expect("cleanup candidate"); + assert_eq!(candidate.cleanliness, "clean"); + assert_eq!(candidate.action, CleanupTargetKind::WorkdirCleanCleanup); + } + + #[tokio::test] + async fn cleanup_plan_reports_pinned_running_dirty_removed_and_redacts_paths() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + let pinned = seed_cleanup_worker(&api, "worker-pinned", "stopped", "pinned"); + let running = seed_cleanup_worker(&api, "worker-running", "running", "normal"); + seed_cleanup_workdir(&api, "workdir-dirty", "present", "dirty"); + seed_cleanup_workdir(&api, "workdir-removed", "removed", "clean"); + seed_cleanup_link(&api, pinned.as_str(), "workdir-dirty"); + seed_cleanup_link(&api, running.as_str(), "workdir-removed"); + + let plan = build_runtime_cleanup_plan(&api, "runtime-test") + .unwrap_or_else(|err| panic!("cleanup plan: {}", err.error)); + let pinned_worker = plan + .workers + .iter() + .find(|candidate| candidate.worker_id == pinned) + .unwrap(); + assert!(pinned_worker.pinned); + assert_eq!( + pinned_worker.blocking_reason.as_deref(), + Some("worker is pinned") + ); + let running_linked_workdir = plan + .workdirs + .iter() + .find(|candidate| candidate.workdir_id == "workdir-removed") + .unwrap(); + assert_eq!( + running_linked_workdir.action, + CleanupTargetKind::WorkdirRecordDelete + ); + assert!(running_linked_workdir.running_linked); + assert_eq!( + running_linked_workdir.blocking_reason.as_deref(), + Some("workdir is linked to a running Worker") + ); + let dirty_workdir = plan + .workdirs + .iter() + .find(|candidate| candidate.workdir_id == "workdir-dirty") + .unwrap(); + assert_eq!(dirty_workdir.action, CleanupTargetKind::WorkdirDirtyDiscard); + assert!(dirty_workdir.pinned_linked); + let serialized = serde_json::to_string(&plan).unwrap(); + assert!(!serialized.contains("/tmp/secret-runtime-path")); + } + + #[tokio::test] + async fn cleanup_execution_rejects_stale_plan_and_pinned_worker_delete() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + let worker = seed_cleanup_worker(&api, "worker-pinned", "stopped", "pinned"); + let plan = build_runtime_cleanup_plan(&api, "runtime-test") + .unwrap_or_else(|err| panic!("cleanup plan: {}", err.error)); + let target = plan + .workers + .iter() + .find(|candidate| candidate.worker_id == worker) + .unwrap() + .target_id + .clone(); + let stale = ExecuteRuntimeCleanupRequest { + expected_plan_revision: "stale".to_string(), + expected_plan_digest: plan.digest.clone(), + worker_target_ids: vec![target.clone()], + workdir_target_ids: Vec::new(), + confirm_dirty_discard_target_ids: Vec::new(), + }; + assert!(execute_runtime_cleanup(&api, "runtime-test", stale).is_err()); + let pinned = ExecuteRuntimeCleanupRequest { + expected_plan_revision: plan.revision, + expected_plan_digest: plan.digest, + worker_target_ids: vec![target], + workdir_target_ids: Vec::new(), + confirm_dirty_discard_target_ids: Vec::new(), + }; + assert!(execute_runtime_cleanup(&api, "runtime-test", pinned).is_err()); + } + + #[tokio::test] + async fn cleanup_execution_requires_dirty_confirmation_and_deletes_removed_record() { + let workspace = tempfile::tempdir().unwrap(); + init_clean_git_workspace(workspace.path()); + let api = test_api(workspace.path()).await; + seed_cleanup_workdir(&api, "workdir-dirty", "present", "dirty"); + seed_cleanup_workdir(&api, "workdir-removed", "removed", "clean"); + let plan = build_runtime_cleanup_plan(&api, "runtime-test") + .unwrap_or_else(|err| panic!("cleanup plan: {}", err.error)); + let dirty_target = plan + .workdirs + .iter() + .find(|candidate| candidate.workdir_id == "workdir-dirty") + .unwrap() + .target_id + .clone(); + let removed_target = plan + .workdirs + .iter() + .find(|candidate| candidate.workdir_id == "workdir-removed") + .unwrap() + .target_id + .clone(); + let missing_confirmation = ExecuteRuntimeCleanupRequest { + expected_plan_revision: plan.revision.clone(), + expected_plan_digest: plan.digest.clone(), + worker_target_ids: Vec::new(), + workdir_target_ids: vec![dirty_target], + confirm_dirty_discard_target_ids: Vec::new(), + }; + assert!(execute_runtime_cleanup(&api, "runtime-test", missing_confirmation).is_err()); + let delete_removed = ExecuteRuntimeCleanupRequest { + expected_plan_revision: plan.revision, + expected_plan_digest: plan.digest, + worker_target_ids: Vec::new(), + workdir_target_ids: vec![removed_target], + confirm_dirty_discard_target_ids: Vec::new(), + }; + let response = execute_runtime_cleanup(&api, "runtime-test", delete_removed) + .unwrap_or_else(|err| panic!("cleanup execution: {}", err.error)); + assert_eq!(response.results[0].status, "deleted"); + assert!( + api.store + .get_workdir_registry(&api.config.workspace_id, "workdir-removed") + .unwrap() + .is_none() + ); + } + async fn test_app(workspace_root: impl Into) -> Router { build_router(test_api(workspace_root).await) } @@ -4371,8 +5940,50 @@ mod tests { let detail = get_json(app.clone(), &detail_path).await; assert_eq!(detail["item"]["working_directory_id"], working_directory_id); - let removed = request_json(app, "DELETE", &detail_path, None, StatusCode::OK).await; - assert_eq!(removed["item"]["status"], "removed"); + let direct_cleanup = request_json( + app.clone(), + "DELETE", + &detail_path, + None, + StatusCode::BAD_REQUEST, + ) + .await; + assert!( + direct_cleanup["message"] + .as_str() + .unwrap_or_default() + .contains("workspace_cleanup_dirty_confirmation_required") + ); + + let cleanup_plan_path = format!( + "/api/w/{TEST_WORKSPACE_ID}/runtimes/{EMBEDDED_WORKER_RUNTIME_ID}/cleanup-plan" + ); + let plan = get_json(app.clone(), &cleanup_plan_path).await; + let target_id = plan["workdirs"] + .as_array() + .unwrap() + .iter() + .find(|candidate| candidate["workdir_id"] == working_directory_id) + .and_then(|candidate| candidate["target_id"].as_str()) + .unwrap() + .to_string(); + let removed = request_json( + app, + "POST", + &format!( + "/api/w/{TEST_WORKSPACE_ID}/runtimes/{EMBEDDED_WORKER_RUNTIME_ID}/cleanup-executions" + ), + Some(serde_json::json!({ + "expected_plan_revision": plan["revision"], + "expected_plan_digest": plan["digest"], + "worker_target_ids": [], + "workdir_target_ids": [target_id], + "confirm_dirty_discard_target_ids": [target_id] + })), + StatusCode::OK, + ) + .await; + assert_eq!(removed["results"][0]["status"], "discarded"); } #[tokio::test] diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index c9ec8b0f..d0ad0842 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -27,6 +27,11 @@ const MIGRATIONS: &[Migration] = &[ name: "align legacy workspace bootstrap with schema v0", apply: align_legacy_bootstrap_schema, }, + Migration { + version: 3, + name: "backend worker workdir registry schema", + apply: create_worker_workdir_registry_tables, + }, ]; struct Migration { @@ -44,11 +49,113 @@ pub struct WorkspaceRecord { pub updated_at: String, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkerRegistryRecord { + pub workspace_id: String, + /// Backend-owned archival Worker id. In v0 it is derived from runtime_id + runtime_worker_id. + pub worker_id: String, + pub runtime_id: String, + pub runtime_worker_id: String, + pub display_name: String, + pub profile: Option, + pub lifecycle_state: String, + /// Retention state is explicit so `pinned` can be represented before prune exists. + pub retention_state: String, + pub transcript_ref: Option, + pub session_ref: Option, + pub summary_ref: Option, + pub diagnostics_ref: Option, + pub created_at: String, + pub updated_at: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkdirRegistryRecord { + pub workspace_id: String, + pub workdir_id: String, + pub runtime_id: String, + pub repository_id: String, + pub selector: Option, + pub resolved_commit: Option, + pub materialization_status: String, + pub cleanliness: String, + /// `backend_managed` rows are authored by this Backend; `runtime_unmanaged` is for diagnostics only. + pub management_kind: String, + pub created_at: String, + pub updated_at: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct WorkerWorkdirLinkRecord { + pub workspace_id: String, + pub worker_id: String, + pub workdir_id: String, + pub role: String, + pub linked_at: String, + pub unlinked_at: Option, +} + #[async_trait] pub trait ControlPlaneStore: Send + Sync { async fn schema_version(&self) -> Result; async fn upsert_workspace(&self, record: &WorkspaceRecord) -> Result<()>; async fn get_workspace(&self, workspace_id: &str) -> Result>; + + fn upsert_worker_registry(&self, record: &WorkerRegistryRecord) -> Result<()>; + fn get_worker_registry( + &self, + workspace_id: &str, + worker_id: &str, + ) -> Result>; + fn get_worker_registry_by_runtime( + &self, + workspace_id: &str, + runtime_id: &str, + runtime_worker_id: &str, + ) -> Result>; + fn list_worker_registry( + &self, + workspace_id: &str, + limit: usize, + ) -> Result>; + fn update_worker_retention( + &self, + workspace_id: &str, + worker_id: &str, + retention_state: &str, + updated_at: &str, + ) -> Result; + fn delete_worker_registry(&self, workspace_id: &str, worker_id: &str) -> Result; + + fn upsert_workdir_registry(&self, record: &WorkdirRegistryRecord) -> Result<()>; + fn get_workdir_registry( + &self, + workspace_id: &str, + workdir_id: &str, + ) -> Result>; + fn list_workdir_registry( + &self, + workspace_id: &str, + limit: usize, + ) -> Result>; + fn list_managed_workdir_registry( + &self, + workspace_id: &str, + limit: usize, + ) -> Result>; + fn delete_workdir_registry(&self, workspace_id: &str, workdir_id: &str) -> Result; + + fn upsert_worker_workdir_link(&self, record: &WorkerWorkdirLinkRecord) -> Result<()>; + fn list_worker_workdir_links( + &self, + workspace_id: &str, + worker_id: &str, + ) -> Result>; + fn list_workdir_worker_links( + &self, + workspace_id: &str, + workdir_id: &str, + ) -> Result>; } #[derive(Clone)] @@ -131,6 +238,421 @@ impl ControlPlaneStore for SqliteWorkspaceStore { .map_err(Error::from) }) } + + fn upsert_worker_registry(&self, record: &WorkerRegistryRecord) -> Result<()> { + self.with_conn(|conn| { + conn.execute( + r#"INSERT INTO worker_registry ( + workspace_id, worker_id, runtime_id, runtime_worker_id, display_name, profile, + lifecycle_state, retention_state, transcript_ref, session_ref, summary_ref, + diagnostics_ref, created_at, updated_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14) + ON CONFLICT(workspace_id, worker_id) DO UPDATE SET + runtime_id = excluded.runtime_id, + runtime_worker_id = excluded.runtime_worker_id, + display_name = excluded.display_name, + profile = excluded.profile, + lifecycle_state = excluded.lifecycle_state, + retention_state = CASE + WHEN worker_registry.retention_state = 'pinned' AND excluded.retention_state = 'normal' + THEN worker_registry.retention_state + ELSE excluded.retention_state + END, + transcript_ref = excluded.transcript_ref, + session_ref = excluded.session_ref, + summary_ref = excluded.summary_ref, + diagnostics_ref = excluded.diagnostics_ref, + updated_at = excluded.updated_at"#, + params![ + record.workspace_id, + record.worker_id, + record.runtime_id, + record.runtime_worker_id, + record.display_name, + record.profile, + record.lifecycle_state, + record.retention_state, + record.transcript_ref, + record.session_ref, + record.summary_ref, + record.diagnostics_ref, + record.created_at, + record.updated_at, + ], + )?; + Ok(()) + }) + } + + fn get_worker_registry( + &self, + workspace_id: &str, + worker_id: &str, + ) -> Result> { + self.with_conn(|conn| { + conn.query_row( + worker_registry_select_sql("WHERE workspace_id = ?1 AND worker_id = ?2").as_str(), + params![workspace_id, worker_id], + read_worker_registry_record, + ) + .optional() + .map_err(Error::from) + }) + } + + fn get_worker_registry_by_runtime( + &self, + workspace_id: &str, + runtime_id: &str, + runtime_worker_id: &str, + ) -> Result> { + self.with_conn(|conn| { + conn.query_row( + worker_registry_select_sql( + "WHERE workspace_id = ?1 AND runtime_id = ?2 AND runtime_worker_id = ?3", + ) + .as_str(), + params![workspace_id, runtime_id, runtime_worker_id], + read_worker_registry_record, + ) + .optional() + .map_err(Error::from) + }) + } + + fn list_worker_registry( + &self, + workspace_id: &str, + limit: usize, + ) -> Result> { + self.with_conn(|conn| { + let sql = worker_registry_select_sql( + "WHERE workspace_id = ?1 ORDER BY updated_at DESC LIMIT ?2", + ); + let mut stmt = conn.prepare(sql.as_str())?; + let rows = stmt.query_map( + params![workspace_id, limit as i64], + read_worker_registry_record, + )?; + rows.collect::, _>>() + .map_err(Error::from) + }) + } + + fn update_worker_retention( + &self, + workspace_id: &str, + worker_id: &str, + retention_state: &str, + updated_at: &str, + ) -> Result { + self.with_conn(|conn| { + let changed = conn.execute( + r#"UPDATE worker_registry + SET retention_state = ?3, updated_at = ?4 + WHERE workspace_id = ?1 AND worker_id = ?2"#, + params![workspace_id, worker_id, retention_state, updated_at], + )?; + Ok(changed > 0) + }) + } + + fn delete_worker_registry(&self, workspace_id: &str, worker_id: &str) -> Result { + self.with_conn(|conn| { + let changed = conn.execute( + "DELETE FROM worker_registry WHERE workspace_id = ?1 AND worker_id = ?2", + params![workspace_id, worker_id], + )?; + Ok(changed > 0) + }) + } + + fn upsert_workdir_registry(&self, record: &WorkdirRegistryRecord) -> Result<()> { + self.with_conn(|conn| { + conn.execute( + r#"INSERT INTO workdir_registry ( + workspace_id, workdir_id, runtime_id, repository_id, selector, resolved_commit, + materialization_status, cleanliness, management_kind, created_at, updated_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11) + ON CONFLICT(workspace_id, workdir_id) DO UPDATE SET + runtime_id = excluded.runtime_id, + repository_id = excluded.repository_id, + selector = excluded.selector, + resolved_commit = excluded.resolved_commit, + materialization_status = excluded.materialization_status, + cleanliness = excluded.cleanliness, + management_kind = excluded.management_kind, + updated_at = excluded.updated_at"#, + params![ + record.workspace_id, + record.workdir_id, + record.runtime_id, + record.repository_id, + record.selector, + record.resolved_commit, + record.materialization_status, + record.cleanliness, + record.management_kind, + record.created_at, + record.updated_at, + ], + )?; + Ok(()) + }) + } + + fn get_workdir_registry( + &self, + workspace_id: &str, + workdir_id: &str, + ) -> Result> { + self.with_conn(|conn| { + conn.query_row( + workdir_registry_select_sql("WHERE workspace_id = ?1 AND workdir_id = ?2").as_str(), + params![workspace_id, workdir_id], + read_workdir_registry_record, + ) + .optional() + .map_err(Error::from) + }) + } + + fn list_workdir_registry( + &self, + workspace_id: &str, + limit: usize, + ) -> Result> { + self.with_conn(|conn| { + let sql = workdir_registry_select_sql( + "WHERE workspace_id = ?1 ORDER BY updated_at DESC LIMIT ?2", + ); + let mut stmt = conn.prepare(sql.as_str())?; + let rows = stmt.query_map( + params![workspace_id, limit as i64], + read_workdir_registry_record, + )?; + rows.collect::, _>>() + .map_err(Error::from) + }) + } + + fn list_managed_workdir_registry( + &self, + workspace_id: &str, + limit: usize, + ) -> Result> { + self.with_conn(|conn| { + let sql = workdir_registry_select_sql( + "WHERE workspace_id = ?1 AND management_kind = 'backend_managed' ORDER BY updated_at DESC LIMIT ?2", + ); + let mut stmt = conn.prepare(sql.as_str())?; + let rows = stmt.query_map(params![workspace_id, limit as i64], read_workdir_registry_record)?; + rows.collect::, _>>() + .map_err(Error::from) + }) + } + + fn delete_workdir_registry(&self, workspace_id: &str, workdir_id: &str) -> Result { + self.with_conn(|conn| { + let changed = conn.execute( + "DELETE FROM workdir_registry WHERE workspace_id = ?1 AND workdir_id = ?2", + params![workspace_id, workdir_id], + )?; + Ok(changed > 0) + }) + } + + fn upsert_worker_workdir_link(&self, record: &WorkerWorkdirLinkRecord) -> Result<()> { + self.with_conn(|conn| { + conn.execute( + r#"INSERT INTO worker_workdir_links ( + workspace_id, worker_id, workdir_id, role, linked_at, unlinked_at + ) VALUES (?1, ?2, ?3, ?4, ?5, ?6) + ON CONFLICT(workspace_id, worker_id, workdir_id, role) DO UPDATE SET + linked_at = excluded.linked_at, + unlinked_at = excluded.unlinked_at"#, + params![ + record.workspace_id, + record.worker_id, + record.workdir_id, + record.role, + record.linked_at, + record.unlinked_at, + ], + )?; + Ok(()) + }) + } + + fn list_worker_workdir_links( + &self, + workspace_id: &str, + worker_id: &str, + ) -> Result> { + self.with_conn(|conn| { + let mut stmt = conn.prepare( + r#"SELECT workspace_id, worker_id, workdir_id, role, linked_at, unlinked_at + FROM worker_workdir_links + WHERE workspace_id = ?1 AND worker_id = ?2 AND unlinked_at IS NULL + ORDER BY linked_at DESC"#, + )?; + let rows = stmt.query_map( + params![workspace_id, worker_id], + read_worker_workdir_link_record, + )?; + rows.collect::, _>>() + .map_err(Error::from) + }) + } + + fn list_workdir_worker_links( + &self, + workspace_id: &str, + workdir_id: &str, + ) -> Result> { + self.with_conn(|conn| { + let mut stmt = conn.prepare( + r#"SELECT workspace_id, worker_id, workdir_id, role, linked_at, unlinked_at + FROM worker_workdir_links + WHERE workspace_id = ?1 AND workdir_id = ?2 AND unlinked_at IS NULL + ORDER BY linked_at DESC"#, + )?; + let rows = stmt.query_map( + params![workspace_id, workdir_id], + read_worker_workdir_link_record, + )?; + rows.collect::, _>>() + .map_err(Error::from) + }) + } +} + +fn read_worker_workdir_link_record( + row: &rusqlite::Row<'_>, +) -> rusqlite::Result { + Ok(WorkerWorkdirLinkRecord { + workspace_id: row.get(0)?, + worker_id: row.get(1)?, + workdir_id: row.get(2)?, + role: row.get(3)?, + linked_at: row.get(4)?, + unlinked_at: row.get(5)?, + }) +} + +fn worker_registry_select_sql(where_clause: &str) -> String { + format!( + "SELECT workspace_id, worker_id, runtime_id, runtime_worker_id, display_name, profile, \ + lifecycle_state, retention_state, transcript_ref, session_ref, summary_ref, diagnostics_ref, \ + created_at, updated_at FROM worker_registry {where_clause}" + ) +} + +fn read_worker_registry_record(row: &rusqlite::Row<'_>) -> rusqlite::Result { + Ok(WorkerRegistryRecord { + workspace_id: row.get(0)?, + worker_id: row.get(1)?, + runtime_id: row.get(2)?, + runtime_worker_id: row.get(3)?, + display_name: row.get(4)?, + profile: row.get(5)?, + lifecycle_state: row.get(6)?, + retention_state: row.get(7)?, + transcript_ref: row.get(8)?, + session_ref: row.get(9)?, + summary_ref: row.get(10)?, + diagnostics_ref: row.get(11)?, + created_at: row.get(12)?, + updated_at: row.get(13)?, + }) +} + +fn workdir_registry_select_sql(where_clause: &str) -> String { + format!( + "SELECT workspace_id, workdir_id, runtime_id, repository_id, selector, resolved_commit, \ + materialization_status, cleanliness, management_kind, created_at, updated_at \ + FROM workdir_registry {where_clause}" + ) +} + +fn read_workdir_registry_record( + row: &rusqlite::Row<'_>, +) -> rusqlite::Result { + Ok(WorkdirRegistryRecord { + workspace_id: row.get(0)?, + workdir_id: row.get(1)?, + runtime_id: row.get(2)?, + repository_id: row.get(3)?, + selector: row.get(4)?, + resolved_commit: row.get(5)?, + materialization_status: row.get(6)?, + cleanliness: row.get(7)?, + management_kind: row.get(8)?, + created_at: row.get(9)?, + updated_at: row.get(10)?, + }) +} + +fn create_worker_workdir_registry_tables(conn: &Connection) -> Result<()> { + conn.execute_batch( + r#" +CREATE TABLE IF NOT EXISTS worker_registry ( + workspace_id TEXT NOT NULL, + worker_id TEXT NOT NULL, + runtime_id TEXT NOT NULL, + runtime_worker_id TEXT NOT NULL, + display_name TEXT NOT NULL, + profile TEXT, + lifecycle_state TEXT NOT NULL, + retention_state TEXT NOT NULL CHECK (retention_state IN ('normal', 'pinned')), + transcript_ref TEXT, + session_ref TEXT, + summary_ref TEXT, + diagnostics_ref TEXT, + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (workspace_id, worker_id), + UNIQUE (workspace_id, runtime_id, runtime_worker_id), + FOREIGN KEY (workspace_id) REFERENCES workspaces(workspace_id) ON DELETE CASCADE +); + +CREATE TABLE IF NOT EXISTS workdir_registry ( + workspace_id TEXT NOT NULL, + workdir_id TEXT NOT NULL, + runtime_id TEXT NOT NULL, + repository_id TEXT NOT NULL, + selector TEXT, + resolved_commit TEXT, + materialization_status TEXT NOT NULL CHECK (materialization_status IN ('pending', 'present', 'missing', 'removed', 'failed')), + cleanliness TEXT NOT NULL CHECK (cleanliness IN ('clean', 'dirty', 'unknown')), + management_kind TEXT NOT NULL CHECK (management_kind IN ('backend_managed', 'runtime_unmanaged')), + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + PRIMARY KEY (workspace_id, workdir_id), + FOREIGN KEY (workspace_id) REFERENCES workspaces(workspace_id) ON DELETE CASCADE +); + +CREATE TABLE IF NOT EXISTS worker_workdir_links ( + workspace_id TEXT NOT NULL, + worker_id TEXT NOT NULL, + workdir_id TEXT NOT NULL, + role TEXT NOT NULL, + linked_at TEXT NOT NULL, + unlinked_at TEXT, + PRIMARY KEY (workspace_id, worker_id, workdir_id, role), + FOREIGN KEY (workspace_id, worker_id) REFERENCES worker_registry(workspace_id, worker_id) ON DELETE CASCADE, + FOREIGN KEY (workspace_id, workdir_id) REFERENCES workdir_registry(workspace_id, workdir_id) ON DELETE CASCADE +); + +CREATE INDEX IF NOT EXISTS idx_worker_registry_workspace_updated + ON worker_registry(workspace_id, updated_at DESC); +CREATE INDEX IF NOT EXISTS idx_workdir_registry_workspace_updated + ON workdir_registry(workspace_id, updated_at DESC); +CREATE INDEX IF NOT EXISTS idx_worker_workdir_links_worker + ON worker_workdir_links(workspace_id, worker_id, linked_at DESC); +"#, + )?; + Ok(()) } fn configure_sqlite(conn: &Connection) -> Result<()> { @@ -488,7 +1010,7 @@ mod tests { let db = dir.path().join("control-plane.sqlite"); let store = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 2); + assert_eq!(store.schema_version().await.unwrap(), 3); let record = WorkspaceRecord { workspace_id: "local-dev".to_string(), @@ -500,7 +1022,7 @@ mod tests { store.upsert_workspace(&record).await.unwrap(); let reopened = SqliteWorkspaceStore::open(&db).unwrap(); - assert_eq!(reopened.schema_version().await.unwrap(), 2); + assert_eq!(reopened.schema_version().await.unwrap(), 3); assert_eq!( reopened.get_workspace("local-dev").await.unwrap(), Some(record) @@ -527,6 +1049,9 @@ mod tests { "ticket_worker_links", "artifacts", "audit_events", + "worker_registry", + "workdir_registry", + "worker_workdir_links", ] { assert!( tables.contains(expected), @@ -688,7 +1213,7 @@ mod tests { .unwrap(); let store = SqliteWorkspaceStore::from_connection(conn).unwrap(); - assert_eq!(store.schema_version().await.unwrap(), 2); + assert_eq!(store.schema_version().await.unwrap(), 3); store .with_conn(|conn| { @@ -773,6 +1298,115 @@ mod tests { ); } + #[tokio::test] + async fn worker_workdir_registry_round_trips_and_preserves_pinned_retention() { + let temp = tempfile::tempdir().unwrap(); + let db = temp.path().join("workspace.db"); + let store = SqliteWorkspaceStore::open(&db).unwrap(); + let workspace = WorkspaceRecord { + workspace_id: "local-dev".to_string(), + display_name: "Local Dev".to_string(), + state: "active".to_string(), + created_at: "1".to_string(), + updated_at: "1".to_string(), + }; + store.upsert_workspace(&workspace).await.unwrap(); + + let worker = WorkerRegistryRecord { + workspace_id: "local-dev".to_string(), + worker_id: "embedded/browser-1".to_string(), + runtime_id: "embedded".to_string(), + runtime_worker_id: "browser-1".to_string(), + display_name: "Browser 1".to_string(), + profile: Some("builtin:companion".to_string()), + lifecycle_state: "idle".to_string(), + retention_state: "pinned".to_string(), + transcript_ref: Some("runtime://embedded/workers/browser-1/transcript".to_string()), + session_ref: None, + summary_ref: None, + diagnostics_ref: None, + created_at: "2".to_string(), + updated_at: "2".to_string(), + }; + store.upsert_worker_registry(&worker).unwrap(); + let mut runtime_sync_worker = worker.clone(); + runtime_sync_worker.lifecycle_state = "running".to_string(); + runtime_sync_worker.retention_state = "normal".to_string(); + runtime_sync_worker.updated_at = "5".to_string(); + store.upsert_worker_registry(&runtime_sync_worker).unwrap(); + let mut expected_worker = worker.clone(); + expected_worker.lifecycle_state = "running".to_string(); + expected_worker.updated_at = "5".to_string(); + + let workdir = WorkdirRegistryRecord { + workspace_id: "local-dev".to_string(), + workdir_id: "backend-2-repo".to_string(), + runtime_id: "embedded".to_string(), + repository_id: "repo".to_string(), + selector: Some("develop".to_string()), + resolved_commit: Some("abcdef".to_string()), + materialization_status: "removed".to_string(), + cleanliness: "clean".to_string(), + management_kind: "backend_managed".to_string(), + created_at: "2".to_string(), + updated_at: "3".to_string(), + }; + store.upsert_workdir_registry(&workdir).unwrap(); + let unmanaged_workdir = WorkdirRegistryRecord { + workspace_id: "local-dev".to_string(), + workdir_id: "runtime-direct".to_string(), + runtime_id: "embedded".to_string(), + repository_id: "repo".to_string(), + selector: Some("feature".to_string()), + resolved_commit: Some("123456".to_string()), + materialization_status: "present".to_string(), + cleanliness: "unknown".to_string(), + management_kind: "runtime_unmanaged".to_string(), + created_at: "3".to_string(), + updated_at: "4".to_string(), + }; + store.upsert_workdir_registry(&unmanaged_workdir).unwrap(); + + let link = WorkerWorkdirLinkRecord { + workspace_id: "local-dev".to_string(), + worker_id: worker.worker_id.clone(), + workdir_id: workdir.workdir_id.clone(), + role: "primary_cwd".to_string(), + linked_at: "4".to_string(), + unlinked_at: None, + }; + store.upsert_worker_workdir_link(&link).unwrap(); + + assert_eq!( + store + .get_worker_registry_by_runtime("local-dev", "embedded", "browser-1") + .unwrap(), + Some(expected_worker.clone()) + ); + assert_eq!( + store + .get_workdir_registry("local-dev", "backend-2-repo") + .unwrap(), + Some(workdir.clone()) + ); + assert_eq!( + store.list_workdir_registry("local-dev", 10).unwrap(), + vec![unmanaged_workdir.clone(), workdir.clone()] + ); + assert_eq!( + store + .list_managed_workdir_registry("local-dev", 10) + .unwrap(), + vec![workdir] + ); + assert_eq!( + store + .list_worker_workdir_links("local-dev", "embedded/browser-1") + .unwrap(), + vec![link] + ); + } + fn table_names(conn: &Connection) -> BTreeSet { let mut stmt = conn .prepare( diff --git a/web/workspace/src/lib/workspace-sidebar/types.ts b/web/workspace/src/lib/workspace-sidebar/types.ts index 92eeaf98..30cd2b8e 100644 --- a/web/workspace/src/lib/workspace-sidebar/types.ts +++ b/web/workspace/src/lib/workspace-sidebar/types.ts @@ -85,6 +85,8 @@ export type Worker = { workspace: { visibility: string; identity: string }; state: string; status: string; + pinned?: boolean; + retention_state?: string; last_seen_at?: string | null; implementation: { kind: string; display_hint: string }; capabilities: WorkerCapabilities; @@ -124,6 +126,7 @@ export type WorkingDirectorySummary = { resolved_commit: string; resolved_tree?: string | null; status: string; + management_kind?: "backend_managed" | "runtime_unmanaged" | string | null; cleanup_policy: string; cleanup_target: { kind: string; @@ -144,6 +147,70 @@ export type BrowserWorkingDirectoryListResponse = { diagnostics: Diagnostic[]; }; +export type CleanupTargetKind = + | "worker_delete" + | "workdir_clean_cleanup" + | "workdir_dirty_discard" + | "workdir_record_delete"; + +export type CleanupWorkerCandidate = { + target_id: string; + action: CleanupTargetKind; + worker_id: string; + runtime_worker_id: string; + runtime_id: string; + reason: string; + blocking_reason?: string | null; + pinned: boolean; + retention_state: string; + lifecycle_state: string; + linked_workdir_ids: string[]; + running_linked: boolean; + estimated_reclaim_bytes?: number | null; +}; + +export type CleanupWorkdirCandidate = { + target_id: string; + action: CleanupTargetKind; + workdir_id: string; + runtime_id: string; + repository_id: string; + reason: string; + blocking_reason?: string | null; + linked_worker_ids: string[]; + linked_running_worker_ids: string[]; + running_linked: boolean; + pinned_linked: boolean; + file_status: string; + cleanliness: string; + estimated_reclaim_bytes?: number | null; +}; + +export type RuntimeCleanupPlanResponse = { + workspace_id: string; + runtime_id: string; + generated_at: string; + revision: string; + digest: string; + workers: CleanupWorkerCandidate[]; + workdirs: CleanupWorkdirCandidate[]; + diagnostics: Diagnostic[]; +}; + +export type RuntimeCleanupExecutionResponse = { + workspace_id: string; + runtime_id: string; + executed_at: string; + results: { + target_id: string; + action: CleanupTargetKind; + status: string; + message: string; + }[]; + plan_after: RuntimeCleanupPlanResponse; + diagnostics: Diagnostic[]; +}; + export type BrowserWorkerWorkingDirectorySelection = { working_directory_id: string; relative_cwd?: string | null; diff --git a/web/workspace/src/routes/w/[workspaceId]/runtimes/[runtimeId]/workdirs/+page.svelte b/web/workspace/src/routes/w/[workspaceId]/runtimes/[runtimeId]/workdirs/+page.svelte index a7692aab..8186640f 100644 --- a/web/workspace/src/routes/w/[workspaceId]/runtimes/[runtimeId]/workdirs/+page.svelte +++ b/web/workspace/src/routes/w/[workspaceId]/runtimes/[runtimeId]/workdirs/+page.svelte @@ -1,8 +1,16 @@ @@ -37,6 +106,68 @@ {:else if data.workdirs.items.length === 0}

No workdirs are visible for this Runtime.

{:else} +
+
+

Manual cleanup preview

+

Select explicit Workdir targets. Raw Runtime paths are intentionally not shown.

+
+ {#if data.cleanupPlanError} +

{data.cleanupPlanError}

+ {:else if cleanupCandidates.length === 0 && workerCleanupCandidates.length === 0} +

No cleanup candidates.

+ {:else} +
+ {#each workerCleanupCandidates as candidate (candidate.target_id)} + + {/each} + {#each cleanupCandidates as candidate (candidate.target_id)} + + {/each} +
+ + {#if cleanupStatus}

{cleanupStatus}

{/if} + {/if} +
+
diff --git a/web/workspace/src/routes/w/[workspaceId]/runtimes/[runtimeId]/workdirs/+page.ts b/web/workspace/src/routes/w/[workspaceId]/runtimes/[runtimeId]/workdirs/+page.ts index ab9122d0..079cdc60 100644 --- a/web/workspace/src/routes/w/[workspaceId]/runtimes/[runtimeId]/workdirs/+page.ts +++ b/web/workspace/src/routes/w/[workspaceId]/runtimes/[runtimeId]/workdirs/+page.ts @@ -3,12 +3,13 @@ import type { BrowserWorkingDirectoryListResponse, ListResponse, Runtime, + RuntimeCleanupPlanResponse, } from "$lib/workspace-sidebar/types"; import type { PageLoad } from "./$types"; export const load: PageLoad = async ({ fetch, params }) => { const runtimeId = params.runtimeId; - const [runtimes, workdirs] = await Promise.all([ + const [runtimes, workdirs, cleanupPlan] = await Promise.all([ loadJson>(fetch, workspaceApiPath(params.workspaceId, "/runtimes")), loadJson( fetch, @@ -17,6 +18,10 @@ export const load: PageLoad = async ({ fetch, params }) => { `/runtimes/${encodeURIComponent(runtimeId)}/working-directories`, ), ), + loadJson( + fetch, + workspaceApiPath(params.workspaceId, `/runtimes/${encodeURIComponent(runtimeId)}/cleanup-plan`), + ), ]); return { @@ -26,5 +31,7 @@ export const load: PageLoad = async ({ fetch, params }) => { runtimesError: runtimes.error, workdirs: workdirs.data, workdirsError: workdirs.error, + cleanupPlan: cleanupPlan.data, + cleanupPlanError: cleanupPlan.error, }; }; diff --git a/web/workspace/src/routes/w/[workspaceId]/workers/+page.svelte b/web/workspace/src/routes/w/[workspaceId]/workers/+page.svelte index 377cbae2..f1ada28d 100644 --- a/web/workspace/src/routes/w/[workspaceId]/workers/+page.svelte +++ b/web/workspace/src/routes/w/[workspaceId]/workers/+page.svelte @@ -1,9 +1,30 @@