test(FN-4826): complete Step 4 — cover handoff telemetry paths
Fusion-Task-Id: FN-4826 Fusion-Task-Lineage: b08a58c8-66ab-43d4-bdf9-69d2b56b7bcc
This commit is contained in:
committed by
gsxdsm
parent
3453e77396
commit
1553f1dd0d
@@ -46,7 +46,7 @@ describe("MeshLeaseManager owning-node handoff integration", () => {
|
||||
localNodeId: "node-local",
|
||||
getHandoffPolicy: async () => policy,
|
||||
nodeHealthMonitor: {
|
||||
getNodeHealth: () => "offline",
|
||||
getNodeHealth: () => (ownerNodeId === "node-local" ? "online" : "offline"),
|
||||
} as unknown as NodeHealthMonitor,
|
||||
});
|
||||
return manager.recoverAbandonedLease(taskId, "test-owner-unavailable", { preserveProgress: true });
|
||||
@@ -54,28 +54,51 @@ describe("MeshLeaseManager owning-node handoff integration", () => {
|
||||
|
||||
it("applies handoff policy matrix for peer-owned leases", async () => {
|
||||
await seedLease("node-peer");
|
||||
let baselineEventIds = new Set(taskStore.getRunAuditEvents({ taskId, limit: 200 }).map((event) => event.id));
|
||||
expect(await runCase("block", "node-peer")).toBe(false);
|
||||
let task = await taskStore.getTask(taskId);
|
||||
expect(task?.checkedOutBy).toBe("agent-1");
|
||||
let newEvents = taskStore.getRunAuditEvents({ taskId, limit: 200 }).filter((event) => !baselineEventIds.has(event.id));
|
||||
expect(newEvents.some((event) => event.mutationType === "node:handoff:parked" && event.metadata?.source === "mesh-lease.recover" && event.metadata?.decisionReason === "handoff_blocked_by_policy")).toBe(true);
|
||||
expect(newEvents.some((event) => event.mutationType === "node:lease:recovered")).toBe(false);
|
||||
|
||||
await seedLease("node-peer");
|
||||
baselineEventIds = new Set(taskStore.getRunAuditEvents({ taskId, limit: 200 }).map((event) => event.id));
|
||||
expect(await runCase("reassign-to-local", "node-peer")).toBe(true);
|
||||
task = await taskStore.getTask(taskId);
|
||||
expect(task?.checkedOutBy ?? null).toBeNull();
|
||||
newEvents = taskStore.getRunAuditEvents({ taskId, limit: 200 }).filter((event) => !baselineEventIds.has(event.id));
|
||||
const localRecoveryEvent = newEvents.find((event) => event.mutationType === "node:lease:recovered");
|
||||
expect(localRecoveryEvent).toBeTruthy();
|
||||
expect(localRecoveryEvent?.metadata?.source).toBe("mesh-lease.recover");
|
||||
expect(typeof localRecoveryEvent?.metadata?.epoch).toBe("number");
|
||||
expect(String(localRecoveryEvent?.metadata?.recoveryReason ?? "")).toContain("test-owner-unavailable");
|
||||
|
||||
await seedLease("node-peer");
|
||||
baselineEventIds = new Set(taskStore.getRunAuditEvents({ taskId, limit: 200 }).map((event) => event.id));
|
||||
expect(await runCase("reassign-any-healthy", "node-peer")).toBe(true);
|
||||
task = await taskStore.getTask(taskId);
|
||||
expect(task?.checkedOutBy ?? null).toBeNull();
|
||||
newEvents = taskStore.getRunAuditEvents({ taskId, limit: 200 }).filter((event) => !baselineEventIds.has(event.id));
|
||||
const anyRecoveryEvent = newEvents.find((event) => event.mutationType === "node:lease:recovered");
|
||||
expect(anyRecoveryEvent).toBeTruthy();
|
||||
expect(anyRecoveryEvent?.metadata?.source).toBe("mesh-lease.recover");
|
||||
expect(typeof anyRecoveryEvent?.metadata?.epoch).toBe("number");
|
||||
expect(String(anyRecoveryEvent?.metadata?.recoveryReason ?? "")).toContain("test-owner-unavailable");
|
||||
});
|
||||
|
||||
it("recovers self-owned leases regardless of policy", async () => {
|
||||
for (const policy of ["block", "reassign-to-local", "reassign-any-healthy"] as const) {
|
||||
await seedLease("node-local");
|
||||
const baselineEventIds = new Set(taskStore.getRunAuditEvents({ taskId, limit: 200 }).map((event) => event.id));
|
||||
const recovered = await runCase(policy, "node-local");
|
||||
expect(recovered).toBe(true);
|
||||
const task = await taskStore.getTask(taskId);
|
||||
expect(task?.checkedOutBy ?? null).toBeNull();
|
||||
const newEvents = taskStore.getRunAuditEvents({ taskId, limit: 200 }).filter((event) => !baselineEventIds.has(event.id));
|
||||
const recoveryEvents = newEvents.filter((event) => event.mutationType === "node:lease:recovered");
|
||||
expect(recoveryEvents).toHaveLength(1);
|
||||
expect(recoveryEvents[0]?.metadata?.source).toBe("mesh-lease.recover");
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
@@ -59,6 +59,7 @@ function createMockStore(task: Task, settings: Record<string, unknown> = {}): Ta
|
||||
moveTask: vi.fn().mockResolvedValue(undefined),
|
||||
parseFileScopeFromPrompt: vi.fn().mockResolvedValue([]),
|
||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||
recordRunAuditEvent: vi.fn().mockResolvedValue(undefined),
|
||||
getRootDir: vi.fn().mockReturnValue("/tmp/test"),
|
||||
getTasksDir: vi.fn().mockReturnValue("/tmp/test/.fusion/tasks"),
|
||||
on: vi.fn(),
|
||||
@@ -421,6 +422,23 @@ describe("Scheduler node routing", () => {
|
||||
task.id,
|
||||
"Owning-node handoff parked dispatch: handoff_blocked_by_policy",
|
||||
);
|
||||
const handoffEvent = vi.mocked(store.recordRunAuditEvent).mock.calls
|
||||
.map(([event]) => event)
|
||||
.find((event) => event.mutationType === "node:handoff:parked");
|
||||
expect(handoffEvent).toEqual(expect.objectContaining({
|
||||
domain: "database",
|
||||
mutationType: "node:handoff:parked",
|
||||
target: task.id,
|
||||
metadata: {
|
||||
taskId: task.id,
|
||||
ownerNodeId: "node-owner",
|
||||
ownerNodeHealth: "offline",
|
||||
localNodeId: "local",
|
||||
handoffPolicy: "block",
|
||||
decisionReason: "handoff_blocked_by_policy",
|
||||
source: "scheduler.dispatch",
|
||||
},
|
||||
}));
|
||||
});
|
||||
|
||||
it("forces local dispatch when owning-node handoff returns reassign-local", async () => {
|
||||
@@ -442,6 +460,23 @@ describe("Scheduler node routing", () => {
|
||||
effectiveNodeSource: "local",
|
||||
}));
|
||||
expect(store.logEntry).toHaveBeenCalledWith(task.id, "Owning-node handoff applied: owner_offline_local_takes_over");
|
||||
const handoffEvent = vi.mocked(store.recordRunAuditEvent).mock.calls
|
||||
.map(([event]) => event)
|
||||
.find((event) => event.mutationType === "node:handoff:reassign-local");
|
||||
expect(handoffEvent).toEqual(expect.objectContaining({
|
||||
domain: "database",
|
||||
mutationType: "node:handoff:reassign-local",
|
||||
target: task.id,
|
||||
metadata: {
|
||||
taskId: task.id,
|
||||
ownerNodeId: "node-owner",
|
||||
ownerNodeHealth: "offline",
|
||||
localNodeId: "local",
|
||||
handoffPolicy: "reassign-to-local",
|
||||
decisionReason: "owner_offline_local_takes_over",
|
||||
source: "scheduler.dispatch",
|
||||
},
|
||||
}));
|
||||
});
|
||||
|
||||
it("keeps non-local routing when owning-node handoff returns reassign-any", async () => {
|
||||
@@ -463,6 +498,23 @@ describe("Scheduler node routing", () => {
|
||||
effectiveNodeSource: "task-override",
|
||||
}));
|
||||
expect(store.logEntry).toHaveBeenCalledWith(task.id, "Owning-node handoff applied: owner_offline_any_healthy_eligible");
|
||||
const handoffEvent = vi.mocked(store.recordRunAuditEvent).mock.calls
|
||||
.map(([event]) => event)
|
||||
.find((event) => event.mutationType === "node:handoff:reassign-any");
|
||||
expect(handoffEvent).toEqual(expect.objectContaining({
|
||||
domain: "database",
|
||||
mutationType: "node:handoff:reassign-any",
|
||||
target: task.id,
|
||||
metadata: {
|
||||
taskId: task.id,
|
||||
ownerNodeId: "node-owner",
|
||||
ownerNodeHealth: "offline",
|
||||
localNodeId: "local",
|
||||
handoffPolicy: "reassign-any-healthy",
|
||||
decisionReason: "owner_offline_any_healthy_eligible",
|
||||
source: "scheduler.dispatch",
|
||||
},
|
||||
}));
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -406,6 +406,9 @@ export class MeshLeaseManager {
|
||||
});
|
||||
meshLeaseManagerLog.log(`mesh-lease: handoff parked taskId=${task.id} reason=${handoffDecision.reason}`);
|
||||
await (this.options.taskStore as any).recordRunAuditEvent?.({
|
||||
taskId: task.id,
|
||||
agentId: "mesh-lease-manager",
|
||||
runId: generateSyntheticRunId("mesh-lease", task.id),
|
||||
domain: "database",
|
||||
mutationType: "node:handoff:parked",
|
||||
target: task.id,
|
||||
@@ -470,6 +473,9 @@ export class MeshLeaseManager {
|
||||
}
|
||||
|
||||
await (this.options.taskStore as any).recordRunAuditEvent?.({
|
||||
taskId: task.id,
|
||||
agentId: "mesh-lease-manager",
|
||||
runId: generateSyntheticRunId("mesh-lease", task.id),
|
||||
domain: "database",
|
||||
mutationType: "node:lease:recovered",
|
||||
target: task.id,
|
||||
|
||||
@@ -1075,6 +1075,9 @@ export class Scheduler {
|
||||
schedulerLog.log(`Task ${task.id} dispatch blocked — ${reason}`);
|
||||
await this.store.logEntry(task.id, reason);
|
||||
await (this.store as any).recordRunAuditEvent?.({
|
||||
taskId: freshTask.id,
|
||||
agentId: "scheduler",
|
||||
runId: generateSyntheticRunId("scheduler", freshTask.id),
|
||||
domain: "database",
|
||||
mutationType: "node:handoff:parked",
|
||||
target: freshTask.id,
|
||||
@@ -1097,6 +1100,9 @@ export class Scheduler {
|
||||
|
||||
await this.store.logEntry(task.id, `Owning-node handoff applied: ${handoffDecision.reason}`);
|
||||
await (this.store as any).recordRunAuditEvent?.({
|
||||
taskId: freshTask.id,
|
||||
agentId: "scheduler",
|
||||
runId: generateSyntheticRunId("scheduler", freshTask.id),
|
||||
domain: "database",
|
||||
mutationType: handoffDecision.action === "reassign-local" ? "node:handoff:reassign-local" : "node:handoff:reassign-any",
|
||||
target: freshTask.id,
|
||||
|
||||
Reference in New Issue
Block a user