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:
committed by
gsxdsm
parent
ad4b236985
commit
199f317813
5
.changeset/FN-4822-central-task-claim-mutex.md
Normal file
5
.changeset/FN-4822-central-task-claim-mutex.md
Normal 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.
|
||||
@@ -1429,7 +1429,7 @@ export class AgentStore extends EventEmitter {
|
||||
const isSameNodeHolder = existingNodeId !== null && requestNodeId === existingNodeId;
|
||||
const expectedEpoch = isSameAgentHolder && isSameNodeHolder ? existingEpoch : undefined;
|
||||
|
||||
const centralResult = this.claimStore.tryClaimTask({
|
||||
const centralResult = await Promise.resolve(this.claimStore.tryClaimTask({
|
||||
projectId: this.claimProjectId,
|
||||
taskId,
|
||||
nodeId: requestNodeId,
|
||||
@@ -1437,7 +1437,7 @@ export class AgentStore extends EventEmitter {
|
||||
runId: leaseContext?.runId ?? task.checkoutRunId ?? null,
|
||||
renewedAt: nextRenewedAt,
|
||||
expectedEpoch,
|
||||
});
|
||||
}));
|
||||
|
||||
if (!centralResult.ok) {
|
||||
throw new CheckoutConflictError(taskId, centralResult.current.ownerAgentId, agentId);
|
||||
|
||||
@@ -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 { AgentStore, CheckoutConflictError, TaskStore, createCentralDatabase } from "@fusion/core";
|
||||
import { AgentStore, CheckoutConflictError, TaskStore, createCentralDatabase, type CentralDatabase } from "@fusion/core";
|
||||
|
||||
function makeTmpDir(): string {
|
||||
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 () => {
|
||||
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([
|
||||
storeA.checkoutTask(agentA, taskId, { runId: "run-a" }),
|
||||
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);
|
||||
|
||||
installBarrier();
|
||||
const [third, fourth] = await Promise.allSettled([
|
||||
storeA.checkoutTask(agentA, taskId, { runId: "run-c" }),
|
||||
storeB.checkoutTask(agentB, taskId, { runId: "run-d" }),
|
||||
|
||||
Reference in New Issue
Block a user