From fd57b5b7a1a781caee3a34215e394a9e87d76de3 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Thu, 11 Jun 2026 08:09:13 -0700 Subject: [PATCH] fix(FN-000): prevent duplicate handoff merge work Address PR #1580 feedback by preserving active same-key handoff work, cancelling stale merge/manual-hold work across re-handoffs, and avoiding stable fallback run id collisions after terminal work. --- .../__tests__/merge-request-record.test.ts | 91 +++++++++++++++++++ packages/core/src/store.ts | 63 +++++++++++-- 2 files changed, 146 insertions(+), 8 deletions(-) diff --git a/packages/core/src/__tests__/merge-request-record.test.ts b/packages/core/src/__tests__/merge-request-record.test.ts index 2050f62565..9abdc9bd41 100644 --- a/packages/core/src/__tests__/merge-request-record.test.ts +++ b/packages/core/src/__tests__/merge-request-record.test.ts @@ -250,6 +250,97 @@ describe("TaskStore merge request record + completion handoff marker", () => { ]); }); + it("cancels previous active handoff work when a re-handoff uses a new run id", async () => { + const taskId = await createTask(); + await store.moveTask(taskId, "todo"); + await store.moveTask(taskId, "in-progress"); + + await store.handoffToReview(taskId, { + ownerAgentId: "agent-test", + evidence: { reason: "fn_task_done", runId: "run-handoff-1", agentId: "agent-test" }, + now: "2026-05-30T00:00:00.000Z", + }); + await store.handoffToReview(taskId, { + ownerAgentId: "agent-test", + evidence: { reason: "fn_task_done", runId: "run-handoff-2", agentId: "agent-test" }, + now: "2026-05-30T00:00:01.000Z", + }); + + expect(store.listWorkflowWorkItemsForTask(taskId, { kinds: ["merge"] })).toEqual([ + expect.objectContaining({ + runId: "run-handoff-1", + state: "cancelled", + lastError: "superseded-by-completion-handoff", + }), + expect.objectContaining({ + runId: "run-handoff-2", + state: "runnable", + }), + ]); + }); + + it("cancels opposite handoff kind when autoMerge flips between handoffs", async () => { + const taskId = await createTask(); + await store.moveTask(taskId, "todo"); + await store.moveTask(taskId, "in-progress"); + + await store.handoffToReview(taskId, { + ownerAgentId: "agent-test", + evidence: { reason: "fn_task_done", runId: "run-merge", agentId: "agent-test" }, + now: "2026-05-30T00:00:00.000Z", + }); + await store.updateTask(taskId, { autoMerge: false }); + await store.handoffToReview(taskId, { + ownerAgentId: "agent-test", + evidence: { reason: "fn_task_done", runId: "run-manual", agentId: "agent-test" }, + now: "2026-05-30T00:00:01.000Z", + }); + + expect(store.listWorkflowWorkItemsForTask(taskId)).toEqual([ + expect.objectContaining({ + runId: "run-merge", + kind: "merge", + state: "cancelled", + lastError: "superseded-by-completion-handoff", + }), + expect.objectContaining({ + runId: "run-manual", + kind: "manual-hold", + state: "manual-required", + }), + ]); + }); + + it("does not reset running handoff work to runnable on same-run replay", async () => { + const taskId = await createTask(); + await store.moveTask(taskId, "todo"); + await store.moveTask(taskId, "in-progress"); + + await store.handoffToReview(taskId, { + ownerAgentId: "agent-test", + evidence: { reason: "fn_task_done", runId: "run-handoff", agentId: "agent-test" }, + now: "2026-05-30T00:00:00.000Z", + }); + const [mergeWork] = store.listWorkflowWorkItemsForTask(taskId, { kinds: ["merge"] }); + store.transitionWorkflowWorkItem(mergeWork.id, "running", { + leaseOwner: "worker-a", + leaseExpiresAt: "2026-05-30T00:05:00.000Z", + now: "2026-05-30T00:00:01.000Z", + }); + + await store.handoffToReview(taskId, { + ownerAgentId: "agent-test", + evidence: { reason: "fn_task_done", runId: "run-handoff", agentId: "agent-test" }, + now: "2026-05-30T00:00:02.000Z", + }); + + expect(store.getWorkflowWorkItem(mergeWork.id)).toMatchObject({ + state: "running", + leaseOwner: "worker-a", + leaseExpiresAt: "2026-05-30T00:05:00.000Z", + }); + }); + it("creates manual hold workflow work instead of merge work when autoMerge is false", async () => { const taskId = await createTask(); await store.updateTask(taskId, { autoMerge: false }); diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index b6376bc7b6..537ccf74e9 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -8859,6 +8859,10 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} return state === "succeeded" || state === "failed" || state === "cancelled" || state === "exhausted"; } + private isActiveWorkflowWorkItemState(state: WorkflowWorkItemState): boolean { + return state === "runnable" || state === "running" || state === "held" || state === "retrying" || state === "manual-required"; + } + private workflowStateForMergeRequestState(state: MergeRequestState): WorkflowWorkItemState { const states: Record = { queued: "runnable", @@ -9022,15 +9026,57 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} opts: { runId?: string; now?: string; source?: string } = {}, ): WorkflowWorkItem { const autoMerge = task.autoMerge !== false; + const runId = opts.runId ?? `completion-handoff:${task.id}:${randomUUID()}`; + const nodeId = autoMerge ? "merge-gate" : "merge-manual-hold"; + const kind: WorkflowWorkItemKind = autoMerge ? "merge" : "manual-hold"; + const existing = this.getWorkflowWorkItemByIdentity(runId, task.id, nodeId, kind); + if (existing && this.isActiveWorkflowWorkItemState(existing.state)) { + this.cancelActiveWorkflowWorkItemsForTask(task.id, { + kinds: ["merge", "manual-hold"], + excludeIds: [existing.id], + now: opts.now, + lastError: "superseded-by-completion-handoff", + }); + this.insertCompletionHandoffWorkflowWorkAudit(task, existing, autoMerge, opts.source); + return existing; + } + + this.cancelActiveWorkflowWorkItemsForTask(task.id, { + kinds: ["merge", "manual-hold"], + now: opts.now, + lastError: "superseded-by-completion-handoff", + }); const item = this.upsertWorkflowWorkItem({ - runId: opts.runId ?? `completion-handoff:${task.id}`, + runId, taskId: task.id, - nodeId: autoMerge ? "merge-gate" : "merge-manual-hold", - kind: autoMerge ? "merge" : "manual-hold", + nodeId, + kind, state: autoMerge ? "runnable" : "manual-required", blockedReason: autoMerge ? null : "autoMerge:false", now: opts.now, }); + this.insertCompletionHandoffWorkflowWorkAudit(task, item, autoMerge, opts.source); + return item; + } + + private getWorkflowWorkItemByIdentity( + runId: string, + taskId: string, + nodeId: string, + kind: WorkflowWorkItemKind, + ): WorkflowWorkItem | null { + const row = this.db + .prepare("SELECT * FROM workflow_work_items WHERE runId = ? AND taskId = ? AND nodeId = ? AND kind = ?") + .get(runId, taskId, nodeId, kind) as WorkflowWorkItemRow | undefined; + return row ? this.rowToWorkflowWorkItem(row) : null; + } + + private insertCompletionHandoffWorkflowWorkAudit( + task: Pick, + item: WorkflowWorkItem, + autoMerge: boolean, + source?: string, + ): void { this.insertRunAuditEventRow({ taskId: task.id, runId: item.runId, @@ -9040,13 +9086,12 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} metadata: { taskId: task.id, autoMerge, - source: opts.source ?? "completion-handoff", + source: source ?? "completion-handoff", workItemId: item.id, nodeId: item.nodeId, state: item.state, }, }); - return item; } upsertWorkflowWorkItem(input: WorkflowWorkItemUpsertInput): WorkflowWorkItem { @@ -9190,11 +9235,13 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} cancelActiveWorkflowWorkItemsForTask( taskId: string, - opts: { kinds?: WorkflowWorkItemKind[]; now?: string; lastError?: string | null } = {}, + opts: { kinds?: WorkflowWorkItemKind[]; now?: string; lastError?: string | null; excludeIds?: string[] } = {}, ): WorkflowWorkItem[] { return this.db.transactionImmediate(() => { - const activeStates: WorkflowWorkItemState[] = ["runnable", "running", "held", "retrying", "manual-required"]; - const items = this.listWorkflowWorkItemsForTask(taskId, opts).filter((item) => activeStates.includes(item.state)); + const excludeIds = new Set(opts.excludeIds ?? []); + const items = this.listWorkflowWorkItemsForTask(taskId, opts).filter((item) => + this.isActiveWorkflowWorkItemState(item.state) && !excludeIds.has(item.id) + ); return items.map((item) => this.transitionWorkflowWorkItem(item.id, "cancelled", { now: opts.now,