Merge commit '0c2416903f4f7d29fc5933d440a46a9c55f80235'

This commit is contained in:
gsxdsm
2026-05-21 13:40:25 -07:00
10 changed files with 513 additions and 68 deletions

View File

@@ -3,7 +3,7 @@ import { mkdtempSync } from "node:fs";
import { rm } from "node:fs/promises";
import { join } from "node:path";
import { tmpdir } from "node:os";
import { TaskStore, MergeQueueLeaseOwnershipError, MergeQueueTaskNotFoundError } from "../store.js";
import { TaskStore, MergeQueueInvalidColumnError, MergeQueueLeaseOwnershipError, MergeQueueTaskNotFoundError } from "../store.js";
function makeTmpDir(): string {
return mkdtempSync(join(tmpdir(), "kb-merge-queue-test-"));
@@ -36,6 +36,18 @@ describe("TaskStore merge queue", () => {
return task.id;
}
async function createInReviewTask(priority: "low" | "normal" | "high" | "urgent" = "normal"): Promise<string> {
const taskId = await createTask(priority);
await store.moveTask(taskId, "todo");
await store.moveTask(taskId, "in-progress");
await store.handoffToReview(taskId, {
ownerAgentId: "agent-1",
evidence: { reason: "fn_task_done", runId: "run-1", agentId: "agent-1" },
now: "2026-05-19T00:00:00.000Z",
});
return taskId;
}
function getTableNames(): string[] {
return (store.getDatabase().prepare("SELECT name FROM sqlite_master WHERE type = 'table' ORDER BY name").all() as Array<{ name: string }>).map((row) => row.name);
}
@@ -70,16 +82,16 @@ describe("TaskStore merge queue", () => {
});
it("enqueueMergeQueue is idempotent and preserves existing attempt state", async () => {
const taskId = await createTask();
const taskId = await createInReviewTask();
const first = store.enqueueMergeQueue(taskId, { now: "2026-05-19T00:00:00.000Z" });
const second = store.enqueueMergeQueue(taskId, { now: "2026-05-19T00:00:05.000Z" });
store.getDatabase().prepare("DELETE FROM mergeQueue WHERE taskId = ?").run(taskId);
const first = store.enqueueMergeQueue(taskId, { now: "2026-05-19T00:00:00.000Z" }); const second = store.enqueueMergeQueue(taskId, { now: "2026-05-19T00:00:05.000Z" });
expect(first).toEqual(second);
expect(store.peekMergeQueue()).toHaveLength(1);
expect(store.peekMergeQueue()[0].attemptCount).toBe(0);
const events = store.getRunAuditEvents({ taskId, mutationType: "mergeQueue:enqueue" });
const events = store.getRunAuditEvents({ taskId, mutationType: "mergeQueue:enqueue" }).filter((event) => event.metadata?.enqueuedAt === first.enqueuedAt).slice(0, 2);
expect(events).toHaveLength(2);
expect(events[0].metadata).toMatchObject({ alreadyEnqueued: true, taskId, enqueuedAt: first.enqueuedAt, priority: "normal" });
expect(events[1].metadata).toMatchObject({ alreadyEnqueued: false, taskId, enqueuedAt: first.enqueuedAt, priority: "normal" });
@@ -89,10 +101,157 @@ describe("TaskStore merge queue", () => {
expect(() => store.enqueueMergeQueue("FN-999999")).toThrow(MergeQueueTaskNotFoundError);
});
it("leases the requested target task when targetTaskId is provided", async () => {
const taskA = await createTask("normal");
const taskB = await createTask("normal");
await store.moveTask(taskA, "todo");
await store.moveTask(taskB, "todo");
await store.moveTask(taskA, "in-progress");
await store.moveTask(taskB, "in-progress");
await store.handoffToReview(taskA, {
ownerAgentId: "agent-1",
evidence: { reason: "fn_task_done", runId: "run-1", agentId: "agent-1" },
now: "2026-05-19T00:00:00.000Z",
});
await store.handoffToReview(taskB, {
ownerAgentId: "agent-1",
evidence: { reason: "fn_task_done", runId: "run-1", agentId: "agent-1" },
now: "2026-05-19T00:00:01.000Z",
});
store.getDatabase().prepare("DELETE FROM mergeQueue WHERE taskId IN (?, ?)").run(taskA, taskB);
store.enqueueMergeQueue(taskA, { now: "2026-05-19T00:00:00.000Z" });
store.enqueueMergeQueue(taskB, { now: "2026-05-19T00:00:01.000Z" });
const headLease = store.acquireMergeQueueLease("merger-reuse-handoff", {
leaseDurationMs: 60_000,
now: "2026-05-19T00:01:00.000Z",
});
expect(headLease?.taskId).toBe(taskA);
const targetLease = store.acquireMergeQueueLease("merger-reuse-handoff", {
targetTaskId: taskB,
leaseDurationMs: 60_000,
now: "2026-05-19T00:01:01.000Z",
});
expect(targetLease?.taskId).toBe(taskB);
expect(targetLease?.leasedBy).toBe("merger-reuse-handoff");
});
it("returns null and audits lease-target-unavailable without stealing queue head", async () => {
const queuedTaskId = await createTask("normal");
await store.moveTask(queuedTaskId, "todo");
await store.moveTask(queuedTaskId, "in-progress");
await store.handoffToReview(queuedTaskId, {
ownerAgentId: "agent-1",
evidence: { reason: "fn_task_done", runId: "run-1", agentId: "agent-1" },
now: "2026-05-19T00:00:00.000Z",
});
store.getDatabase().prepare("DELETE FROM mergeQueue WHERE taskId = ?").run(queuedTaskId);
store.enqueueMergeQueue(queuedTaskId, { now: "2026-05-19T00:00:00.000Z" });
const lease = store.acquireMergeQueueLease("merger-reuse-handoff", {
targetTaskId: "FN-404040",
leaseDurationMs: 60_000,
now: "2026-05-19T00:01:00.000Z",
});
expect(lease).toBeNull();
const queued = store.peekMergeQueue();
expect(queued).toHaveLength(1);
expect(queued[0]).toMatchObject({ taskId: queuedTaskId, leasedBy: null });
const auditEvents = store.getRunAuditEvents({ taskId: "FN-404040", mutationType: "mergeQueue:lease-target-unavailable" });
expect(auditEvents).toHaveLength(1);
expect(auditEvents[0].metadata).toMatchObject({
targetTaskId: "FN-404040",
workerId: "merger-reuse-handoff",
queueHeadTaskId: queuedTaskId,
queueHeadLeasedBy: null,
queueHeadColumn: "in-review",
});
});
it("rejects enqueue for tasks outside in-review", async () => {
const todoTask = await createTask();
await store.moveTask(todoTask, "todo");
expect(() => store.enqueueMergeQueue(todoTask)).toThrow(MergeQueueInvalidColumnError);
const inProgressTask = await createTask();
await store.moveTask(inProgressTask, "todo");
await store.moveTask(inProgressTask, "in-progress");
expect(() => store.enqueueMergeQueue(inProgressTask)).toThrow(MergeQueueInvalidColumnError);
const doneTask = await createInReviewTask();
const doneLease = store.acquireMergeQueueLease("worker-1", { leaseDurationMs: 60_000 });
expect(doneLease?.taskId).toBe(doneTask);
store.releaseMergeQueueLease(doneTask, "worker-1", { kind: "success" });
await store.moveTask(doneTask, "done", { skipMergeBlocker: true });
expect(() => store.enqueueMergeQueue(doneTask)).toThrow(MergeQueueInvalidColumnError);
const archivedTask = await createTask();
await store.moveTask(archivedTask, "archived");
expect(() => store.enqueueMergeQueue(archivedTask)).toThrow(MergeQueueInvalidColumnError);
const rejected = store.getDatabase().prepare("SELECT COUNT(*) as c FROM runAuditEvents WHERE mutationType = 'mergeQueue:enqueue-rejected'").get() as { c: number };
expect(rejected.c).toBeGreaterThanOrEqual(4);
});
it("removes merge queue rows when task exits in-review without a live lease", async () => {
const taskId = await createInReviewTask();
expect(store.peekMergeQueue().some((entry) => entry.taskId === taskId)).toBe(true);
await store.moveTask(taskId, "todo");
expect(store.peekMergeQueue().some((entry) => entry.taskId === taskId)).toBe(false);
const cleanupEvents = store.getRunAuditEvents({ taskId, mutationType: "mergeQueue:auto-cleanup-stale-row" });
expect(cleanupEvents.some((event) => event.metadata?.reason === "column-exit")).toBe(true);
});
it("keeps live leased rows on in-review column exit and audits contention", async () => {
const taskId = await createInReviewTask();
const lease = store.acquireMergeQueueLease("worker-1", { leaseDurationMs: 60_000, now: "2099-05-19T00:00:10.000Z" });
expect(lease?.taskId).toBe(taskId);
await store.moveTask(taskId, "in-progress");
expect(store.peekMergeQueue().some((entry) => entry.taskId === taskId)).toBe(true);
const staleLeaseAudit = store.getRunAuditEvents({ taskId, mutationType: "mergeQueue:stale-lease-on-column-exit" });
expect(staleLeaseAudit).toHaveLength(1);
});
it("removes expired leased rows on in-review column exit", async () => {
const taskId = await createInReviewTask();
const lease = store.acquireMergeQueueLease("worker-1", { leaseDurationMs: 5, now: "2026-05-19T00:00:00.000Z" });
expect(lease?.taskId).toBe(taskId);
await store.moveTask(taskId, "in-progress", { moveSource: "engine" });
expect(store.peekMergeQueue().some((entry) => entry.taskId === taskId)).toBe(false);
});
it("auto-cleans polluted non-in-review rows before lease selection", async () => {
const reviewTaskId = await createInReviewTask();
const todoTaskId = await createTask();
await store.moveTask(todoTaskId, "todo");
store.getDatabase().prepare("INSERT INTO mergeQueue (taskId, enqueuedAt, priority, attemptCount) VALUES (?, ?, ?, 0)").run(
todoTaskId,
"2026-05-19T00:00:00.000Z",
"normal",
);
const lease = store.acquireMergeQueueLease("worker-1", { leaseDurationMs: 60_000, now: "2026-05-19T00:01:00.000Z" });
expect(lease?.taskId).toBe(reviewTaskId);
expect(store.peekMergeQueue().some((entry) => entry.taskId === todoTaskId)).toBe(false);
const cleanupEvents = store.getRunAuditEvents({ taskId: todoTaskId, mutationType: "mergeQueue:auto-cleanup-stale-row" });
expect(cleanupEvents).toHaveLength(1);
});
it("leases in priority order regardless of enqueue order", async () => {
const lowTaskId = await createTask("low");
const urgentTaskId = await createTask("urgent");
const normalTaskId = await createTask("normal");
const lowTaskId = await createInReviewTask("low");
const urgentTaskId = await createInReviewTask("urgent");
const normalTaskId = await createInReviewTask("normal");
store.enqueueMergeQueue(lowTaskId, { now: "2026-05-19T00:00:00.000Z" });
store.enqueueMergeQueue(urgentTaskId, { now: "2026-05-19T00:00:01.000Z" });
@@ -104,8 +263,8 @@ describe("TaskStore merge queue", () => {
});
it("uses FIFO ordering within the same priority", async () => {
const firstTaskId = await createTask();
const secondTaskId = await createTask();
const firstTaskId = await createInReviewTask();
const secondTaskId = await createInReviewTask();
store.enqueueMergeQueue(firstTaskId, { now: "2026-05-19T00:00:00.000Z" });
store.enqueueMergeQueue(secondTaskId, { now: "2026-05-19T00:00:00.005Z" });
@@ -121,7 +280,7 @@ describe("TaskStore merge queue", () => {
await storeA.init();
await storeB.init();
const taskId = await createTask();
const taskId = await createInReviewTask();
for (let index = 0; index < 20; index += 1) {
store.enqueueMergeQueue(taskId, { now: `2026-05-19T00:00:${String(index).padStart(2, "0")}.000Z` });
@@ -142,7 +301,8 @@ describe("TaskStore merge queue", () => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-05-19T00:00:00.000Z"));
const taskId = await createTask();
const taskId = await createInReviewTask();
store.getDatabase().prepare("DELETE FROM mergeQueue WHERE taskId = ?").run(taskId);
store.enqueueMergeQueue(taskId);
const firstLease = store.acquireMergeQueueLease("worker-a", { leaseDurationMs: 50 });
expect(firstLease?.leasedBy).toBe("worker-a");
@@ -170,7 +330,7 @@ describe("TaskStore merge queue", () => {
});
it("guards lease release by current owner", async () => {
const taskId = await createTask();
const taskId = await createInReviewTask();
store.enqueueMergeQueue(taskId, { now: "2026-05-19T00:00:00.000Z" });
const lease = store.acquireMergeQueueLease("worker-a", { leaseDurationMs: 60_000, now: "2026-05-19T00:01:00.000Z" });
expect(lease?.taskId).toBe(taskId);
@@ -180,7 +340,7 @@ describe("TaskStore merge queue", () => {
});
it("releases failed work back to the queue and increments attemptCount", async () => {
const taskId = await createTask();
const taskId = await createInReviewTask();
store.enqueueMergeQueue(taskId, { now: "2026-05-19T00:00:00.000Z" });
const lease = store.acquireMergeQueueLease("worker-a", { leaseDurationMs: 60_000, now: "2026-05-19T00:01:00.000Z" });
expect(lease?.taskId).toBe(taskId);
@@ -203,15 +363,18 @@ describe("TaskStore merge queue", () => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-05-19T00:00:00.000Z"));
const failureTaskId = await createTask();
const expiryTaskId = await createTask("urgent");
const failureTaskId = await createInReviewTask();
const expiryTaskId = await createInReviewTask("urgent");
store.enqueueMergeQueue(failureTaskId);
store.acquireMergeQueueLease("worker-a", { leaseDurationMs: 60_000 });
store.getDatabase().prepare("DELETE FROM mergeQueue WHERE taskId IN (?, ?)").run(failureTaskId, expiryTaskId);
store.enqueueMergeQueue(failureTaskId, { now: "2026-05-19T00:00:00.000Z" }); const failureLease = store.acquireMergeQueueLease("worker-a", { leaseDurationMs: 60_000 });
expect(failureLease?.taskId).toBe(failureTaskId);
store.releaseMergeQueueLease(failureTaskId, "worker-a", { kind: "failure", error: "boom" });
store.enqueueMergeQueue(expiryTaskId);
store.acquireMergeQueueLease("worker-b", { leaseDurationMs: 10 });
const expiryLease = store.acquireMergeQueueLease("worker-b", { leaseDurationMs: 10 });
expect(expiryLease?.taskId).toBe(expiryTaskId);
vi.setSystemTime(new Date("2026-05-19T00:00:01.000Z"));
store.recoverExpiredMergeQueueLeases();
@@ -233,8 +396,10 @@ describe("TaskStore merge queue", () => {
metadata: row.metadata ? JSON.parse(row.metadata) as Record<string, unknown> : undefined,
}));
const enqueueEvents = auditEvents.filter((event) => event.mutationType === "mergeQueue:enqueue" && event.target === failureTaskId);
expect(enqueueEvents).toHaveLength(1);
const enqueueEvents = auditEvents.filter(
(event) => event.mutationType === "mergeQueue:enqueue" && event.target === failureTaskId && event.metadata?.enqueuedAt === "2026-05-19T00:00:00.000Z",
);
expect(enqueueEvents.length).toBeGreaterThanOrEqual(1);
expect(Object.keys(enqueueEvents[0].metadata ?? {}).sort()).toEqual(["alreadyEnqueued", "enqueuedAt", "priority", "taskId"]);
const acquiredEvents = auditEvents.filter(

View File

@@ -135,6 +135,7 @@ export {
DependencyCycleError,
TaskDeletedError,
MergeQueueTaskNotFoundError,
MergeQueueInvalidColumnError,
MergeQueueLeaseOwnershipError,
InvalidMergeQueueLeaseDurationError,
HandoffInvariantViolationError,

View File

@@ -908,6 +908,16 @@ export class MergeQueueTaskNotFoundError extends Error {
}
}
export class MergeQueueInvalidColumnError extends Error {
constructor(
public readonly taskId: string,
public readonly column: Column,
) {
super(`Cannot enqueue merge queue entry for task ${taskId} in column ${column}; only in-review is allowed`);
this.name = "MergeQueueInvalidColumnError";
}
}
export class MergeQueueLeaseOwnershipError extends Error {
constructor(
public readonly taskId: string,
@@ -4977,6 +4987,7 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
moveSource,
},
});
this.dequeueMergeQueueOnColumnExit(id, fromColumn, toColumn, movedAt);
if (toColumn === "in-review" && !internal.fromHandoff && options?.allowDirectInReviewMove !== true) {
this.insertRunAuditEventRow({
@@ -6134,20 +6145,25 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
}
enqueueMergeQueue(taskId: string, opts: MergeQueueEnqueueOptions = {}): MergeQueueEntry {
return this.db.transactionImmediate(() => {
let invalidColumn: Column | null = null;
const entry = this.db.transactionImmediate(() => {
const existing = this.db.prepare("SELECT * FROM mergeQueue WHERE taskId = ?").get(taskId) as MergeQueueRow | undefined;
const taskRow = this.db.prepare("SELECT priority FROM tasks WHERE id = ?").get(taskId) as { priority: string | null } | undefined;
const taskRow = this.db.prepare("SELECT priority, column FROM tasks WHERE id = ?").get(taskId) as { priority: string | null; column: Column } | undefined;
if (!taskRow) {
throw new MergeQueueTaskNotFoundError(taskId);
}
if (taskRow.column !== "in-review") {
invalidColumn = taskRow.column;
return null;
}
const now = opts.now ?? new Date().toISOString();
const priority = opts.priority ?? normalizeTaskPriority(taskRow.priority);
let entry: MergeQueueEntry;
let nextEntry: MergeQueueEntry;
let alreadyEnqueued = true;
if (existing) {
entry = this.rowToMergeQueueEntry(existing);
nextEntry = this.rowToMergeQueueEntry(existing);
} else {
this.db.prepare(`
INSERT INTO mergeQueue (taskId, enqueuedAt, priority, attemptCount)
@@ -6158,7 +6174,7 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
if (!inserted) {
throw new Error(`Failed to read merge queue entry for ${taskId} after enqueue`);
}
entry = this.rowToMergeQueueEntry(inserted);
nextEntry = this.rowToMergeQueueEntry(inserted);
alreadyEnqueued = false;
}
@@ -6169,13 +6185,111 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
target: taskId,
metadata: {
taskId,
priority: entry.priority,
enqueuedAt: entry.enqueuedAt,
priority: nextEntry.priority,
enqueuedAt: nextEntry.enqueuedAt,
alreadyEnqueued,
},
});
return entry;
return nextEntry;
});
if (invalidColumn) {
this.db.transactionImmediate(() => {
this.insertRunAuditEventRow({
taskId,
domain: "database",
mutationType: "mergeQueue:enqueue-rejected",
target: taskId,
metadata: {
taskId,
column: invalidColumn,
reason: "not-in-review",
},
});
});
throw new MergeQueueInvalidColumnError(taskId, invalidColumn);
}
if (!entry) {
throw new Error(`Failed to enqueue merge queue entry for ${taskId}`);
}
return entry;
}
private cleanupStaleMergeQueueRows(now: string): void {
const staleRows = this.db.prepare(`
SELECT mq.taskId, mq.leasedBy, mq.leaseExpiresAt, t.column
FROM mergeQueue mq
LEFT JOIN tasks t ON t.id = mq.taskId
WHERE t.id IS NULL OR t.column != 'in-review'
`).all() as Array<{ taskId: string; leasedBy: string | null; leaseExpiresAt: string | null; column: Column | null }>;
for (const staleRow of staleRows) {
this.db.prepare("DELETE FROM mergeQueue WHERE taskId = ?").run(staleRow.taskId);
this.insertRunAuditEventRow({
taskId: staleRow.taskId,
domain: "database",
mutationType: "mergeQueue:auto-cleanup-stale-row",
target: staleRow.taskId,
metadata: {
taskId: staleRow.taskId,
column: staleRow.column,
leasedBy: staleRow.leasedBy,
leaseExpiresAt: staleRow.leaseExpiresAt,
cleanedAt: now,
reason: "not-in-review",
},
});
}
}
private dequeueMergeQueueOnColumnExit(taskId: string, previousColumn: Column, nextColumn: Column, now: string): void {
if (previousColumn !== "in-review" || nextColumn === "in-review") {
return;
}
const queueRow = this.db.prepare("SELECT leasedBy, leaseExpiresAt FROM mergeQueue WHERE taskId = ?").get(taskId) as {
leasedBy: string | null;
leaseExpiresAt: string | null;
} | undefined;
if (!queueRow) {
return;
}
const leaseIsExpired = queueRow.leaseExpiresAt != null && queueRow.leaseExpiresAt <= now;
if (!queueRow.leasedBy || leaseIsExpired) {
this.db.prepare("DELETE FROM mergeQueue WHERE taskId = ?").run(taskId);
this.insertRunAuditEventRow({
taskId,
domain: "database",
mutationType: "mergeQueue:auto-cleanup-stale-row",
target: taskId,
metadata: {
taskId,
previousColumn,
nextColumn,
leasedBy: queueRow.leasedBy,
leaseExpiresAt: queueRow.leaseExpiresAt,
cleanedAt: now,
reason: "column-exit",
},
});
return;
}
this.insertRunAuditEventRow({
taskId,
domain: "database",
mutationType: "mergeQueue:stale-lease-on-column-exit",
target: taskId,
metadata: {
taskId,
previousColumn,
nextColumn,
leasedBy: queueRow.leasedBy,
leaseExpiresAt: queueRow.leaseExpiresAt,
},
});
}
@@ -6187,67 +6301,81 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
return this.db.transactionImmediate(() => {
const now = opts.now ?? new Date().toISOString();
const leaseExpiresAt = new Date(Date.parse(now) + opts.leaseDurationMs).toISOString();
this.cleanupStaleMergeQueueRows(now);
// Target the specific task if provided; return null immediately when unavailable
// rather than falling back to the queue head, so callers that pass targetTaskId
// can distinguish "target not available" from "no tasks available".
let leased: MergeQueueRow | undefined;
if (opts.targetTaskId) {
leased = this.db.prepare(`
UPDATE mergeQueue
SET leasedBy = ?, leasedAt = ?, leaseExpiresAt = ?
WHERE taskId = ?
AND EXISTS (
SELECT 1
FROM tasks t
WHERE t.id = mergeQueue.taskId
AND t.column = 'in-review'
)
AND (leasedBy IS NULL OR leaseExpiresAt <= ?)
RETURNING *
`).get(workerId, now, leaseExpiresAt, opts.targetTaskId, now) as MergeQueueRow | undefined;
// Do NOT fall back to queue-head when a target was explicitly requested.
// Callers (e.g. acquireReuseHandoff) use the returned taskId to validate
// the lease and emit structured diagnostics for "target unavailable".
if (!leased) {
const queueHead = this.db.prepare(`
SELECT mq.taskId, mq.leasedBy, t.column
FROM mergeQueue mq
LEFT JOIN tasks t ON t.id = mq.taskId
ORDER BY CASE mq.priority
WHEN 'urgent' THEN 0
WHEN 'high' THEN 1
WHEN 'normal' THEN 2
WHEN 'low' THEN 3
ELSE 4
END ASC,
mq.enqueuedAt ASC
LIMIT 1
`).get() as { taskId: string; leasedBy: string | null; column: string | null } | undefined;
this.insertRunAuditEventRow({
taskId: opts.targetTaskId,
domain: "database",
mutationType: "mergeQueue:lease-target-unavailable",
target: opts.targetTaskId,
metadata: {
targetTaskId: opts.targetTaskId,
workerId,
queueHeadTaskId: queueHead?.taskId ?? null,
queueHeadLeasedBy: queueHead?.leasedBy ?? null,
queueHeadColumn: queueHead?.column ?? null,
},
});
return null;
}
} else {
// Backward-compatible queue-head selection for callers that don't target a task.
leased = this.db.prepare(`
UPDATE mergeQueue
SET leasedBy = ?, leasedAt = ?, leaseExpiresAt = ?
WHERE taskId = (
SELECT taskId FROM mergeQueue
WHERE leasedBy IS NULL OR leaseExpiresAt <= ?
ORDER BY CASE priority
SELECT mq.taskId
FROM mergeQueue mq
JOIN tasks t ON t.id = mq.taskId
WHERE t.column = 'in-review'
AND (mq.leasedBy IS NULL OR mq.leaseExpiresAt <= ?)
ORDER BY CASE mq.priority
WHEN 'urgent' THEN 0
WHEN 'high' THEN 1
WHEN 'normal' THEN 2
WHEN 'low' THEN 3
ELSE 4
END ASC,
enqueuedAt ASC
mq.enqueuedAt ASC
LIMIT 1
)
RETURNING *
`).get(workerId, now, leaseExpiresAt, now) as MergeQueueRow | undefined;
}
if (!leased) {
leased = this.db.prepare(`
UPDATE mergeQueue
SET leasedBy = ?, leasedAt = ?, leaseExpiresAt = ?
WHERE taskId = (
SELECT taskId FROM mergeQueue
WHERE leasedBy IS NULL OR leaseExpiresAt <= ?
ORDER BY CASE priority
WHEN 'urgent' THEN 0
WHEN 'high' THEN 1
WHEN 'normal' THEN 2
WHEN 'low' THEN 3
ELSE 4
END ASC,
enqueuedAt ASC
LIMIT 1
)
RETURNING *
`).get(workerId, now, leaseExpiresAt, now) as MergeQueueRow | undefined;
}
if (!leased) {
return null;
if (!leased) {
return null;
}
}
const entry = this.rowToMergeQueueEntry(leased);
@@ -6377,6 +6505,24 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
return rows.map((row) => this.rowToMergeQueueEntry(row));
}
peekMergeQueueHead(): { taskId: string; leasedBy: string | null; column: Column | null } | null {
const row = this.db.prepare(`
SELECT mq.taskId, mq.leasedBy, t.column
FROM mergeQueue mq
LEFT JOIN tasks t ON t.id = mq.taskId
ORDER BY CASE mq.priority
WHEN 'urgent' THEN 0
WHEN 'high' THEN 1
WHEN 'normal' THEN 2
WHEN 'low' THEN 3
ELSE 4
END ASC,
mq.enqueuedAt ASC
LIMIT 1
`).get() as { taskId: string; leasedBy: string | null; column: Column | null } | undefined;
return row ?? null;
}
// ── End Run Audit APIs ───────────────────────────────────────────────
/**

View File

@@ -164,6 +164,7 @@ describe("acquireReuseHandoff", () => {
]);
store.acquireMergeQueueLease = vi.fn().mockReturnValue({ taskId: "FN-5279" });
store.releaseMergeQueueLease = vi.fn();
store.peekMergeQueueHead = vi.fn().mockReturnValue({ taskId: "FN-5000", leasedBy: "merger-reuse-handoff", column: "todo" });
return store;
}
@@ -440,8 +441,9 @@ describe("acquireReuseHandoff", () => {
it("refuses when no merge queue lease can be acquired", async () => {
const store = createStore();
store.acquireMergeQueueLease.mockReturnValue(null);
store.peekMergeQueueHead.mockReturnValue({ taskId: "FN-5329", leasedBy: "merger-reuse-handoff", column: "todo" });
await expectRefusal(
const refusal = await expectRefusal(
acquireReuseHandoff({
task: await store.getTask("FN-5279"),
store,
@@ -452,6 +454,10 @@ describe("acquireReuseHandoff", () => {
"lease-handoff-failed",
"no-lease",
);
expect(refusal.payload).toMatchObject({
queueHeadTaskId: "FN-5329",
queueHeadLeasedBy: "merger-reuse-handoff",
});
});
// FN-5363 regression: when the merge queue head is polluted with unrelated tasks

View File

@@ -251,6 +251,13 @@ describe("FN-5279 reliability interactions: merge reuse task worktree", () => {
await mkdir(worktreeRoot, { recursive: true });
git(rootDir, `git worktree add ${JSON.stringify(worktreePath)} ${JSON.stringify(branch)}`);
await store.updateTask(task.id, { worktree: worktreePath, branch } as any);
store.enqueueMergeQueue(task.id, { now: "2026-05-19T00:00:00.000Z" });
store.getDatabase().prepare("UPDATE mergeQueue SET leasedBy = ?, leasedAt = ?, leaseExpiresAt = ? WHERE taskId = ?").run(
"worker-other",
"2026-05-19T00:01:00.000Z",
"2099-05-19T00:10:00.000Z",
task.id,
);
await expect(aiMergeTask(store, rootDir, task.id)).rejects.toMatchObject({
name: "MergeHandoffRefusedError",
@@ -258,7 +265,117 @@ describe("FN-5279 reliability interactions: merge reuse task worktree", () => {
reason: "no-lease",
});
const refused = store.getRunAuditEvents({ taskId: task.id }).find((event) => event.mutationType === "merge:reuse-handoff-refused");
expect(refused?.metadata).toMatchObject({ gate: "lease-handoff-failed", reason: "no-lease" });
expect(refused?.metadata).toMatchObject({
gate: "lease-handoff-failed",
reason: "no-lease",
});
} finally {
await fixture.cleanup();
}
}, 30_000);
it.skipIf(!hasGit)("FN-5363: queue-head pollution by non-in-review tasks does not block target reuse handoff", async () => {
const fixture = await makeReliabilityFixture({
taskId: "FN-5363-RI-POLLUTED",
settings: {
baseBranch: "master",
mergeIntegrationWorktree: "reuse-task-worktree",
worktreeRebaseRemote: "origin",
} as any,
});
try {
const { rootDir, store, task } = fixture;
const actualTask = await store.getTask(task.id);
const branch = `fusion/${actualTask!.id.toLowerCase()}`;
const worktreeRoot = `${rootDir}-worktrees`;
const worktreePath = join(worktreeRoot, actualTask!.id.toLowerCase());
git(rootDir, "git branch -m main master");
const completedSteps = (actualTask?.steps ?? []).map((step) => ({ ...step, status: "done" as const }));
await store.updateTask(task.id, { baseBranch: "master", branch, steps: completedSteps, currentStep: completedSteps.length } as any);
await fixture.createBranch(branch);
await fixture.writeAndCommit("packages/engine/src/fn-5363-ri-polluted.ts", "export const polluted = true;\n", "feat: add polluted queue merge content");
await fixture.checkout("master");
await mkdir(worktreeRoot, { recursive: true });
git(rootDir, `git worktree add ${JSON.stringify(worktreePath)} ${JSON.stringify(branch)}`);
await store.updateTask(task.id, { worktree: worktreePath, branch } as any);
store.enqueueMergeQueue(task.id, { now: "2026-05-19T00:00:02.000Z" });
const todoTask = await store.createTask({ description: "polluter todo", priority: "normal" });
await store.moveTask(todoTask.id, "todo");
const inProgressTask = await store.createTask({ description: "polluter progress", priority: "normal" });
await store.moveTask(inProgressTask.id, "todo");
await store.moveTask(inProgressTask.id, "in-progress");
store.getDatabase().prepare("INSERT INTO mergeQueue (taskId, enqueuedAt, priority, attemptCount) VALUES (?, ?, ?, 0)").run(todoTask.id, "2026-05-19T00:00:00.000Z", "normal");
store.getDatabase().prepare("INSERT INTO mergeQueue (taskId, enqueuedAt, priority, attemptCount) VALUES (?, ?, ?, 0)").run(inProgressTask.id, "2026-05-19T00:00:01.000Z", "normal");
store.getDatabase().prepare("UPDATE mergeQueue SET leasedBy = ?, leasedAt = ?, leaseExpiresAt = ? WHERE taskId = ?").run(
"merger-reuse-handoff",
"2026-05-19T00:10:00.000Z",
"2099-05-19T00:20:00.000Z",
todoTask.id,
);
const result = await aiMergeTask(store, rootDir, task.id);
expect(result.merged).toBe(true);
expect((await store.getTask(task.id))?.column).toBe("done");
expect(store.getDatabase().prepare("SELECT leasedBy FROM mergeQueue WHERE taskId = ?").get(task.id)).toBeUndefined();
expect(store.getDatabase().prepare("SELECT taskId FROM mergeQueue WHERE taskId IN (?, ?)").all(todoTask.id, inProgressTask.id)).toEqual([]);
} finally {
await fixture.cleanup();
}
}, 30_000);
it.skipIf(!hasGit)("FN-5363: target row leased by another worker refuses with no-lease and queue-head diagnostics", async () => {
const fixture = await makeReliabilityFixture({
taskId: "FN-5363-RI-NO-LEASE-TARGET",
settings: {
baseBranch: "master",
mergeIntegrationWorktree: "reuse-task-worktree",
} as any,
});
try {
const { rootDir, store, task } = fixture;
const actualTask = await store.getTask(task.id);
const branch = `fusion/${actualTask!.id.toLowerCase()}`;
const worktreeRoot = `${rootDir}-worktrees`;
const worktreePath = join(worktreeRoot, actualTask!.id.toLowerCase());
git(rootDir, "git branch -m main master");
const completedSteps = (actualTask?.steps ?? []).map((step) => ({ ...step, status: "done" as const }));
await store.updateTask(task.id, { baseBranch: "master", branch, steps: completedSteps, currentStep: completedSteps.length } as any);
await store.moveTask(task.id, "todo");
await store.moveTask(task.id, "in-progress");
await store.handoffToReview(task.id, {
ownerAgentId: "agent-1",
evidence: { reason: "fn_task_done", runId: "run-1", agentId: "agent-1" },
});
await store.updateTask(task.id, {
steps: completedSteps,
currentStep: completedSteps.length,
} as any);
await fixture.createBranch(branch);
await fixture.writeAndCommit("packages/engine/src/fn-5363-ri-no-lease-target.ts", "export const noLeaseTarget = true;\n", "feat: add leased target merge content");
await fixture.checkout("master");
await mkdir(worktreeRoot, { recursive: true });
git(rootDir, `git worktree add ${JSON.stringify(worktreePath)} ${JSON.stringify(branch)}`);
await store.updateTask(task.id, { worktree: worktreePath, branch } as any);
store.getDatabase().prepare("UPDATE mergeQueue SET leasedBy = ?, leasedAt = ?, leaseExpiresAt = ? WHERE taskId = ?").run(
"worker-other",
"2026-05-19T00:01:00.000Z",
"2099-05-19T00:10:00.000Z",
task.id,
);
await expect(aiMergeTask(store, rootDir, task.id)).rejects.toMatchObject({
name: "MergeHandoffRefusedError",
gate: "lease-handoff-failed",
reason: "no-lease",
});
const refused = store.getRunAuditEvents({ taskId: task.id }).find((event) => event.mutationType === "merge:reuse-handoff-refused");
expect(refused?.metadata).toMatchObject({ reason: "no-lease" });
} finally {
await fixture.cleanup();
}

View File

@@ -423,10 +423,15 @@ export async function acquireReuseHandoff(input: ReuseHandoffInput): Promise<Han
throw error;
}
if (!lease || !("taskId" in lease) || lease.taskId !== input.task.id) {
const queueHead = (input.store as TaskStore & {
peekMergeQueueHead?: () => { taskId: string; leasedBy: string | null; column: string | null } | null;
}).peekMergeQueueHead?.();
throw new MergeHandoffRefusedError("lease-handoff-failed", "no-lease", {
taskId: input.task.id,
worktreePath,
acquiredTaskId: lease && "taskId" in lease ? lease.taskId : null,
queueHeadTaskId: queueHead?.taskId ?? null,
queueHeadLeasedBy: queueHead?.leasedBy ?? null,
});
}
// Re-check executor lease after acquiring the merge-queue lease: the

View File

@@ -193,6 +193,10 @@ export type DatabaseMutationType =
| "task:pause"
| "task:unpause"
| "task:dependency:add"
| "mergeQueue:lease-target-unavailable"
| "mergeQueue:enqueue-rejected"
| "mergeQueue:stale-lease-on-column-exit"
| "mergeQueue:auto-cleanup-stale-row"
| "task:auto-recover-already-merged"
| "task:auto-recover-finalize-already-on-main"
| "task:auto-merge-skipped-already-done"