feat(FN-4822): complete Step 6 — enforce central-claim race coverage

Fusion-Task-Id: FN-4822
Fusion-Task-Lineage: 08cc29e8-114a-48dc-80de-8d7fd2ce0e69
This commit is contained in:
Fusion (runfusion.ai)
2026-05-16 20:42:01 -07:00
committed by gsxdsm
parent ad4b236985
commit 199f317813
3 changed files with 27 additions and 3 deletions

View File

@@ -0,0 +1,5 @@
---
"@runfusion/fusion": patch
---
Add central `taskClaims` table (central DB schema v13) and route `AgentStore.checkoutTask` through it when a `CentralClaimStore` is configured, providing the authoritative cross-node task-claim mutex required by FN-4819 §2. Single-node behavior is unchanged when no claim store is wired.

View File

@@ -1429,7 +1429,7 @@ export class AgentStore extends EventEmitter {
const isSameNodeHolder = existingNodeId !== null && requestNodeId === existingNodeId; const isSameNodeHolder = existingNodeId !== null && requestNodeId === existingNodeId;
const expectedEpoch = isSameAgentHolder && isSameNodeHolder ? existingEpoch : undefined; const expectedEpoch = isSameAgentHolder && isSameNodeHolder ? existingEpoch : undefined;
const centralResult = this.claimStore.tryClaimTask({ const centralResult = await Promise.resolve(this.claimStore.tryClaimTask({
projectId: this.claimProjectId, projectId: this.claimProjectId,
taskId, taskId,
nodeId: requestNodeId, nodeId: requestNodeId,
@@ -1437,7 +1437,7 @@ export class AgentStore extends EventEmitter {
runId: leaseContext?.runId ?? task.checkoutRunId ?? null, runId: leaseContext?.runId ?? task.checkoutRunId ?? null,
renewedAt: nextRenewedAt, renewedAt: nextRenewedAt,
expectedEpoch, expectedEpoch,
}); }));
if (!centralResult.ok) { if (!centralResult.ok) {
throw new CheckoutConflictError(taskId, centralResult.current.ownerAgentId, agentId); throw new CheckoutConflictError(taskId, centralResult.current.ownerAgentId, agentId);

View File

@@ -3,7 +3,7 @@ import { mkdtempSync } from "node:fs";
import { rm } from "node:fs/promises"; import { rm } from "node:fs/promises";
import { join } from "node:path"; import { join } from "node:path";
import { tmpdir } from "node:os"; import { tmpdir } from "node:os";
import { AgentStore, CheckoutConflictError, TaskStore, createCentralDatabase } from "@fusion/core"; import { AgentStore, CheckoutConflictError, TaskStore, createCentralDatabase, type CentralDatabase } from "@fusion/core";
function makeTmpDir(): string { function makeTmpDir(): string {
return mkdtempSync(join(tmpdir(), "fn-cross-node-claim-test-")); return mkdtempSync(join(tmpdir(), "fn-cross-node-claim-test-"));
@@ -47,6 +47,24 @@ describe("cross-node claim mutex integration", () => {
}); });
it("allows one winner per race and bumps epoch once per successful ownership acquisition", async () => { it("allows one winner per race and bumps epoch once per successful ownership acquisition", async () => {
const originalTryClaim = centralDb.tryClaimTask.bind(centralDb);
const installBarrier = () => {
let waiters = 0;
let releaseBarrier: (() => void) | undefined;
const barrier = new Promise<void>((resolve) => {
releaseBarrier = resolve;
});
centralDb.tryClaimTask = ((input) => {
waiters += 1;
if (waiters === 2) {
releaseBarrier?.();
}
return barrier.then(() => originalTryClaim(input));
}) as CentralDatabase["tryClaimTask"];
};
installBarrier();
const [first, second] = await Promise.allSettled([ const [first, second] = await Promise.allSettled([
storeA.checkoutTask(agentA, taskId, { runId: "run-a" }), storeA.checkoutTask(agentA, taskId, { runId: "run-a" }),
storeB.checkoutTask(agentB, taskId, { runId: "run-b" }), storeB.checkoutTask(agentB, taskId, { runId: "run-b" }),
@@ -67,6 +85,7 @@ describe("cross-node claim mutex integration", () => {
await (winner.checkedOutBy === agentA ? storeA : storeB).releaseTask(winner.checkedOutBy ?? "", taskId); await (winner.checkedOutBy === agentA ? storeA : storeB).releaseTask(winner.checkedOutBy ?? "", taskId);
installBarrier();
const [third, fourth] = await Promise.allSettled([ const [third, fourth] = await Promise.allSettled([
storeA.checkoutTask(agentA, taskId, { runId: "run-c" }), storeA.checkoutTask(agentA, taskId, { runId: "run-c" }),
storeB.checkoutTask(agentB, taskId, { runId: "run-d" }), storeB.checkoutTask(agentB, taskId, { runId: "run-d" }),