diff --git a/.changeset/workflow-work-items.md b/.changeset/workflow-work-items.md new file mode 100644 index 0000000000..bfa8b27ff1 --- /dev/null +++ b/.changeset/workflow-work-items.md @@ -0,0 +1,5 @@ +--- +"@runfusion/fusion": patch +--- + +Add workflow work-item storage primitives for workflow-owned merge migration. diff --git a/docs/plans/2026-06-09-003-refactor-workflow-owned-merge-full-migration-slices-plan.md b/docs/plans/2026-06-09-003-refactor-workflow-owned-merge-full-migration-slices-plan.md new file mode 100644 index 0000000000..50254bac85 --- /dev/null +++ b/docs/plans/2026-06-09-003-refactor-workflow-owned-merge-full-migration-slices-plan.md @@ -0,0 +1,537 @@ +--- +title: "refactor: Workflow-owned merge full migration slices" +type: refactor +status: active +date: 2026-06-09 +depth: deep +origin: docs/plans/2026-06-09-002-refactor-workflow-owned-merge-retry-scheduling-plan.md +--- + +# refactor: Workflow-owned merge full migration slices + +## Summary + +This plan turns the workflow-owned merge/retry/scheduling architecture into a +sequence of PR-sized migration slices. The target state is unchanged from the +origin plan: workflow IR/runtime owns merge policy, retry policy, scheduling +policy, recovery routing, and git operation flow; the engine keeps substrate +responsibilities such as storage, leases, timers, process supervision, routing, +capacity, guard services, and audit plumbing. + +The migration should ship in independently reviewable slices, but the final +cutover must not leave two production control planes. Compatibility projections +are allowed while slices are in flight. Production fallback paths are removed at +the deletion gates. + +## Requirements Trace + +- R1. Built-in workflow IR expresses default merge, retry, scheduling, and + recovery policy explicitly. +- R2. Workflow work state replaces hidden merge queue and retry routing as the + policy authority. +- R3. Scheduler claims generic workflow work; it does not infer task lifecycle + advancement from columns. +- R4. Git and merge operations are workflow node capabilities guarded by shared + repository safety services. +- R5. Retry state is node/run scoped; task retry fields are compatibility + projections only. +- R6. Self-healing publishes typed workflow recovery facts and wakes recovery + nodes; it does not directly mutate merge/retry lifecycle. +- R7. Dashboard/API/CLI state derives from workflow state first. +- R8. Existing invariants remain non-configurable: `autoMerge:false`, hard + cancel on `in-progress -> todo`, file-scope/squash guards, branch-group target + rules, and user pauses. +- R9. Branch-group member integration and group promotion are workflow-owned + subgraphs with separate gates. +- R10. Final deletion tests fail if engine-owned merge/retry/scheduling policy + reappears. + +## Current Baseline + +The starting checkpoint is `docs/workflow-policy-ownership-map.md`. It classifies +today's ownership of merge queue enqueue/dequeue, scheduler in-review policy, +merge request shadow state, git/merge procedures, retry helpers, manual retry, +self-healing recovery, built-in workflow IR, and dashboard projections. + +Keep that map updated through the migration. A slice is not complete if it moves +policy without updating the map or adding the corresponding deletion gate. + +## Slice Strategy + +- Keep each slice mergeable and behavior-preserving unless the slice is an + explicit cutover gate. +- Prefer characterization-first on legacy policy before moving it. +- Introduce workflow-native state and projections before changing production + routing. +- Route one ownership surface at a time, then delete the old owner. +- Run the merge gate for every slice: `pnpm test:gate`. +- Add `pnpm lint` and `pnpm build` for every behavior-bearing slice. +- Add a changeset only when a slice changes published `@runfusion/fusion` + behavior. + +## Migration Slices + +### S0. Ownership Map And Guard + +- **Goal:** Keep the migration inventory explicit and enforced. +- **Status:** Done by the origin PR. +- **Files:** `docs/workflow-policy-ownership-map.md`, + `packages/engine/src/__tests__/workflow-policy-ownership-map.test.ts`, + `packages/engine/vitest.config.ts`. +- **Tests:** `packages/engine/src/__tests__/workflow-policy-ownership-map.test.ts`. +- **Exit gate:** Every known current policy surface is classified as + `substrate`, `workflow-policy`, `capability`, `compat-projection`, or + `delete-after-cutover`. + +### S1. Workflow Work-Item Schema And Store API + +- **Goal:** Add durable workflow work items that can represent runnable, held, + retrying, merge, manual-hold, and recovery work without changing production + routing yet. +- **Depends on:** S0. +- **Files:** `packages/core/src/db.ts`, `packages/core/src/store.ts`, + `packages/core/src/types.ts`, `packages/core/src/index.ts`, + `packages/core/src/__tests__/central-db.test.ts`, + `packages/core/src/__tests__/store-workflow-runtime.test.ts` (new), + `packages/core/src/__tests__/merge-request-record.test.ts`. +- **Decisions:** Work items are keyed by workflow run, task, node, and kind. + Minimum fields: `id`, `runId`, `taskId`, `nodeId`, `kind`, `state`, + `attempt`, `retryAfter`, `leaseOwner`, `leaseExpiresAt`, `lastError`, + `blockedReason`, `createdAt`, `updatedAt`. +- **Test scenarios:** create runnable work; create merge work; transition to + held/retrying/manual-required/succeeded/cancelled/exhausted; reclaim expired + lease; duplicate wakeups are idempotent; completed work cannot be requeued. +- **Exit gate:** Store can find due runnable work without reading task columns + for merge/retry policy. + +### S2. Merge Request Projection Onto Work Items + +- **Goal:** Project existing merge request records into workflow work-item state + so dashboards and schedulers can dual-read before cutover. +- **Depends on:** S1. +- **Files:** `packages/core/src/store.ts`, `packages/core/src/task-merge.ts`, + `packages/core/src/types.ts`, + `packages/core/src/__tests__/merge-request-record.test.ts`, + `packages/core/src/__tests__/store-workflow-runtime.test.ts`, + `packages/engine/src/__tests__/dual-observe-merge-seam.test.ts`. +- **Decisions:** Existing `mergeRequestContractShadowEnabled` remains a + compatibility switch during this slice. Work-item state is the new shape; + merge request rows remain the old projection. +- **Test scenarios:** queued/running/retrying/manual-required/succeeded rows + project to equivalent work items; task hard-cancel cancels active merge work; + exhausted merge request maps to terminal failed work; projection is + idempotent across restart. +- **Exit gate:** Every merge request state has a lossless workflow work-item + equivalent. + +### S3. Generic Scheduler Claim Path + +- **Goal:** Teach `Scheduler` to claim due workflow work items while preserving + existing task dispatch behavior. +- **Depends on:** S1. +- **Files:** `packages/engine/src/scheduler.ts`, + `packages/engine/src/workflow-task-runtime.ts`, + `packages/engine/src/project-engine.ts`, + `packages/engine/src/__tests__/scheduler.test.ts`, + `packages/engine/src/__tests__/scheduler-node-routing.test.ts`, + `packages/engine/src/__tests__/workflow-work-engine-dispatch.test.ts` (new). +- **Decisions:** Scheduler remains substrate. It may apply capacity, routing, + leases, global pause, engine pause, and remote-node dispatch. It must not own + merge eligibility, retry routing, or recovery outcome. +- **Test scenarios:** claim only due runnable work; skip `retryAfter` until due; + hold on capacity without task mutation; user-paused work is not claimed; stale + leases are reclaimable; remote node receives workflow runtime work. +- **Exit gate:** A workflow work item can be dispatched end to end in tests + without constructing a merge queue branch. + +### S4. Built-In Merge/Retry/Recovery IR Regions + +- **Goal:** Add explicit merge, retry, manual hold, branch-group, and recovery + regions to built-in workflow IR. +- **Depends on:** S1, S2. +- **Files:** `packages/core/src/builtin-coding-workflow-ir.ts`, + `packages/core/src/builtin-stepwise-coding-workflow-ir.ts`, + `packages/core/src/builtin-pr-workflow-ir.ts`, + `packages/core/src/builtin-workflows.ts`, + `packages/core/src/workflow-ir-types.ts`, + `packages/core/src/__tests__/builtin-coding-workflow-ir.test.ts`, + `packages/core/src/__tests__/builtin-stepwise-coding-workflow-ir.test.ts`, + `packages/core/src/__tests__/builtin-pr-workflow-ir.test.ts`. +- **Decisions:** Use built-in node kinds for non-authorable primitives: + merge gate, merge attempt, manual merge hold, retry/backoff, branch-group + member integration, group promotion, finalize, and recovery router. +- **Test scenarios:** built-in workflows validate; default coding has a merge + gate; stepwise coding has per-step retry plus merge retry; PR workflow routes + review/fix/merge; `autoMerge:false` routes to manual hold; branch-group member + integration and group promotion are separate nodes. +- **Exit gate:** Built-in IR is the source of truth for all default + merge/retry/recovery policy, even if production handlers are not wired yet. + +### S5. Runtime Work-Item Driver + +- **Goal:** Let `WorkflowTaskRuntime` start from a workflow work item and persist + node/work-item outcomes. +- **Depends on:** S1, S3, S4. +- **Files:** `packages/engine/src/workflow-task-runtime.ts`, + `packages/engine/src/workflow-graph-executor.ts`, + `packages/engine/src/workflow-node-handlers.ts`, + `packages/engine/src/__tests__/workflow-task-runtime.test.ts`, + `packages/engine/src/__tests__/workflow-graph-executor-retry-coding-workflow.test.ts`, + `packages/engine/src/__tests__/workflow-node-handlers.test.ts`. +- **Decisions:** Runtime receives `{ workItemId, runId, taskId, nodeId }` and + returns a typed outcome that updates work item state. Task column updates are + side effects of workflow primitives, not scheduler policy. +- **Test scenarios:** runnable work completes; failing node creates retrying + work; manual hold node creates held work; runtime restart resumes from stored + work; duplicate start of same work item is refused by lease. +- **Exit gate:** Runtime can progress workflow work without old merge queue + callbacks. + +### S6. Git And Merge Capability Extraction + +- **Goal:** Put checkout preparation, branch integration, merge attempt, squash, + finalize, and conflict classification behind workflow node capability modules. +- **Depends on:** S4, S5. +- **Files:** `packages/engine/src/merger.ts`, + `packages/engine/src/merger-ai.ts`, + `packages/engine/src/merger-integration-worktree.ts`, + `packages/engine/src/workflow-merge-nodes.ts` (new), + `packages/engine/src/workflow-node-handlers.ts`, + `packages/engine/src/merge-trait.ts`, + `packages/engine/src/__tests__/interpreter-merge-seam.test.ts`, + `packages/engine/src/__tests__/dual-observe-merge-seam.test.ts`, + `packages/engine/src/__tests__/workflow-merge-nodes.test.ts` (new). +- **Decisions:** This slice does not rewrite low-level merge algorithms. It + extracts orchestration boundaries so workflow nodes call existing guarded + operations. +- **Test scenarios:** merge node calls checkout preparation; file-scope + violation returns workflow failure; already-on-main routes to finalize; + transient merge error returns retry outcome; non-transient conflict routes to + revision/manual hold; no production caller can bypass guard service in tests. +- **Exit gate:** A merge attempt can be driven by a workflow node capability in + tests with the same guard behavior as `merger.ts`. + +### S7. Completion Handoff Creates Merge Work + +- **Goal:** Replace task-moved `in-review` auto-enqueue as the policy authority + with workflow completion handoff creating merge work. +- **Depends on:** S2, S5, S6. +- **Files:** `packages/engine/src/project-engine.ts`, + `packages/engine/src/merger.ts`, + `packages/core/src/store.ts`, + `packages/engine/src/__tests__/workflow-interpreter-cutover.test.ts`, + `packages/engine/src/__tests__/completion-fanout-x-self-healing.test.ts`, + `packages/engine/src/__tests__/merge-reuse-task-worktree.slow.test.ts`. +- **Decisions:** During this slice the old queue can remain as a projection, but + merge work creation happens through workflow handoff. `autoMerge:false` creates + a manual hold work item. +- **Test scenarios:** coding completion creates merge work; `autoMerge:false` + creates manual hold and does not enqueue merge; duplicate handoff is + idempotent; soft-deleted task cancels handoff; startup projection does not + create duplicate merge work. +- **Exit gate:** New task completions produce workflow merge work before any old + queue processing path runs. + +### S8. Workflow-Owned Merge Queue Processing + +- **Goal:** Process merge work items through workflow runtime instead of + `ProjectEngine`'s in-memory merge queue loop. +- **Depends on:** S3, S6, S7. +- **Files:** `packages/engine/src/project-engine.ts`, + `packages/engine/src/scheduler.ts`, + `packages/engine/src/merger.ts`, + `packages/core/src/store.ts`, + `packages/engine/src/__tests__/merger-merge-lifecycle.test.ts`, + `packages/engine/src/__tests__/merger-post-merge.test.ts`, + `packages/engine/src/__tests__/workflow-work-engine-dispatch.test.ts`, + `packages/engine/src/__tests__/workflow-merge-nodes.test.ts`. +- **Decisions:** Keep queue fairness and serialization as substrate leases. The + policy route after success/failure belongs to workflow node outcomes. +- **Test scenarios:** queued merge work claims one at a time; successful merge + finalizes task; transient failure schedules retrying merge work; permanent + conflict routes to revision/manual hold; active merge lease blocks duplicate + processing; hard cancel cancels running merge work. +- **Exit gate:** Production merge processing no longer depends on a hidden + `mergeQueue` dequeue loop. + +### S9. Workflow-Owned Retry State + +- **Goal:** Move retry attempts, budgets, backoff, retry-after, exhaustion, and + manual retry reset into workflow node/work-item state. +- **Depends on:** S5, S8. +- **Files:** `packages/engine/src/workflow-graph-executor.ts`, + `packages/engine/src/workflow-node-handlers.ts`, + `packages/engine/src/retry-with-backoff.ts`, + `packages/engine/src/rate-limit-retry.ts`, + `packages/engine/src/transient-merge-error-classifier.ts`, + `packages/core/src/retry-summary.ts`, + `packages/core/src/manual-retry-reset.ts`, + `packages/engine/src/__tests__/workflow-node-retry-policy.test.ts` (new), + `packages/core/src/__tests__/manual-retry-reset.test.ts`. +- **Decisions:** Task retry fields remain as derived display summaries until + deletion. Manual retry emits a workflow wake and clears only targeted failed + node state. +- **Test scenarios:** implementation node retry stays within budget; merge node + retry does not reset implementation progress; rate-limit error persists due + time; exhausted retry routes to failure/manual hold; manual retry clears only + failed node; retry state survives restart. +- **Exit gate:** No retry branch is controlled solely by task counters. + +### S10. Self-Healing Recovery Events + +- **Goal:** Convert self-healing merge/retry lifecycle mutations into typed + workflow recovery events and node wakes. +- **Depends on:** S5, S8, S9. +- **Files:** `packages/engine/src/self-healing.ts`, + `packages/engine/src/restart-recovery-coordinator.ts`, + `packages/engine/src/recovery-policy.ts`, + `packages/engine/src/workflow-task-runtime.ts`, + `packages/engine/src/__tests__/self-healing.test.ts`, + `packages/engine/src/__tests__/workflow-recovery-events.test.ts` (new), + `packages/engine/src/__tests__/reliability-interactions/in-review-automerge-off.test.ts`, + `packages/engine/src/__tests__/reliability-interactions/workflow-interpreter-cutover.test.ts`. +- **Decisions:** Sweeps detect facts. Recovery nodes decide routes. Non-task + agent/heartbeat cleanup may remain engine-owned when it is not task lifecycle + policy. +- **Test scenarios:** mergeable in-review task gets recovery event, not direct + requeue; stale merge status emits event; transient merge failure emits retry + event; already landed emits finalize event; `autoMerge:false` remains terminal; + duplicate recovery events are deduped. +- **Exit gate:** Self-healing no longer directly requeues, pauses, fails, + unpauses, or moves merge/retry tasks except through guarded workflow + primitives. + +### S11. Branch Group Workflow Subgraphs + +- **Goal:** Move branch-group member integration and group promotion into + workflow-owned merge subgraphs. +- **Depends on:** S6, S8, S10. +- **Files:** `packages/engine/src/group-merge-coordinator.ts`, + `packages/engine/src/merge-trait.ts`, + `packages/engine/src/merger-integration-worktree.ts`, + `packages/core/src/builtin-coding-workflow-ir.ts`, + `packages/engine/src/__tests__/reliability-interactions/shared-branch-group-lifecycle.slow.test.ts`, + `packages/engine/src/__tests__/workflow-branch-group-merge.test.ts` (new). +- **Decisions:** Member-to-shared-branch integration and shared-branch-to-default + promotion are distinct workflow nodes with distinct auto-merge gates. +- **Test scenarios:** shared member integrates while global auto-merge is off + under the scoped exception; group promotion remains blocked when group/global + auto-merge is off; conflicting member integration routes to recovery/revision; + final group promotion runs file-scope and squash guards. +- **Exit gate:** Branch-group coordinator no longer owns task lifecycle + independent of workflow runtime. + +### S12. Dashboard/API/CLI Workflow Projection + +- **Goal:** Surface workflow-native queued, retrying, merging, manual-hold, + failed, stalled, and recovered reasons across user inspection surfaces. +- **Depends on:** S1, S2, S7, S9, S10. +- **Files:** `packages/dashboard/app/components/TaskCard.tsx`, + task detail components, reliability views, task API routes, + CLI task output files, `packages/core/src/retry-summary.ts`, + `packages/core/src/task-merge.ts`, + `packages/dashboard/app/components/__tests__/TaskCard.test.tsx`, + reliability/dashboard API tests. +- **Decisions:** Workflow state wins over stale task fields. Legacy task fields + remain fallback for old rows only. +- **Test scenarios:** merge queued shows workflow merge work, not stalled; + retrying shows attempt and due time; manual hold shows human action required; + recovery event shows reason; completed work hides stale stalled badges; + branch-group merge work identifies target branch. +- **Exit gate:** UI/API/CLI tests prove workflow state is the first projection + source. + +### S13. Scheduler Policy Deletion + +- **Goal:** Delete scheduler branches that infer lifecycle, merge eligibility, + retry routing, or in-review dependency behavior from task columns. +- **Depends on:** S3, S7, S8, S12. +- **Files:** `packages/engine/src/scheduler.ts`, + `packages/core/src/task-merge.ts`, + `packages/engine/src/__tests__/scheduler.test.ts`, + `packages/engine/src/__tests__/scheduler-overlap-requeue.test.ts`, + `packages/engine/src/__tests__/workflow-scheduler-policy-deletion.test.ts` (new). +- **Decisions:** Dependency satisfaction should use completion handoff/workflow + state. In-review scope leases are replaced by workflow work leases and guard + services. +- **Test scenarios:** scheduler cannot satisfy dependency only because a task is + `in-review`; retry due time comes from work item; overlap lease comes from + workflow work; PR monitor behavior remains as watch substrate, not lifecycle + owner. +- **Exit gate:** Search/structure test fails if scheduler reintroduces + task-column merge/retry policy. + +### S14. ProjectEngine Merge Queue Deletion + +- **Goal:** Remove production `ProjectEngine` merge queue policy and retain only + explicit human/manual event entry points plus substrate helpers. +- **Depends on:** S8, S11, S13. +- **Files:** `packages/engine/src/project-engine.ts`, + `packages/engine/src/runtimes/in-process-runtime.ts`, + `packages/core/src/store.ts`, + `packages/engine/src/__tests__/merger-merge-lifecycle.test.ts`, + `packages/engine/src/__tests__/workflow-merge-policy-deletion.test.ts` (new). +- **Decisions:** Manual merge APIs record a workflow event or create due workflow + work; they do not enqueue hidden engine work. +- **Test scenarios:** no startup in-review scan enqueues hidden merge work; + unpause wakes workflow work; manual merge event wakes merge node; stale + `mergeActive` state cannot block workflow work; old queue APIs are absent or + compatibility-only. +- **Exit gate:** No production caller starts merge processing outside workflow + runtime. + +### S15. Self-Healing Policy Deletion + +- **Goal:** Delete self-healing direct lifecycle mutations for merge/retry tasks + after recovery events cover all cases. +- **Depends on:** S10, S11, S14. +- **Files:** `packages/engine/src/self-healing.ts`, + `docs/self-healing-backward-move-audit.md`, + `packages/engine/src/__tests__/self-healing.test.ts`, + `packages/engine/src/__tests__/workflow-recovery-events.test.ts`, + `packages/engine/src/__tests__/workflow-self-healing-policy-deletion.test.ts` (new). +- **Decisions:** Metadata reconciliation and non-task agent cleanup can remain. + Task lifecycle repair becomes recovery events plus workflow node outcomes. +- **Test scenarios:** direct calls to `moveTask(..., "todo")`, + `updateTask({ paused: true })`, merge requeue callbacks, and merge retry resets + are absent for merge/retry surfaces; valid held states are no-ops; recovery + facts carry audit context. +- **Exit gate:** Search tests fail on direct self-healing merge/retry lifecycle + mutation patterns. + +### S16. Legacy Retry Field Demotion + +- **Goal:** Demote task-level retry/merge counters to projections and remove + policy reads that still treat them as authority. +- **Depends on:** S9, S12, S15. +- **Files:** `packages/core/src/types.ts`, `packages/core/src/retry-summary.ts`, + `packages/core/src/manual-retry-reset.ts`, `packages/core/src/store.ts`, + `packages/engine/src/project-engine.ts`, `packages/engine/src/self-healing.ts`, + `packages/core/src/__tests__/manual-retry-reset.test.ts`, + `packages/engine/src/__tests__/workflow-node-retry-policy.test.ts`. +- **Decisions:** Do not remove fields until all compatibility surfaces can read + workflow projections. Removal can be a later cleanup; this slice removes policy + authority. +- **Test scenarios:** retry summaries derive from workflow node/work state; + manual retry emits workflow wake; old task fields changing alone cannot cause + scheduler/recovery/merge action. +- **Exit gate:** Task retry fields are display-only compatibility data. + +### S17. End-To-End Cutover Matrix + +- **Goal:** Prove the full workflow-owned invariant across all known production + surfaces before removing dual-read compatibility. +- **Depends on:** S13, S14, S15, S16. +- **Files:** focused tests across `packages/engine/src/__tests__/`, + reliability interactions under + `packages/engine/src/__tests__/reliability-interactions/`, core store tests, + dashboard projection tests, `docs/testing.md`. +- **Test matrix:** default coding auto-merge; stepwise coding; custom workflow; + PR workflow; plugin workflow extension; `autoMerge:false`; manual retry; + user hard cancel; engine restart during merge work; transient merge failure; + permanent conflict; branch-group member integration; branch-group promotion; + stale recovery; already-landed finalization; dashboard task card/detail; + CLI task output. +- **Exit gate:** `pnpm test:gate`, `pnpm lint`, `pnpm build`, and targeted matrix + suites pass. No old engine merge/retry/scheduling policy path can race workflow + runtime in production. + +### S18. Documentation, Settings, And Release Notes + +- **Goal:** Update architecture, settings, dashboard, CLI, and testing docs for + workflow-owned policy and compatibility projections. +- **Depends on:** S17. +- **Files:** `docs/architecture.md`, `docs/workflow-steps.md`, + `docs/dashboard-guide.md`, `docs/settings-reference.md`, `docs/testing.md`, + `CONCEPTS.md`, `.changeset/.md`. +- **Decisions:** Document the new source of truth, remaining compatibility fields, + recovery event vocabulary, manual hold behavior, branch-group routing, and + deletion gates. +- **Test scenarios:** docs inventory/search tests if applicable; lazy view + inventory unchanged unless dashboard imports change. +- **Exit gate:** User-facing docs use the same state names as API/UI tests, and a + patch changeset exists if published `@runfusion/fusion` behavior changed. + +## Dependency Graph + +```mermaid +flowchart TB + S0 --> S1 + S1 --> S2 + S1 --> S3 + S2 --> S4 + S3 --> S5 + S4 --> S5 + S5 --> S6 + S6 --> S7 + S7 --> S8 + S8 --> S9 + S9 --> S10 + S8 --> S11 + S10 --> S11 + S7 --> S12 + S9 --> S12 + S10 --> S12 + S12 --> S13 + S13 --> S14 + S11 --> S14 + S14 --> S15 + S15 --> S16 + S16 --> S17 + S17 --> S18 +``` + +## Release And Merge Strategy + +- **Preferred PR count:** 18 slices, one PR per slice. +- **Can combine:** S1+S2 if schema and projection are small; S13+S14 if deletion + is purely mechanical after S8. +- **Do not combine:** S8 with S14, or S10 with S15. Route through workflow first, + then delete old owner in a separate reviewable PR. +- **Branch policy:** Each slice branches from current `main`, not from a stale + feature stack. Drop duplicate commits before merging. +- **Changesets:** Add only when published CLI behavior changes. Internal docs, + CI config, and behavior-preserving refactors do not require changesets. + +## Cutover Safety Gates + +- Gate A after S4: built-in IR expresses all planned policy regions. +- Gate B after S8: workflow runtime can process merge work without hidden queue + ownership. +- Gate C after S10: self-healing emits recovery events for merge/retry surfaces. +- Gate D after S12: dashboard/API/CLI read workflow state first. +- Gate E after S17: deletion tests and end-to-end matrix prove no production + legacy control plane remains. + +## Rollback Strategy + +- Before S13, rollback is disabling workflow work dispatch and relying on legacy + projections. +- After S13, rollback is revert-by-slice, not runtime fallback. Do not ship a + production dual-controller fallback after deletion gates begin. +- Keep old fields as compatibility projections through S17 so data downgrade is + not required for ordinary slice rollback. + +## Verification Commands + +- `pnpm test:gate` +- `pnpm lint` +- `pnpm build` +- Targeted suites named in each slice. +- `pnpm test:full` only for explicit final matrix verification or release-adjacent + confidence, not as the normal merge gate. + +## Done Criteria + +- Workflow work items are the durable source of runnable, held, retrying, merge, + and recovery work. +- Built-in workflows express default merge, retry, scheduling, branch-group, and + recovery policy. +- Scheduler dispatches generic workflow work only. +- Git/merge operations are invoked by workflow nodes or explicit human/manual + events recorded as workflow events. +- Retry budgets and manual retry reset are node/run scoped. +- Self-healing publishes recovery facts and wakes workflow recovery nodes. +- Dashboard/API/CLI projections read workflow state first. +- Deletion tests prevent reintroducing engine-owned merge/retry/scheduling + policy. diff --git a/packages/core/src/__tests__/db-migrate.test.ts b/packages/core/src/__tests__/db-migrate.test.ts index fd69f5c1b0..3b06532f27 100644 --- a/packages/core/src/__tests__/db-migrate.test.ts +++ b/packages/core/src/__tests__/db-migrate.test.ts @@ -715,8 +715,8 @@ describe("schema migration", () => { const row = db.prepare("SELECT deletedAt FROM tasks WHERE id = 'FN-legacy'").get() as { deletedAt: string | null }; expect(row.deletedAt).toBeNull(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -749,8 +749,8 @@ describe("schema migration", () => { { id: "WS-001", mode: "prompt", gateMode: "advisory" }, { id: "WS-002", mode: "script", gateMode: "advisory" }, ]); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -800,8 +800,8 @@ describe("schema migration", () => { reviewerContextRetryCount: 0, reviewerFallbackRetryCount: 0, }); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -830,8 +830,8 @@ describe("schema migration", () => { const columns = db.prepare("PRAGMA table_info(milestones)").all() as Array<{ name: string }>; expect(columns.map((column) => column.name)).toContain("acceptanceCriteria"); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -872,8 +872,8 @@ describe("schema migration", () => { const missionColumns = db.prepare("PRAGMA table_info(missions)").all() as Array<{ name: string }>; expect(missionColumns.map((column) => column.name)).toContain("autoMerge"); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -907,8 +907,8 @@ describe("schema migration", () => { { id: "WS-002", mode: "script", enabled: 1, gateMode: "advisory" }, { id: "WS-003", mode: "prompt", enabled: 0, gateMode: "advisory" }, ]); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -945,8 +945,8 @@ describe("schema migration", () => { const indexes = db.prepare("PRAGMA index_list(mission_goals)").all() as Array<{ name: string }>; expect(indexes.some((index) => index.name === "idxMissionGoalsGoalId")).toBe(true); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -1007,7 +1007,7 @@ describe("schema migration", () => { expect(customFieldsColumn).toBeDefined(); expect(customFieldsColumn?.dflt_value).toBe("'{}'"); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -1045,7 +1045,7 @@ describe("schema migration", () => { const indexes = db.prepare("PRAGMA index_list(workflow_settings)").all() as Array<{ name: string }>; expect(indexes.some((index) => index.name === "idx_workflow_settings_project")).toBe(true); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -1127,7 +1127,7 @@ describe("schema migration", () => { expect(indexNames).toContain("idx_cli_sessions_chatSessionId"); expect(indexNames).toContain("idx_cli_sessions_project_state"); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -1159,7 +1159,7 @@ describe("schema migration", () => { .all() as Array<{ name: string }>; expect(columns.map((column) => column.name)).toContain("cliExecutorAdapterId"); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -1169,7 +1169,7 @@ describe("schema migration", () => { const tables = db.prepare("SELECT name FROM sqlite_master WHERE type='table'").all() as Array<{ name: string }>; expect(tables.map((row) => row.name)).toContain("cli_sessions"); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -1226,23 +1226,23 @@ describe("schema migration", () => { .get() as { migrated_fragment_id: string | null }; expect(stepRow.migrated_fragment_id).toBeNull(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); it("migration 109 is idempotent on re-init", () => { const db = new Database(fusionDir); db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); db.close(); // Re-open the same on-disk DB: already at 109, the 109 block must be a no-op. const reopened = new Database(fusionDir); reopened.init(); - expect(reopened.getSchemaVersion()).toBe(114); - expect(reopened.getSchemaVersion()).toBe(114); + expect(reopened.getSchemaVersion()).toBe(115); + expect(reopened.getSchemaVersion()).toBe(115); const workflowColumns = reopened.prepare("PRAGMA table_info(workflows)").all() as Array<{ name: string }>; expect(workflowColumns.filter((c) => c.name === "kind")).toHaveLength(1); const stepColumns = reopened.prepare("PRAGMA table_info(workflow_steps)").all() as Array<{ name: string }>; diff --git a/packages/core/src/__tests__/db.test.ts b/packages/core/src/__tests__/db.test.ts index 23d1f377b7..a4f89a3caa 100644 --- a/packages/core/src/__tests__/db.test.ts +++ b/packages/core/src/__tests__/db.test.ts @@ -334,8 +334,8 @@ describe("Database", () => { }); it("seeds schema version", () => { - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); }); it("includes tokenUsageCacheWriteTokens on freshly initialized tasks table", () => { @@ -394,8 +394,8 @@ describe("Database", () => { it("is idempotent - calling init() twice does not fail", () => { expect(() => db.init()).not.toThrow(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); }); it("does not overwrite existing config on re-init", () => { // Update the config @@ -1465,8 +1465,8 @@ describe("schema migrations", () => { db.init(); // Verify version bumped to 29 (includes v1→v2 through v26→v29) - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); // Verify new columns exist and existing data is intact const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>; @@ -1491,16 +1491,16 @@ describe("schema migrations", () => { const db = new Database(fusionDir); db.init(); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); // Re-init should not fail db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); // Re-init should not fail db.init(); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); db.close(); }); @@ -1535,8 +1535,8 @@ describe("schema migrations", () => { db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>; expect(cols.map((col) => col.name)).toContain("priority"); @@ -1577,8 +1577,8 @@ describe("schema migrations", () => { db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>; const colNames = cols.map((col) => col.name); @@ -1650,8 +1650,8 @@ describe("schema migrations", () => { db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>; const colNames = cols.map((col) => col.name); @@ -1891,8 +1891,8 @@ describe("schema migrations", () => { db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); const cols = db.prepare("PRAGMA table_info(chat_messages)").all() as Array<{ name: string }>; expect(cols.map((col) => col.name)).toContain("attachments"); @@ -1966,8 +1966,8 @@ describe("schema migrations", () => { db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); const tables = db.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name = 'agentRatings'").all() as Array<{ name: string }>; expect(tables).toEqual([{ name: "agentRatings" }]); @@ -1991,8 +1991,8 @@ describe("schema migrations", () => { db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); const tables = db.prepare("SELECT name FROM sqlite_master WHERE type='table' AND name = 'mission_events'").all() as Array<{ name: string }>; expect(tables).toEqual([{ name: "mission_events" }]); @@ -2096,8 +2096,8 @@ describe("schema migrations", () => { db.init(); // Verify version bumped to 29 - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); // Verify new columns exist and existing data is intact const cols = db.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>; @@ -2316,8 +2316,8 @@ describe("schema migrations", () => { localDb.init(); - expect(localDb.getSchemaVersion()).toBe(114); - expect(localDb.getSchemaVersion()).toBe(114); + expect(localDb.getSchemaVersion()).toBe(115); + expect(localDb.getSchemaVersion()).toBe(115); const columns = localDb.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>; expect(columns.map((column) => column.name)).toContain("tokenUsageCacheWriteTokens"); @@ -2628,8 +2628,8 @@ describe("createDatabase factory", () => { const db = createDatabase(fusionDir); db.init(); - expect(db.getSchemaVersion()).toBe(114); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); + expect(db.getSchemaVersion()).toBe(115); expect(db.getLastModified()).toBeGreaterThan(0); db.close(); @@ -2783,8 +2783,8 @@ describe("migration v77 task token budget columns", () => { migrated = new Database(fusion); migrated.init(); - expect(migrated.getSchemaVersion()).toBe(114); - expect(migrated.getSchemaVersion()).toBe(114); + expect(migrated.getSchemaVersion()).toBe(115); + expect(migrated.getSchemaVersion()).toBe(115); const rows = migrated.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>; const names = new Set(rows.map((row) => row.name)); expect(names.has("tokenBudgetSoftAlertedAt")).toBe(true); @@ -2815,8 +2815,8 @@ describe("migration v106 adds tasks.transitionPending (FN-1417)", () => { const fresh = new Database(fusion); try { fresh.init(); - expect(fresh.getSchemaVersion()).toBe(114); - expect(fresh.getSchemaVersion()).toBe(114); + expect(fresh.getSchemaVersion()).toBe(115); + expect(fresh.getSchemaVersion()).toBe(115); const names = new Set( (fresh.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>).map((r) => r.name), ); @@ -2844,8 +2844,8 @@ describe("migration v106 adds tasks.transitionPending (FN-1417)", () => { migrated = new Database(fusion); migrated.init(); - expect(migrated.getSchemaVersion()).toBe(114); - expect(migrated.getSchemaVersion()).toBe(114); + expect(migrated.getSchemaVersion()).toBe(115); + expect(migrated.getSchemaVersion()).toBe(115); const names = new Set( (migrated.prepare("PRAGMA table_info(tasks)").all() as Array<{ name: string }>).map((r) => r.name), ); @@ -2871,8 +2871,8 @@ describe("migration v107 adds workflow_run_branches + index (FN-1417)", () => { const fresh = new Database(fusion); try { fresh.init(); - expect(fresh.getSchemaVersion()).toBe(114); - expect(fresh.getSchemaVersion()).toBe(114); + expect(fresh.getSchemaVersion()).toBe(115); + expect(fresh.getSchemaVersion()).toBe(115); const table = fresh .prepare("SELECT name FROM sqlite_master WHERE type='table' AND name = 'workflow_run_branches'") .get() as { name: string } | undefined; @@ -2906,8 +2906,8 @@ describe("migration v107 adds workflow_run_branches + index (FN-1417)", () => { migrated = new Database(fusion); migrated.init(); - expect(migrated.getSchemaVersion()).toBe(114); - expect(migrated.getSchemaVersion()).toBe(114); + expect(migrated.getSchemaVersion()).toBe(115); + expect(migrated.getSchemaVersion()).toBe(115); const table = migrated .prepare("SELECT name FROM sqlite_master WHERE type='table' AND name = 'workflow_run_branches'") .get() as { name: string } | undefined; @@ -2948,8 +2948,8 @@ describe("migration v67 drops orphan project auth tables", () => { migrated = new Database(fusion); migrated.init(); - expect(migrated.getSchemaVersion()).toBe(114); - expect(migrated.getSchemaVersion()).toBe(114); + expect(migrated.getSchemaVersion()).toBe(115); + expect(migrated.getSchemaVersion()).toBe(115); const tables = migrated .prepare("SELECT name FROM sqlite_master WHERE type='table' AND name LIKE 'project_auth_%'") .all() as Array<{ name: string }>; @@ -2976,8 +2976,8 @@ describe("migration v67 drops orphan project auth tables", () => { try { fresh.init(); - expect(fresh.getSchemaVersion()).toBe(114); - expect(fresh.getSchemaVersion()).toBe(114); + expect(fresh.getSchemaVersion()).toBe(115); + expect(fresh.getSchemaVersion()).toBe(115); const tables = fresh .prepare("SELECT name FROM sqlite_master WHERE type='table' AND name LIKE 'project_auth_%'") .all() as Array<{ name: string }>; diff --git a/packages/core/src/__tests__/goals-schema.test.ts b/packages/core/src/__tests__/goals-schema.test.ts index d1859c8df3..5ad25567d1 100644 --- a/packages/core/src/__tests__/goals-schema.test.ts +++ b/packages/core/src/__tests__/goals-schema.test.ts @@ -91,6 +91,6 @@ describe("goals schema", () => { }); it("reports schema version 101", () => { - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); }); }); diff --git a/packages/core/src/__tests__/insight-store.test.ts b/packages/core/src/__tests__/insight-store.test.ts index cc42db9ec3..cd834d81a9 100644 --- a/packages/core/src/__tests__/insight-store.test.ts +++ b/packages/core/src/__tests__/insight-store.test.ts @@ -1000,7 +1000,7 @@ describe("Migration: pre-33 DB upgrade", () => { // Step 1: Create a fresh database at v33 (runs all migrations up to 33) const db1 = createDatabase(legacyDir); db1.init(); - expect(db1.getSchemaVersion()).toBe(114); + expect(db1.getSchemaVersion()).toBe(115); db1.close(); // Step 2: Manually downgrade to version 32 and drop insight tables @@ -1035,7 +1035,7 @@ describe("Migration: pre-33 DB upgrade", () => { expect(tableNamesBefore).not.toContain("project_insight_runs"); // Now run init — this triggers the v32→v33 migration db3.init(); - expect(db3.getSchemaVersion()).toBe(114); + expect(db3.getSchemaVersion()).toBe(115); // Step 4: Verify insight tables exist after migration const tablesAfter = db3.prepare( @@ -1066,12 +1066,12 @@ describe("Migration: pre-33 DB upgrade", () => { try { const db1 = createDatabase(testDir); db1.init(); - expect(db1.getSchemaVersion()).toBe(114); + expect(db1.getSchemaVersion()).toBe(115); db1.close(); const db2 = createDatabase(testDir); expect(() => db2.init()).not.toThrow(); - expect(db2.getSchemaVersion()).toBe(114); + expect(db2.getSchemaVersion()).toBe(115); db2.close(); } finally { rmSync(testDir, { recursive: true, force: true }); @@ -1085,7 +1085,7 @@ describe("Migration: pre-33 DB upgrade", () => { // Step 1: Create a fresh DB and run migrations const db1 = createDatabase(compatDir); db1.init(); - expect(db1.getSchemaVersion()).toBe(114); + expect(db1.getSchemaVersion()).toBe(115); // Step 2: Strip lifecycle and cancelledAt columns by recreating the // table without them. This simulates a DB that was created before the diff --git a/packages/core/src/__tests__/merge-request-record.test.ts b/packages/core/src/__tests__/merge-request-record.test.ts index 364dbdeac8..018b79c15e 100644 --- a/packages/core/src/__tests__/merge-request-record.test.ts +++ b/packages/core/src/__tests__/merge-request-record.test.ts @@ -38,7 +38,7 @@ describe("TaskStore merge request record + completion handoff marker", () => { .all() as Array<{ name: string }>; expect(tableRows).toEqual([{ name: "completion_handoff_markers" }, { name: "merge_requests" }]); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); }); it("upserts merge request records", async () => { diff --git a/packages/core/src/__tests__/mission-store.test.ts b/packages/core/src/__tests__/mission-store.test.ts index 52fdebcac2..0f3d2b09c5 100644 --- a/packages/core/src/__tests__/mission-store.test.ts +++ b/packages/core/src/__tests__/mission-store.test.ts @@ -3746,7 +3746,7 @@ describe("MissionStore", () => { describe("Loop State & Validator Run Schema (v31)", () => { it("schema version is 101 after migration", () => { - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); }); it("mission_features table has loop state columns", () => { diff --git a/packages/core/src/__tests__/run-audit.test.ts b/packages/core/src/__tests__/run-audit.test.ts index ce0035b235..087f71d8fc 100644 --- a/packages/core/src/__tests__/run-audit.test.ts +++ b/packages/core/src/__tests__/run-audit.test.ts @@ -584,7 +584,7 @@ describe("Run Audit", () => { }); it("schema version is bumped to 40", () => { - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); }); }); }); diff --git a/packages/core/src/__tests__/store-merge-queue.test.ts b/packages/core/src/__tests__/store-merge-queue.test.ts index 6f5fd1551c..2db85dc1d2 100644 --- a/packages/core/src/__tests__/store-merge-queue.test.ts +++ b/packages/core/src/__tests__/store-merge-queue.test.ts @@ -60,7 +60,7 @@ describe("TaskStore merge queue", () => { expect.arrayContaining(["idx_mergeQueue_lease_ready", "idx_mergeQueue_leaseExpiresAt"]), ); - expect(store.getDatabase().getSchemaVersion()).toBe(114); + expect(store.getDatabase().getSchemaVersion()).toBe(115); }); it("migrates a legacy v88 database and preserves task rows", async () => { diff --git a/packages/core/src/__tests__/store-workflow-runtime.test.ts b/packages/core/src/__tests__/store-workflow-runtime.test.ts new file mode 100644 index 0000000000..0bee6ddf5f --- /dev/null +++ b/packages/core/src/__tests__/store-workflow-runtime.test.ts @@ -0,0 +1,232 @@ +import { afterEach, beforeEach, describe, expect, it } from "vitest"; +import { mkdtempSync } from "node:fs"; +import { rm } from "node:fs/promises"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { SCHEMA_VERSION } from "../db.js"; +import { TaskStore } from "../store.js"; + +function makeTmpDir(): string { + return mkdtempSync(join(tmpdir(), "kb-workflow-runtime-test-")); +} + +describe("TaskStore workflow work items", () => { + let rootDir: string; + let globalDir: string; + let store: TaskStore; + + beforeEach(async () => { + rootDir = makeTmpDir(); + globalDir = join(rootDir, ".fusion-global"); + store = new TaskStore(rootDir, globalDir); + await store.init(); + }); + + afterEach(async () => { + store.close(); + await rm(rootDir, { recursive: true, force: true, maxRetries: 5, retryDelay: 50 }); + }); + + async function createTaskId(): Promise { + const task = await store.createTask({ description: "workflow work item test" }); + return task.id; + } + + it("creates workflow work-item tables on fresh schema", () => { + const db = store.getDatabase(); + const table = db + .prepare("SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'workflow_work_items'") + .get() as { name: string } | undefined; + const indexes = db + .prepare("SELECT name FROM sqlite_master WHERE type = 'index' AND tbl_name = 'workflow_work_items' ORDER BY name") + .all() as Array<{ name: string }>; + + expect(table).toEqual({ name: "workflow_work_items" }); + expect(indexes.map((row) => row.name)).toEqual( + expect.arrayContaining([ + "idx_workflow_work_items_due", + "idx_workflow_work_items_leaseExpiresAt", + "idx_workflow_work_items_task_run", + ]), + ); + expect(db.getSchemaVersion()).toBe(SCHEMA_VERSION); + }); + + it("upserts by run, task, node, and kind without duplicating work", async () => { + const taskId = await createTaskId(); + + const created = store.upsertWorkflowWorkItem({ + runId: "run-1", + taskId, + nodeId: "merge.node", + kind: "merge", + now: "2026-06-09T00:00:00.000Z", + }); + const updated = store.upsertWorkflowWorkItem({ + runId: "run-1", + taskId, + nodeId: "merge.node", + kind: "merge", + state: "held", + blockedReason: "shared branch is assembling", + now: "2026-06-09T00:00:01.000Z", + }); + + expect(updated).toMatchObject({ + id: created.id, + runId: "run-1", + taskId, + nodeId: "merge.node", + kind: "merge", + state: "held", + attempt: 0, + blockedReason: "shared branch is assembling", + }); + + const rows = store + .getDatabase() + .prepare("SELECT COUNT(*) AS count FROM workflow_work_items WHERE runId = ? AND taskId = ?") + .get("run-1", taskId) as { count: number }; + expect(rows.count).toBe(1); + }); + + it("lists due runnable and retrying work independently of task column", async () => { + const taskId = await createTaskId(); + await store.moveTask(taskId, "todo"); + await store.moveTask(taskId, "in-progress"); + await store.moveTask(taskId, "in-review"); + const now = "2026-06-09T00:00:00.000Z"; + + const runnable = store.upsertWorkflowWorkItem({ + runId: "run-1", + taskId, + nodeId: "plan.node", + kind: "task", + state: "runnable", + now, + }); + const futureRetry = store.upsertWorkflowWorkItem({ + runId: "run-1", + taskId, + nodeId: "retry.node", + kind: "retry", + state: "retrying", + retryAfter: "2026-06-09T00:05:00.000Z", + now, + }); + store.upsertWorkflowWorkItem({ + runId: "run-1", + taskId, + nodeId: "hold.node", + kind: "manual-hold", + state: "held", + now, + }); + + expect(store.listDueWorkflowWorkItems({ now }).map((item) => item.id)).toEqual([runnable.id]); + expect(store.listDueWorkflowWorkItems({ now: "2026-06-09T00:05:00.000Z" }).map((item) => item.id)).toEqual([ + runnable.id, + futureRetry.id, + ]); + }); + + it("acquires due leases and exposes expired running leases for reclaim", async () => { + const taskId = await createTaskId(); + const item = store.upsertWorkflowWorkItem({ + runId: "run-lease", + taskId, + nodeId: "merge.node", + kind: "merge", + state: "runnable", + now: "2026-06-09T00:00:00.000Z", + }); + + const leased = store.acquireWorkflowWorkItemLease(item.id, "worker-a", { + now: "2026-06-09T00:00:00.000Z", + leaseDurationMs: 60_000, + }); + expect(leased).toMatchObject({ + id: item.id, + state: "running", + leaseOwner: "worker-a", + leaseExpiresAt: "2026-06-09T00:01:00.000Z", + }); + + expect( + store.acquireWorkflowWorkItemLease(item.id, "worker-b", { + now: "2026-06-09T00:00:30.000Z", + leaseDurationMs: 60_000, + }), + ).toBeNull(); + expect(store.listDueWorkflowWorkItems({ now: "2026-06-09T00:00:30.000Z" })).toEqual([]); + + expect(store.listDueWorkflowWorkItems({ now: "2026-06-09T00:01:00.000Z" }).map((due) => due.id)).toEqual([item.id]); + const reclaimed = store.acquireWorkflowWorkItemLease(item.id, "worker-b", { + now: "2026-06-09T00:01:00.000Z", + leaseDurationMs: 60_000, + }); + expect(reclaimed).toMatchObject({ + id: item.id, + state: "running", + leaseOwner: "worker-b", + leaseExpiresAt: "2026-06-09T00:02:00.000Z", + }); + }); + + it("preserves lease and retry metadata on idempotent duplicate upserts", async () => { + const taskId = await createTaskId(); + const item = store.upsertWorkflowWorkItem({ + runId: "run-idempotent", + taskId, + nodeId: "retry.node", + kind: "retry", + state: "retrying", + retryAfter: "2026-06-09T00:05:00.000Z", + leaseOwner: "worker-a", + leaseExpiresAt: "2026-06-09T00:06:00.000Z", + lastError: "temporary failure", + now: "2026-06-09T00:00:00.000Z", + }); + + const duplicate = store.upsertWorkflowWorkItem({ + runId: "run-idempotent", + taskId, + nodeId: "retry.node", + kind: "retry", + now: "2026-06-09T00:01:00.000Z", + }); + + expect(duplicate).toMatchObject({ + id: item.id, + state: "retrying", + retryAfter: "2026-06-09T00:05:00.000Z", + leaseOwner: "worker-a", + leaseExpiresAt: "2026-06-09T00:06:00.000Z", + lastError: "temporary failure", + updatedAt: "2026-06-09T00:01:00.000Z", + }); + }); + + it("does not requeue terminal work", async () => { + const taskId = await createTaskId(); + const item = store.upsertWorkflowWorkItem({ + runId: "run-terminal", + taskId, + nodeId: "merge.node", + kind: "merge", + state: "runnable", + }); + + store.transitionWorkflowWorkItem(item.id, "succeeded", { now: "2026-06-09T00:00:01.000Z" }); + + expect(() => + store.upsertWorkflowWorkItem({ + runId: "run-terminal", + taskId, + nodeId: "merge.node", + kind: "merge", + state: "runnable", + }), + ).toThrow(/terminal \(succeeded\) and cannot be requeued as runnable/); + }); +}); diff --git a/packages/core/src/__tests__/task-documents.test.ts b/packages/core/src/__tests__/task-documents.test.ts index b7b71047c0..af7c770510 100644 --- a/packages/core/src/__tests__/task-documents.test.ts +++ b/packages/core/src/__tests__/task-documents.test.ts @@ -51,7 +51,7 @@ describe("TaskStore task documents", () => { expect(tableNames.has("task_documents")).toBe(true); expect(tableNames.has("task_document_revisions")).toBe(true); - expect(db.getSchemaVersion()).toBe(114); + expect(db.getSchemaVersion()).toBe(115); const index = db .prepare( diff --git a/packages/core/src/db.ts b/packages/core/src/db.ts index eb50c0ae20..25b640cabe 100644 --- a/packages/core/src/db.ts +++ b/packages/core/src/db.ts @@ -162,7 +162,7 @@ export function isFts5CorruptionError(error: unknown): boolean { // ── Schema Definition ──────────────────────────────────────────────── -const SCHEMA_VERSION = 114; +const SCHEMA_VERSION = 115; const TASKS_FTS_AUTOMERGE = 8; const TASKS_FTS_CRISISMERGE = 16; @@ -602,6 +602,27 @@ CREATE TABLE IF NOT EXISTS completion_handoff_markers ( ); CREATE INDEX IF NOT EXISTS idx_completion_handoff_markers_acceptedAt ON completion_handoff_markers(acceptedAt); +CREATE TABLE IF NOT EXISTS workflow_work_items ( + id TEXT PRIMARY KEY, + runId TEXT NOT NULL, + taskId TEXT NOT NULL REFERENCES tasks(id) ON DELETE CASCADE, + nodeId TEXT NOT NULL, + kind TEXT NOT NULL, + state TEXT NOT NULL, + attempt INTEGER NOT NULL DEFAULT 0, + retryAfter TEXT, + leaseOwner TEXT, + leaseExpiresAt TEXT, + lastError TEXT, + blockedReason TEXT, + createdAt TEXT NOT NULL, + updatedAt TEXT NOT NULL, + UNIQUE(runId, taskId, nodeId, kind) +); +CREATE INDEX IF NOT EXISTS idx_workflow_work_items_due ON workflow_work_items(state, retryAfter, createdAt); +CREATE INDEX IF NOT EXISTS idx_workflow_work_items_leaseExpiresAt ON workflow_work_items(leaseExpiresAt); +CREATE INDEX IF NOT EXISTS idx_workflow_work_items_task_run ON workflow_work_items(taskId, runId); + -- Per-branch run state for concurrent workflow fan-out/join (U13, KTD-11/R21). -- Reconstructible per ADR-0001: a crashed parallel run resumes each branch from -- its persisted node; completed branches are not re-run. Additive-only. @@ -4634,6 +4655,40 @@ export class Database { }); } + // Migration 115: Workflow-owned merge/retry/scheduling S1. + // Adds durable workflow work items so runnable, held, retrying, merge, + // manual-hold, and recovery work can be claimed generically before legacy + // merge queue and retry policy are deleted. + if (version < 115) { + this.applyMigration(115, () => { + this.db.exec(` + CREATE TABLE IF NOT EXISTS workflow_work_items ( + id TEXT PRIMARY KEY, + runId TEXT NOT NULL, + taskId TEXT NOT NULL REFERENCES tasks(id) ON DELETE CASCADE, + nodeId TEXT NOT NULL, + kind TEXT NOT NULL, + state TEXT NOT NULL, + attempt INTEGER NOT NULL DEFAULT 0, + retryAfter TEXT, + leaseOwner TEXT, + leaseExpiresAt TEXT, + lastError TEXT, + blockedReason TEXT, + createdAt TEXT NOT NULL, + updatedAt TEXT NOT NULL, + UNIQUE(runId, taskId, nodeId, kind) + ); + CREATE INDEX IF NOT EXISTS idx_workflow_work_items_due + ON workflow_work_items(state, retryAfter, createdAt); + CREATE INDEX IF NOT EXISTS idx_workflow_work_items_leaseExpiresAt + ON workflow_work_items(leaseExpiresAt); + CREATE INDEX IF NOT EXISTS idx_workflow_work_items_task_run + ON workflow_work_items(taskId, runId); + `); + }); + } + } /** diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 326972e0d8..4b7d074933 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -1,5 +1,5 @@ -export { COLUMNS, DEFAULT_COLUMN, isColumn, normalizeColumn, COLUMN_LABELS, COLUMN_DESCRIPTIONS, VALID_TRANSITIONS, DEFAULT_SETTINGS, DEFAULT_GLOBAL_SETTINGS, DEFAULT_PROJECT_SETTINGS, GLOBAL_SETTINGS_KEYS, PROJECT_SETTINGS_KEYS, isGlobalSettingsKey, isProjectSettingsKey, isMergeRequestContractShadowEnabled, resolvePersistAgentThinkingLog, THINKING_LEVELS, THEME_MODES, COLOR_THEMES, SUPPORTED_LOCALES, DEFAULT_LOCALE, isLocale, WORKFLOW_STEP_TEMPLATES, AGENT_PERMISSIONS, PERMANENT_AGENT_ACTION_CATEGORIES, AGENT_PERMISSION_POLICY_ACTION_CATEGORIES, AGENT_PROVISIONING_APPROVAL_MODES, SANDBOX_PROVISIONING_APPROVAL_MODES, AGENT_PERMISSION_POLICY_PRESET_IDS, LEGACY_AGENT_PERMISSION_POLICY_ACTION_CATEGORY_ALIASES, APPROVAL_REQUEST_STATUSES, APPROVAL_REQUEST_AUDIT_EVENT_TYPES, normalizeApprovalRequestActionCategory, isValidApprovalRequestTransition, agentToConfigSnapshot, diffConfigSnapshots, isEphemeralAgent, hasAgentIdentity, CheckoutConflictError, DEFAULT_HEARTBEAT_PROCEDURE_PATH, getDefaultHeartbeatProcedurePath, EXECUTION_MODES, DEFAULT_EXECUTION_MODE, TASK_PRIORITIES, DEFAULT_TASK_PRIORITY, HIGH_FANOUT_BLOCKER_TODO_THRESHOLD, STALE_HIGH_FANOUT_BLOCKER_AGE_THRESHOLD_MS, DASHBOARD_USER_ID, normalizeMessageParticipant, validateMessageMetadata, validateDockerNodeConfig, sanitizeDockerNodeConfigForResponse, normalizeMergeIntegrationWorktreeMode, normalizeMergeAdvanceAutoSyncMode, MERGE_ADVANCE_AUTO_SYNC_MODES, normalizeMergeConflictStrategy, normalizeMergeStrategyOverlapBehavior, normalizePostMergeAuditMode, POST_MERGE_AUDIT_MODES, normalizeMergeAuditAutoRecovery, MERGE_AUDIT_AUTO_RECOVERY_MODES, normalizeMergerMode, MERGER_MODES, normalizeAutoRecovery, AUTO_RECOVERY_MODES, buildResearchDocumentKey, REPO_OVERRIDE_RE, SHARED_STATE_SNAPSHOT_VERSION, sanitizeCliAgentSettings, sanitizeCliAgentsSettings, CLI_AGENT_ADAPTER_IDS, CLI_AGENT_AUTONOMY_MODES } from "./types.js"; -export type { Column, ColumnId, IssueInfo, IssueState, TaskSourceIssue, PrInfo, PrConflictState, PrConflictDiagnostics, PrCheckState, PrCheckStatus, PrStatus, BranchGroup, BranchGroupCreateInput, BranchGroupUpdate, BranchGroupPrState, Task, TaskTokenUsage, TaskAttachment, TaskComment, TaskCommentInput, TaskDocument, TaskDocumentRevision, TaskDocumentCreateInput, TaskDocumentWithTask, TaskCreateInput, MeshReplicatedTaskCreatePayload, MeshReplicatedTaskApplyResult, TaskSource, SourceType, TaskDetail, RetrySummary, InboxTask, TodoList, TodoItem, TodoListCreateInput, TodoListUpdateInput, TodoItemCreateInput, TodoItemUpdateInput, TodoListWithItems, AgentLogEntry, AgentLogType, AgentRole, BoardConfig, DistributedTaskIdReserveInput, DistributedTaskIdReserveResult, DistributedTaskIdCommitInput, DistributedTaskIdCommitResult, DistributedTaskIdAbortInput, DistributedTaskIdAbortResult, DistributedTaskIdStateInput, DistributedTaskIdStateResult, AutostashOrphanRecord, AutostashOutcome, MergeDetails, MergeResult, MergeIntegrationWorktreeMode, MergeAdvanceAutoSyncMode, MergeConflictStrategy, CanonicalMergeConflictStrategy, MergeStrategyOverlapBehavior, PostMergeAuditMode, MergeAuditAutoRecoveryMode, MergerMode, MergerSettings, AutoRecoveryMode, AutoRecoveryFailureClass, AutoRecoverySettings, DirectMergeCommitStrategy, Settings, GlobalSettings, ProjectSettings, SecretsEnvConfig, WebSearchBackend, ResearchEnabledSources, ResearchGlobalDefaults, ResearchProjectLimits, ResearchProjectSettings, SandboxBackendName, SandboxFailureMode, SandboxPolicy, SandboxProjectSettings, EvalFollowUpPolicy, EvalProjectSettings, ResolvedEvalSettings, SettingsScope, DaemonTokenSettings, TaskStep, StepStatus, TaskLogEntry, RunMutationContext, ActivityLogEntry, ActivityEventType, ThinkingLevel, ThemeMode, ColorTheme, Locale, ExecutionMode, TaskPriority, MergeQueueEntry, MergeQueueEnqueueOptions, MergeQueueAcquireOptions, MergeQueueReleaseOutcome, MergeRequestState, MergeRequestRecord, CompletionHandoffMarker, HandoffEvidence, HandoffToReviewOptions, UnavailableNodePolicy, OwningNodeHandoffPolicy, PlanningQuestion, PlanningSummary, PlanningResponse, PlanningQuestionType, ArchivedTaskEntry, BatchStatusRequest, BatchStatusResponse, BatchStatusEntry, BatchStatusResult, GithubIssueAction, ModelPreset, WorkflowStep, WorkflowStepMode, WorkflowStepGateMode, WorkflowStepPhase, WorkflowStepInput, WorkflowStepResult, WorkflowStepTemplate, Agent, OrgTreeNode, AgentState, AgentDetail, AgentCreateInput, AgentUpdateInput, AgentApiKey, AgentApiKeyCreateResult, AgentCapability, AgentPromptTemplate, AgentPromptsConfig, AgentPermission, PermanentAgentActionCategory, PermanentAgentSensitiveActionCategory, PermanentAgentGatingContext, AgentPermissionPolicy, AgentPermissionPolicyRules, AgentPermissionPolicyActionCategory, AgentProvisioningApprovalMode, SandboxProvisioningApprovalMode, LegacyAgentPermissionPolicyActionCategory, ApprovalRequestActionCategoryInput, ApprovalRequestActionCategory, AgentPermissionPolicyDisposition, AgentPermissionPolicyPresetId, ApprovalRequestStatus, ApprovalRequestAuditEventType, ApprovalRequestActorSnapshot, ApprovalRequestTargetAction, ApprovalRequestAuditEvent, ApprovalRequest, ApprovalRequestCreateInput, ApprovalRequestDecisionInput, ApprovalRequestCompletionInput, ApprovalRequestListInput, TaskAssignSource, AgentAccessState, AgentHeartbeatConfig, AgentBudgetConfig, AgentBudgetStatus, InstructionsBundleConfig, MessageResponseMode, AgentHeartbeatEvent, AgentHeartbeatRun, BlockedStateSnapshot, HeartbeatInvocationSource, AgentTaskSession, AgentRating, AgentRatingSummary, AgentRatingInput, AgentConfigSnapshot, RevisionFieldDiff, AgentConfigRevision, AgentStats, ReflectionTrigger, ReflectionMetrics, AgentReflection, AgentPerformanceSummary, NtfyNotificationEvent, NotificationEvent, NotificationPayload, NotificationProviderConfig, CustomProvider, SteeringComment, ParticipantType, MessageType, Message, MessageCreateInput, MessageFilter, MessageMetadata, MessageReplyReference, Mailbox, CheckoutLease, CheckoutClaimPrecondition, TaskClaimRow, CentralClaimStore, RunAuditDomain, RunAuditEvent, RunAuditEventInput, RunAuditEventFilter, AgentMemoryInclusionMode, HeartbeatPromptTemplate, HeartbeatScopeDisciplineMode, WorktrunkSettings, WorktrunkOnFailure, TaskBranchContext, CliAgentSettings } from "./types.js"; +export { COLUMNS, DEFAULT_COLUMN, isColumn, normalizeColumn, COLUMN_LABELS, COLUMN_DESCRIPTIONS, VALID_TRANSITIONS, DEFAULT_SETTINGS, DEFAULT_GLOBAL_SETTINGS, DEFAULT_PROJECT_SETTINGS, GLOBAL_SETTINGS_KEYS, PROJECT_SETTINGS_KEYS, isGlobalSettingsKey, isProjectSettingsKey, isMergeRequestContractShadowEnabled, resolvePersistAgentThinkingLog, THINKING_LEVELS, THEME_MODES, COLOR_THEMES, SUPPORTED_LOCALES, DEFAULT_LOCALE, isLocale, WORKFLOW_STEP_TEMPLATES, AGENT_PERMISSIONS, PERMANENT_AGENT_ACTION_CATEGORIES, AGENT_PERMISSION_POLICY_ACTION_CATEGORIES, AGENT_PROVISIONING_APPROVAL_MODES, SANDBOX_PROVISIONING_APPROVAL_MODES, AGENT_PERMISSION_POLICY_PRESET_IDS, LEGACY_AGENT_PERMISSION_POLICY_ACTION_CATEGORY_ALIASES, APPROVAL_REQUEST_STATUSES, APPROVAL_REQUEST_AUDIT_EVENT_TYPES, normalizeApprovalRequestActionCategory, isValidApprovalRequestTransition, agentToConfigSnapshot, diffConfigSnapshots, isEphemeralAgent, hasAgentIdentity, CheckoutConflictError, DEFAULT_HEARTBEAT_PROCEDURE_PATH, getDefaultHeartbeatProcedurePath, EXECUTION_MODES, DEFAULT_EXECUTION_MODE, TASK_PRIORITIES, DEFAULT_TASK_PRIORITY, WORKFLOW_WORK_ITEM_KINDS, WORKFLOW_WORK_ITEM_STATES, HIGH_FANOUT_BLOCKER_TODO_THRESHOLD, STALE_HIGH_FANOUT_BLOCKER_AGE_THRESHOLD_MS, DASHBOARD_USER_ID, normalizeMessageParticipant, validateMessageMetadata, validateDockerNodeConfig, sanitizeDockerNodeConfigForResponse, normalizeMergeIntegrationWorktreeMode, normalizeMergeAdvanceAutoSyncMode, MERGE_ADVANCE_AUTO_SYNC_MODES, normalizeMergeConflictStrategy, normalizeMergeStrategyOverlapBehavior, normalizePostMergeAuditMode, POST_MERGE_AUDIT_MODES, normalizeMergeAuditAutoRecovery, MERGE_AUDIT_AUTO_RECOVERY_MODES, normalizeMergerMode, MERGER_MODES, normalizeAutoRecovery, AUTO_RECOVERY_MODES, buildResearchDocumentKey, REPO_OVERRIDE_RE, SHARED_STATE_SNAPSHOT_VERSION, sanitizeCliAgentSettings, sanitizeCliAgentsSettings, CLI_AGENT_ADAPTER_IDS, CLI_AGENT_AUTONOMY_MODES } from "./types.js"; +export type { Column, ColumnId, IssueInfo, IssueState, TaskSourceIssue, PrInfo, PrConflictState, PrConflictDiagnostics, PrCheckState, PrCheckStatus, PrStatus, BranchGroup, BranchGroupCreateInput, BranchGroupUpdate, BranchGroupPrState, Task, TaskTokenUsage, TaskAttachment, TaskComment, TaskCommentInput, TaskDocument, TaskDocumentRevision, TaskDocumentCreateInput, TaskDocumentWithTask, TaskCreateInput, MeshReplicatedTaskCreatePayload, MeshReplicatedTaskApplyResult, TaskSource, SourceType, TaskDetail, RetrySummary, InboxTask, TodoList, TodoItem, TodoListCreateInput, TodoListUpdateInput, TodoItemCreateInput, TodoItemUpdateInput, TodoListWithItems, AgentLogEntry, AgentLogType, AgentRole, BoardConfig, DistributedTaskIdReserveInput, DistributedTaskIdReserveResult, DistributedTaskIdCommitInput, DistributedTaskIdCommitResult, DistributedTaskIdAbortInput, DistributedTaskIdAbortResult, DistributedTaskIdStateInput, DistributedTaskIdStateResult, AutostashOrphanRecord, AutostashOutcome, MergeDetails, MergeResult, MergeIntegrationWorktreeMode, MergeAdvanceAutoSyncMode, MergeConflictStrategy, CanonicalMergeConflictStrategy, MergeStrategyOverlapBehavior, PostMergeAuditMode, MergeAuditAutoRecoveryMode, MergerMode, MergerSettings, AutoRecoveryMode, AutoRecoveryFailureClass, AutoRecoverySettings, DirectMergeCommitStrategy, Settings, GlobalSettings, ProjectSettings, SecretsEnvConfig, WebSearchBackend, ResearchEnabledSources, ResearchGlobalDefaults, ResearchProjectLimits, ResearchProjectSettings, SandboxBackendName, SandboxFailureMode, SandboxPolicy, SandboxProjectSettings, EvalFollowUpPolicy, EvalProjectSettings, ResolvedEvalSettings, SettingsScope, DaemonTokenSettings, TaskStep, StepStatus, TaskLogEntry, RunMutationContext, ActivityLogEntry, ActivityEventType, ThinkingLevel, ThemeMode, ColorTheme, Locale, ExecutionMode, TaskPriority, MergeQueueEntry, MergeQueueEnqueueOptions, MergeQueueAcquireOptions, MergeQueueReleaseOutcome, MergeRequestState, MergeRequestRecord, CompletionHandoffMarker, WorkflowWorkItem, WorkflowWorkItemDueFilter, WorkflowWorkItemKind, WorkflowWorkItemState, WorkflowWorkItemTransitionPatch, WorkflowWorkItemUpsertInput, HandoffEvidence, HandoffToReviewOptions, UnavailableNodePolicy, OwningNodeHandoffPolicy, PlanningQuestion, PlanningSummary, PlanningResponse, PlanningQuestionType, ArchivedTaskEntry, BatchStatusRequest, BatchStatusResponse, BatchStatusEntry, BatchStatusResult, GithubIssueAction, ModelPreset, WorkflowStep, WorkflowStepMode, WorkflowStepGateMode, WorkflowStepPhase, WorkflowStepInput, WorkflowStepResult, WorkflowStepTemplate, Agent, OrgTreeNode, AgentState, AgentDetail, AgentCreateInput, AgentUpdateInput, AgentApiKey, AgentApiKeyCreateResult, AgentCapability, AgentPromptTemplate, AgentPromptsConfig, AgentPermission, PermanentAgentActionCategory, PermanentAgentSensitiveActionCategory, PermanentAgentGatingContext, AgentPermissionPolicy, AgentPermissionPolicyRules, AgentPermissionPolicyActionCategory, AgentProvisioningApprovalMode, SandboxProvisioningApprovalMode, LegacyAgentPermissionPolicyActionCategory, ApprovalRequestActionCategoryInput, ApprovalRequestActionCategory, AgentPermissionPolicyDisposition, AgentPermissionPolicyPresetId, ApprovalRequestStatus, ApprovalRequestAuditEventType, ApprovalRequestActorSnapshot, ApprovalRequestTargetAction, ApprovalRequestAuditEvent, ApprovalRequest, ApprovalRequestCreateInput, ApprovalRequestDecisionInput, ApprovalRequestCompletionInput, ApprovalRequestListInput, TaskAssignSource, AgentAccessState, AgentHeartbeatConfig, AgentBudgetConfig, AgentBudgetStatus, InstructionsBundleConfig, MessageResponseMode, AgentHeartbeatEvent, AgentHeartbeatRun, BlockedStateSnapshot, HeartbeatInvocationSource, AgentTaskSession, AgentRating, AgentRatingSummary, AgentRatingInput, AgentConfigSnapshot, RevisionFieldDiff, AgentConfigRevision, AgentStats, ReflectionTrigger, ReflectionMetrics, AgentReflection, AgentPerformanceSummary, NtfyNotificationEvent, NotificationEvent, NotificationPayload, NotificationProviderConfig, CustomProvider, SteeringComment, ParticipantType, MessageType, Message, MessageCreateInput, MessageFilter, MessageMetadata, MessageReplyReference, Mailbox, CheckoutLease, CheckoutClaimPrecondition, TaskClaimRow, CentralClaimStore, RunAuditDomain, RunAuditEvent, RunAuditEventInput, RunAuditEventFilter, AgentMemoryInclusionMode, HeartbeatPromptTemplate, HeartbeatScopeDisciplineMode, WorktrunkSettings, WorktrunkOnFailure, TaskBranchContext, CliAgentSettings } from "./types.js"; export { AGENT_VALID_TRANSITIONS, DUPLICATE_OF_METADATA_KEY } from "./types.js"; export { resolveEntryPointBranchAssignment, diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index 7dff4a89dd..a7556d52b4 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -3,7 +3,7 @@ import { randomUUID } from "node:crypto"; import { mkdir, readdir, readFile, writeFile, rename, unlink } from "node:fs/promises"; import { join } from "node:path"; import { existsSync, watch, type FSWatcher } from "node:fs"; -import type { Task, TaskDetail, TaskCreateInput, TaskAttachment, AgentLogEntry, BoardConfig, Column, ColumnId, CheckoutClaimPrecondition, MergeResult, Settings, GlobalSettings, ProjectSettings, ActivityLogEntry, ActivityEventType, TaskDocument, TaskDocumentRevision, TaskDocumentCreateInput, TaskDocumentWithTask, InboxTask, TaskLogEntry, RunMutationContext, RunAuditEvent, RunAuditEventInput, RunAuditEventFilter, ArchivedTaskEntry, ArchiveAgentLogMode, TaskPriority, SourceType, WorkflowStepTemplate, Agent, AutostashOrphanRecord, TaskCommitAssociation, TaskCommitAssociationMatchSource, TaskCommitAssociationConfidence, GithubIssueAction, MergeQueueEntry, MergeQueueEnqueueOptions, MergeQueueAcquireOptions, MergeQueueReleaseOutcome, HandoffToReviewOptions, GoalCitation, GoalCitationFilter, GoalCitationInput, GoalCitationSurface, BranchGroup, BranchGroupCreateInput, BranchGroupUpdate, TaskBranchAssignmentMode, MergeRequestRecord, MergeRequestState, CompletionHandoffMarker, PrEntity, PrEntityCreateInput, PrEntityUpdate, PrEntityState, PrThreadState, PrThreadOutcome, PrConflictState, PrChecksRollup, PrReviewDecision } from "./types.js"; +import type { Task, TaskDetail, TaskCreateInput, TaskAttachment, AgentLogEntry, BoardConfig, Column, ColumnId, CheckoutClaimPrecondition, MergeResult, Settings, GlobalSettings, ProjectSettings, ActivityLogEntry, ActivityEventType, TaskDocument, TaskDocumentRevision, TaskDocumentCreateInput, TaskDocumentWithTask, InboxTask, TaskLogEntry, RunMutationContext, RunAuditEvent, RunAuditEventInput, RunAuditEventFilter, ArchivedTaskEntry, ArchiveAgentLogMode, TaskPriority, SourceType, WorkflowStepTemplate, Agent, AutostashOrphanRecord, TaskCommitAssociation, TaskCommitAssociationMatchSource, TaskCommitAssociationConfidence, GithubIssueAction, MergeQueueEntry, MergeQueueEnqueueOptions, MergeQueueAcquireOptions, MergeQueueReleaseOutcome, HandoffToReviewOptions, GoalCitation, GoalCitationFilter, GoalCitationInput, GoalCitationSurface, BranchGroup, BranchGroupCreateInput, BranchGroupUpdate, TaskBranchAssignmentMode, MergeRequestRecord, MergeRequestState, CompletionHandoffMarker, WorkflowWorkItem, WorkflowWorkItemDueFilter, WorkflowWorkItemKind, WorkflowWorkItemState, WorkflowWorkItemTransitionPatch, WorkflowWorkItemUpsertInput, PrEntity, PrEntityCreateInput, PrEntityUpdate, PrEntityState, PrThreadState, PrThreadOutcome, PrConflictState, PrChecksRollup, PrReviewDecision } from "./types.js"; import { createActivityLogSnapshot, createRunAuditSnapshot, createTaskMetadataSnapshot, toTaskMetadataRecord, validateSnapshotEnvelope, type ActivityLogSnapshot, type RunAuditSnapshot, type TaskMetadataSnapshot } from "./shared-mesh-state.js"; import { VALID_TRANSITIONS, COLUMNS, DEFAULT_SETTINGS, isGlobalOnlySettingsKey, WORKFLOW_STEP_TEMPLATES, validateDocumentKey } from "./types.js"; import { DEFAULT_PROJECT_SETTINGS } from "./settings-schema.js"; @@ -618,6 +618,23 @@ interface CompletionHandoffMarkerRow { source: string; } +interface WorkflowWorkItemRow { + id: string; + runId: string; + taskId: string; + nodeId: string; + kind: string; + state: string; + attempt: number; + retryAfter: string | null; + leaseOwner: string | null; + leaseExpiresAt: string | null; + lastError: string | null; + blockedReason: string | null; + createdAt: string; + updatedAt: string; +} + /** Database row shape for the config table. */ interface ConfigRow { nextId: number; @@ -8793,6 +8810,59 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} }; } + private normalizeWorkflowWorkItemKind(value: string): WorkflowWorkItemKind { + switch (value) { + case "task": + case "merge": + case "retry": + case "manual-hold": + case "recovery": + return value; + default: + return "task"; + } + } + + private normalizeWorkflowWorkItemState(value: string): WorkflowWorkItemState { + switch (value) { + case "runnable": + case "running": + case "held": + case "retrying": + case "manual-required": + case "succeeded": + case "failed": + case "cancelled": + case "exhausted": + return value; + default: + return "runnable"; + } + } + + private isTerminalWorkflowWorkItemState(state: WorkflowWorkItemState): boolean { + return state === "succeeded" || state === "failed" || state === "cancelled" || state === "exhausted"; + } + + private rowToWorkflowWorkItem(row: WorkflowWorkItemRow): WorkflowWorkItem { + return { + id: row.id, + runId: row.runId, + taskId: row.taskId, + nodeId: row.nodeId, + kind: this.normalizeWorkflowWorkItemKind(row.kind), + state: this.normalizeWorkflowWorkItemState(row.state), + attempt: row.attempt, + retryAfter: row.retryAfter, + leaseOwner: row.leaseOwner, + leaseExpiresAt: row.leaseExpiresAt, + lastError: row.lastError, + blockedReason: row.blockedReason, + createdAt: row.createdAt, + updatedAt: row.updatedAt, + }; + } + private isValidMergeRequestTransition(from: MergeRequestState, to: MergeRequestState): boolean { if (from === to) return true; const allowed: Record> = { @@ -8882,6 +8952,190 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} return row ? this.rowToMergeRequestRecord(row) : null; } + upsertWorkflowWorkItem(input: WorkflowWorkItemUpsertInput): WorkflowWorkItem { + return this.db.transactionImmediate(() => { + const existing = this.db + .prepare("SELECT * FROM workflow_work_items WHERE runId = ? AND taskId = ? AND nodeId = ? AND kind = ?") + .get(input.runId, input.taskId, input.nodeId, input.kind) as WorkflowWorkItemRow | undefined; + const now = input.now ?? new Date().toISOString(); + const existingState = existing ? this.normalizeWorkflowWorkItemState(existing.state) : null; + const state = input.state ?? existingState ?? "runnable"; + if (existingState && this.isTerminalWorkflowWorkItemState(existingState) && existingState !== state) { + throw new Error( + `Workflow work item ${existing?.id ?? input.id ?? input.nodeId} is terminal (${existingState}) and cannot be requeued as ${state}`, + ); + } + + const id = existing?.id ?? input.id ?? randomUUID(); + this.db + .prepare( + `INSERT INTO workflow_work_items ( + id, runId, taskId, nodeId, kind, state, attempt, retryAfter, + leaseOwner, leaseExpiresAt, lastError, blockedReason, createdAt, updatedAt + ) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT(runId, taskId, nodeId, kind) DO UPDATE SET + state = excluded.state, + attempt = excluded.attempt, + retryAfter = excluded.retryAfter, + leaseOwner = excluded.leaseOwner, + leaseExpiresAt = excluded.leaseExpiresAt, + lastError = excluded.lastError, + blockedReason = excluded.blockedReason, + updatedAt = excluded.updatedAt`, + ) + .run( + id, + input.runId, + input.taskId, + input.nodeId, + input.kind, + state, + input.attempt ?? existing?.attempt ?? 0, + input.retryAfter === undefined ? existing?.retryAfter ?? null : input.retryAfter, + input.leaseOwner === undefined ? existing?.leaseOwner ?? null : input.leaseOwner, + input.leaseExpiresAt === undefined ? existing?.leaseExpiresAt ?? null : input.leaseExpiresAt, + input.lastError === undefined ? existing?.lastError ?? null : input.lastError, + input.blockedReason === undefined ? existing?.blockedReason ?? null : input.blockedReason, + existing?.createdAt ?? now, + now, + ); + + const row = this.db.prepare("SELECT * FROM workflow_work_items WHERE id = ?").get(id) as WorkflowWorkItemRow | undefined; + if (!row) throw new Error(`Failed to upsert workflow work item ${id}`); + this.insertRunAuditEventRow({ + taskId: row.taskId, + runId: row.runId, + domain: "database", + mutationType: "workflowWorkItem:upsert", + target: row.id, + metadata: { id: row.id, nodeId: row.nodeId, kind: row.kind, state: row.state, attempt: row.attempt }, + }); + return this.rowToWorkflowWorkItem(row); + }); + } + + transitionWorkflowWorkItem( + id: string, + state: WorkflowWorkItemState, + patch: WorkflowWorkItemTransitionPatch = {}, + ): WorkflowWorkItem { + return this.db.transactionImmediate(() => { + const now = patch.now ?? new Date().toISOString(); + const existing = this.db.prepare("SELECT * FROM workflow_work_items WHERE id = ?").get(id) as WorkflowWorkItemRow | undefined; + if (!existing) throw new Error(`Workflow work item ${id} not found`); + const fromState = this.normalizeWorkflowWorkItemState(existing.state); + if (this.isTerminalWorkflowWorkItemState(fromState) && fromState !== state) { + throw new Error(`Workflow work item ${id} is terminal (${fromState}) and cannot transition to ${state}`); + } + + this.db + .prepare( + `UPDATE workflow_work_items + SET state = ?, + attempt = ?, + retryAfter = ?, + leaseOwner = ?, + leaseExpiresAt = ?, + lastError = ?, + blockedReason = ?, + updatedAt = ? + WHERE id = ?`, + ) + .run( + state, + patch.attempt ?? existing.attempt, + patch.retryAfter === undefined ? existing.retryAfter : patch.retryAfter, + patch.leaseOwner === undefined ? existing.leaseOwner : patch.leaseOwner, + patch.leaseExpiresAt === undefined ? existing.leaseExpiresAt : patch.leaseExpiresAt, + patch.lastError === undefined ? existing.lastError : patch.lastError, + patch.blockedReason === undefined ? existing.blockedReason : patch.blockedReason, + now, + id, + ); + + const updated = this.db.prepare("SELECT * FROM workflow_work_items WHERE id = ?").get(id) as WorkflowWorkItemRow | undefined; + if (!updated) throw new Error(`Workflow work item ${id} disappeared`); + this.insertRunAuditEventRow({ + taskId: updated.taskId, + runId: updated.runId, + domain: "database", + mutationType: "workflowWorkItem:transition", + target: updated.id, + metadata: { id: updated.id, fromState, toState: state, attempt: updated.attempt }, + }); + return this.rowToWorkflowWorkItem(updated); + }); + } + + getWorkflowWorkItem(id: string): WorkflowWorkItem | null { + const row = this.db.prepare("SELECT * FROM workflow_work_items WHERE id = ?").get(id) as WorkflowWorkItemRow | undefined; + return row ? this.rowToWorkflowWorkItem(row) : null; + } + + listDueWorkflowWorkItems(filter: WorkflowWorkItemDueFilter = {}): WorkflowWorkItem[] { + const now = filter.now ?? new Date().toISOString(); + const states = filter.states?.length ? filter.states : ["runnable", "retrying"]; + const conditions = [ + `((state IN (${states.map(() => "?").join(", ")}) AND (leaseExpiresAt IS NULL OR leaseExpiresAt <= ?)) OR (state = 'running' AND leaseExpiresAt IS NOT NULL AND leaseExpiresAt <= ?))`, + "(retryAfter IS NULL OR retryAfter <= ?)", + ]; + const params: unknown[] = [...states, now, now, now]; + if (filter.kinds?.length) { + conditions.push(`kind IN (${filter.kinds.map(() => "?").join(", ")})`); + params.push(...filter.kinds); + } + params.push(filter.limit ?? 100); + + const rows = this.db + .prepare( + `SELECT * + FROM workflow_work_items + WHERE ${conditions.join(" AND ")} + ORDER BY retryAfter IS NOT NULL, retryAfter ASC, createdAt ASC + LIMIT ?`, + ) + .all(...params) as WorkflowWorkItemRow[]; + return rows.map((row) => this.rowToWorkflowWorkItem(row)); + } + + acquireWorkflowWorkItemLease( + id: string, + leaseOwner: string, + opts: { leaseDurationMs: number; now?: string }, + ): WorkflowWorkItem | null { + return this.db.transactionImmediate(() => { + const now = opts.now ?? new Date().toISOString(); + const leaseExpiresAt = new Date(new Date(now).getTime() + opts.leaseDurationMs).toISOString(); + const result = this.db + .prepare( + `UPDATE workflow_work_items + SET state = 'running', + leaseOwner = ?, + leaseExpiresAt = ?, + updatedAt = ? + WHERE id = ? + AND state IN ('runnable', 'retrying', 'running') + AND (retryAfter IS NULL OR retryAfter <= ?) + AND (leaseExpiresAt IS NULL OR leaseExpiresAt <= ?)`, + ) + .run(leaseOwner, leaseExpiresAt, now, id, now, now); + if (result.changes === 0) return null; + + const row = this.db.prepare("SELECT * FROM workflow_work_items WHERE id = ?").get(id) as WorkflowWorkItemRow | undefined; + if (!row) throw new Error(`Workflow work item ${id} disappeared`); + this.insertRunAuditEventRow({ + taskId: row.taskId, + runId: row.runId, + domain: "database", + mutationType: "workflowWorkItem:lease-acquired", + target: row.id, + metadata: { id: row.id, leaseOwner: row.leaseOwner, leaseExpiresAt: row.leaseExpiresAt }, + }); + return this.rowToWorkflowWorkItem(row); + }); + } + setCompletionHandoffAcceptedMarker( taskId: string, opts: { source: string; acceptedAt?: string }, diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index 9ee6a8c411..a26d15a01a 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -86,6 +86,80 @@ export const MERGE_REQUEST_STATES = [ export type MergeRequestState = (typeof MERGE_REQUEST_STATES)[number]; +export const WORKFLOW_WORK_ITEM_KINDS = [ + "task", + "merge", + "retry", + "manual-hold", + "recovery", +] as const; + +export type WorkflowWorkItemKind = (typeof WORKFLOW_WORK_ITEM_KINDS)[number]; + +export const WORKFLOW_WORK_ITEM_STATES = [ + "runnable", + "running", + "held", + "retrying", + "manual-required", + "succeeded", + "failed", + "cancelled", + "exhausted", +] as const; + +export type WorkflowWorkItemState = (typeof WORKFLOW_WORK_ITEM_STATES)[number]; + +export interface WorkflowWorkItem { + id: string; + runId: string; + taskId: string; + nodeId: string; + kind: WorkflowWorkItemKind; + state: WorkflowWorkItemState; + attempt: number; + retryAfter: string | null; + leaseOwner: string | null; + leaseExpiresAt: string | null; + lastError: string | null; + blockedReason: string | null; + createdAt: string; + updatedAt: string; +} + +export interface WorkflowWorkItemUpsertInput { + id?: string; + runId: string; + taskId: string; + nodeId: string; + kind: WorkflowWorkItemKind; + state?: WorkflowWorkItemState; + attempt?: number; + retryAfter?: string | null; + leaseOwner?: string | null; + leaseExpiresAt?: string | null; + lastError?: string | null; + blockedReason?: string | null; + now?: string; +} + +export interface WorkflowWorkItemTransitionPatch { + attempt?: number; + retryAfter?: string | null; + leaseOwner?: string | null; + leaseExpiresAt?: string | null; + lastError?: string | null; + blockedReason?: string | null; + now?: string; +} + +export interface WorkflowWorkItemDueFilter { + now?: string; + limit?: number; + kinds?: WorkflowWorkItemKind[]; + states?: WorkflowWorkItemState[]; +} + export interface MergeQueueEntry { taskId: string; enqueuedAt: string;