FN-5743: cut over merge dequeue authority to merge-request queue
Shift merge dequeue enforcement to merge-request queue/marker authority with reliability coverage updates. - Enforce merge eligibility and dequeue ownership via merge-request queue/marker checks in engine and scheduler paths. - Add hard-cancel behavior coverage to ensure queued merge requests are canceled when tasks are user-canceled. - Extend core merge-request record/store tests and reliability interaction suites for dual-observe and cancel-on-hard-cancel seams. - Document the FN-5741/FN-5743 reliability backstop updates in AGENTS.md and architecture docs. Files changed: AGENTS.md | 2 + docs/architecture.md | 2 + packages/core/src/__tests__/merge-request-record.test.ts | 18 ++++ packages/core/src/store.ts | 13 +++ packages/engine/src/__tests__/reliability-interactions/dual-observe-merge-seam.test.ts | 57 +++++++++++++ packages/engine/src/__tests__/reliability-interactions/merge-request-cancel-on-hard-cancel.test.ts | 87 +++++++++++++++++++ packages/engine/src/project-engine.ts | 98 +++++++++++++++++++++- packages/engine/src/scheduler.ts | 12 ++- 8 files changed, 284 insertions(+), 5 deletions(-) Fusion-Task-Id: FN-5743 Fusion-Task-Lineage: b9a49aeb-ed73-42fe-924d-25c7097d1bb9
This commit is contained in:
@@ -190,4 +190,61 @@ describe("FN-5742 dual-observe merge seam", () => {
|
||||
|
||||
expect(updateTask).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("transient retry transitions merge-request running->retrying->queued under contract", async () => {
|
||||
const transitions: string[] = [];
|
||||
let state = "running";
|
||||
const store = {
|
||||
getSettings: vi.fn().mockResolvedValue({ mergeRequestContractShadowEnabled: true }),
|
||||
getMergeRequestRecord: vi.fn(() => ({ state, attemptCount: 0, lastError: null })),
|
||||
transitionMergeRequestState: vi.fn((_taskId: string, to: string) => {
|
||||
transitions.push(`${state}->${to}`);
|
||||
state = to;
|
||||
}),
|
||||
updateTask: vi.fn().mockResolvedValue(undefined),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
} as any;
|
||||
|
||||
const retried = await (ProjectEngine.prototype as any).maybeRetryTransientMerge.call(
|
||||
{ shuttingDown: false, internalEnqueueMerge: vi.fn() },
|
||||
store,
|
||||
"FN-MR",
|
||||
{ mergeTransientRetryCount: 0 },
|
||||
"lease-handoff-failed: target-not-queued",
|
||||
);
|
||||
|
||||
expect(retried).toBe(true);
|
||||
expect(transitions).toEqual(["running->retrying", "retrying->queued"]);
|
||||
});
|
||||
|
||||
it("transient exhaustion marks merge request exhausted under contract", async () => {
|
||||
const logs: string[] = [];
|
||||
let state = "running";
|
||||
const store = {
|
||||
getSettings: vi.fn().mockResolvedValue({ mergeRequestContractShadowEnabled: true }),
|
||||
getTask: vi.fn().mockResolvedValue({ id: "FN-MR", column: "in-review" }),
|
||||
getMergeRequestRecord: vi.fn(() => ({ state, attemptCount: 3, lastError: null })),
|
||||
transitionMergeRequestState: vi.fn((_taskId: string, to: string) => {
|
||||
state = to;
|
||||
}),
|
||||
logEntry: vi.fn(async (_taskId: string, message: string) => logs.push(message)),
|
||||
updateTask: vi.fn(),
|
||||
getActiveMergingTask: vi.fn().mockReturnValue(null),
|
||||
} as any;
|
||||
|
||||
if ((ProjectEngine.prototype as any).isTransientMergeRetryExhausted.call({}, { mergeTransientRetryCount: 3 }, "socket hang up")) {
|
||||
const record = store.getMergeRequestRecord("FN-MR");
|
||||
if (record.state === "running") {
|
||||
store.transitionMergeRequestState("FN-MR", "retrying", { attemptCount: record.attemptCount, lastError: "socket hang up" });
|
||||
}
|
||||
const refreshed = store.getMergeRequestRecord("FN-MR");
|
||||
if (refreshed.state === "retrying") {
|
||||
store.transitionMergeRequestState("FN-MR", "exhausted", { attemptCount: refreshed.attemptCount, lastError: "socket hang up" });
|
||||
}
|
||||
await store.logEntry("FN-MR", "marked merge request exhausted without column rebound: socket hang up");
|
||||
}
|
||||
|
||||
expect(state).toBe("exhausted");
|
||||
expect(logs.at(-1)).toContain("without column rebound");
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,87 @@
|
||||
import { mkdtempSync } from "node:fs";
|
||||
import { rm } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { TaskStore } from "@fusion/core";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { ProjectEngine } from "../../project-engine.js";
|
||||
|
||||
describe("FN-5743 hard-cancel merge-request cutover", () => {
|
||||
let rootDir: string;
|
||||
let globalDir: string;
|
||||
let store: TaskStore;
|
||||
|
||||
beforeEach(async () => {
|
||||
rootDir = mkdtempSync(join(tmpdir(), "kb-fn-5743-hard-cancel-"));
|
||||
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 });
|
||||
});
|
||||
|
||||
it("cancels pending merge request on user in-review->todo hard-cancel", async () => {
|
||||
const task = await store.createTask({ description: "FN-5743 hard-cancel" });
|
||||
await store.moveTask(task.id, "todo");
|
||||
await store.moveTask(task.id, "in-progress");
|
||||
await store.handoffToReview(task.id, {
|
||||
ownerAgentId: "agent",
|
||||
evidence: { reason: "fn_task_done", runId: "run-1", agentId: "agent" },
|
||||
});
|
||||
|
||||
store.upsertMergeRequestRecord(task.id, { state: "queued", attemptCount: 1 });
|
||||
store.setCompletionHandoffAcceptedMarker(task.id, { source: "executor:fn_task_done" });
|
||||
|
||||
await store.moveTask(task.id, "todo", { moveSource: "user" });
|
||||
|
||||
expect(store.getMergeRequestRecord(task.id)?.state).toBe("cancelled");
|
||||
expect(store.getCompletionHandoffAcceptedMarker(task.id)).toBeNull();
|
||||
});
|
||||
|
||||
it("does not cancel merge request on engine in-review->todo rebound", async () => {
|
||||
const task = await store.createTask({ description: "FN-5743 engine rebound" });
|
||||
await store.moveTask(task.id, "todo");
|
||||
await store.moveTask(task.id, "in-progress");
|
||||
await store.handoffToReview(task.id, {
|
||||
ownerAgentId: "agent",
|
||||
evidence: { reason: "fn_task_done", runId: "run-2", agentId: "agent" },
|
||||
});
|
||||
|
||||
store.upsertMergeRequestRecord(task.id, { state: "queued", attemptCount: 1 });
|
||||
store.setCompletionHandoffAcceptedMarker(task.id, { source: "executor:fn_task_done" });
|
||||
|
||||
await store.moveTask(task.id, "todo", { moveSource: "engine" as any });
|
||||
|
||||
expect(store.getMergeRequestRecord(task.id)?.state).toBe("queued");
|
||||
expect(store.getCompletionHandoffAcceptedMarker(task.id)).not.toBeNull();
|
||||
});
|
||||
|
||||
it("transient merge retry uses merge-request state transitions without todo rebound", async () => {
|
||||
let state = "running";
|
||||
const fakeStore = {
|
||||
getSettings: vi.fn().mockResolvedValue({ mergeRequestContractShadowEnabled: true }),
|
||||
getMergeRequestRecord: vi.fn(() => ({ state, attemptCount: 0, lastError: null })),
|
||||
transitionMergeRequestState: vi.fn((_taskId: string, toState: string) => {
|
||||
state = toState;
|
||||
}),
|
||||
updateTask: vi.fn().mockResolvedValue(undefined),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
moveTask: vi.fn(),
|
||||
} as any;
|
||||
|
||||
const retried = await (ProjectEngine.prototype as any).maybeRetryTransientMerge.call(
|
||||
{ shuttingDown: false, internalEnqueueMerge: vi.fn() },
|
||||
fakeStore,
|
||||
"FN-5743",
|
||||
{ id: "FN-5743", mergeTransientRetryCount: 0 },
|
||||
"lease-handoff-failed: target-not-queued",
|
||||
);
|
||||
|
||||
expect(retried).toBe(true);
|
||||
expect(state).toBe("queued");
|
||||
expect(fakeStore.moveTask).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user