fix(FN-4826): complete Step 6 — stabilize audit emission and tests

Fusion-Task-Id: FN-4826
Fusion-Task-Lineage: b08a58c8-66ab-43d4-bdf9-69d2b56b7bcc
This commit is contained in:
Fusion (runfusion.ai)
2026-05-16 23:41:51 -07:00
committed by gsxdsm
parent 9e5a516997
commit 01f0bf625a
4 changed files with 120 additions and 96 deletions

View File

@@ -69,9 +69,11 @@ describe("MeshLeaseManager", () => {
expect(ok).toBe(true);
expect(moveTask).toHaveBeenCalledWith("FN-1", "todo", expect.any(Object));
const event = recordRunAuditEvent.mock.calls[0]?.[0] as RunAuditEventInput;
expect(event.mutationType).toBe("task:auto-recover-node-unreachable");
expect(event.metadata).toMatchObject({
const event = recordRunAuditEvent.mock.calls
.map((call) => call[0] as RunAuditEventInput)
.find((candidate) => candidate.mutationType === "task:auto-recover-node-unreachable");
expect(event?.mutationType).toBe("task:auto-recover-node-unreachable");
expect(event?.metadata).toMatchObject({
ownerNodeId: "node-a",
ownerNodeHealth: "offline",
previousOwnerAgentId: "agent-1",
@@ -107,8 +109,10 @@ describe("MeshLeaseManager", () => {
const ok = await manager.recoverAbandonedLease("FN-1", "scheduler detected stale todo lease");
expect(ok).toBe(true);
const event = recordRunAuditEvent.mock.calls[0]?.[0] as RunAuditEventInput;
expect(event.metadata).toMatchObject({
const event = recordRunAuditEvent.mock.calls
.map((call) => call[0] as RunAuditEventInput)
.find((candidate) => candidate.mutationType === "task:auto-recover-node-unreachable");
expect(event?.metadata).toMatchObject({
ownerNodeHealth: "error",
decisionPath: "lease-recovered-in-place",
newColumn: "todo",
@@ -136,9 +140,11 @@ describe("MeshLeaseManager", () => {
const ok = await manager.recoverAbandonedLease("FN-1", "scheduler detected stale todo lease");
expect(ok).toBe(false);
const event = recordRunAuditEvent.mock.calls[0]?.[0] as RunAuditEventInput;
expect(event.mutationType).toBe("task:auto-recover-node-unreachable");
expect(event.metadata).toMatchObject({
const event = recordRunAuditEvent.mock.calls
.map((call) => call[0] as RunAuditEventInput)
.find((candidate) => candidate.mutationType === "task:auto-recover-node-unreachable");
expect(event?.mutationType).toBe("task:auto-recover-node-unreachable");
expect(event?.metadata).toMatchObject({
decisionPath: "lease-parked-by-handoff-policy",
recoveryReason: "handoff-policy-park",
handoffPolicy: "block",
@@ -170,7 +176,8 @@ describe("MeshLeaseManager", () => {
const manager = new MeshLeaseManager({ taskStore, agentStore });
const ok = await manager.recoverAbandonedLease("FN-1", "stale-heartbeat");
expect(ok).toBe(true);
expect(recordRunAuditEvent).not.toHaveBeenCalled();
expect(recordRunAuditEvent.mock.calls.some((call) => call[0].mutationType === "task:auto-recover-node-unreachable")).toBe(false);
expect(recordRunAuditEvent.mock.calls.some((call) => call[0].mutationType === "node:lease:recovered")).toBe(true);
});
it("emits lease-released when central + local recovery succeeds", async () => {

View File

@@ -73,10 +73,10 @@ describe("Scheduler node-unreachable audit", () => {
await scheduler.schedule();
await scheduler.schedule();
expect(recordRunAuditEvent).toHaveBeenCalledTimes(1);
const event = recordRunAuditEvent.mock.calls[0][0] as RunAuditEventInput;
expect(event.mutationType).toBe("task:auto-recover-node-unreachable");
expect(event.metadata).toMatchObject({
const events = recordRunAuditEvent.mock.calls.map(([event]) => event as RunAuditEventInput);
const event = events.find((candidate) => candidate.mutationType === "task:auto-recover-node-unreachable");
expect(event?.mutationType).toBe("task:auto-recover-node-unreachable");
expect(event?.metadata).toMatchObject({
handoffAction: "park",
decisionPath: "scheduler-handoff-park",
ownerNodeId: "node-owner",
@@ -100,8 +100,9 @@ describe("Scheduler node-unreachable audit", () => {
await scheduler.schedule();
const event = recordRunAuditEvent.mock.calls[0][0] as RunAuditEventInput;
expect(event.metadata).toMatchObject({
const events = recordRunAuditEvent.mock.calls.map(([event]) => event as RunAuditEventInput);
const event = events.find((candidate) => candidate.mutationType === "task:auto-recover-node-unreachable");
expect(event?.metadata).toMatchObject({
handoffAction: "reassign-local",
decisionPath: "scheduler-handoff-reassign-local",
dispatchNodeBefore: "node-task",
@@ -118,9 +119,9 @@ describe("Scheduler node-unreachable audit", () => {
await scheduler.schedule();
expect(recordRunAuditEvent).toHaveBeenCalledTimes(1);
const event = recordRunAuditEvent.mock.calls[0][0] as RunAuditEventInput;
expect(event.metadata).toMatchObject({
const events = recordRunAuditEvent.mock.calls.map(([event]) => event as RunAuditEventInput);
const event = events.find((candidate) => candidate.mutationType === "task:auto-recover-node-unreachable");
expect(event?.metadata).toMatchObject({
handoffAction: "park",
decisionPath: "scheduler-handoff-park",
ownerNodeId: "node-owner",

View File

@@ -405,29 +405,33 @@ export class MeshLeaseManager {
handoffReason: handoffDecision.reason,
});
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,
metadata: {
try {
await this.options.taskStore.recordRunAuditEvent?.({
taskId: task.id,
ownerNodeId,
ownerNodeHealth:
currentOwnerNodeHealth === "offline" ||
currentOwnerNodeHealth === "error" ||
currentOwnerNodeHealth === "online"
? currentOwnerNodeHealth
: "unknown",
localNodeId,
handoffPolicy,
decisionReason: handoffDecision.reason,
source: "mesh-lease.recover",
recoveryReason: reason,
},
});
agentId: "mesh-lease-manager",
runId: generateSyntheticRunId("mesh-lease", task.id),
domain: "database",
mutationType: "node:handoff:parked",
target: task.id,
metadata: {
taskId: task.id,
ownerNodeId,
ownerNodeHealth:
currentOwnerNodeHealth === "offline" ||
currentOwnerNodeHealth === "error" ||
currentOwnerNodeHealth === "online"
? currentOwnerNodeHealth
: "unknown",
localNodeId,
handoffPolicy,
decisionReason: handoffDecision.reason,
source: "mesh-lease.recover",
recoveryReason: reason,
},
});
} catch (error) {
meshLeaseManagerLog.warn(`mesh-lease: failed to emit node:handoff:parked for taskId=${task.id}: ${error instanceof Error ? error.message : String(error)}`);
}
return false;
}
}
@@ -472,25 +476,29 @@ 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,
metadata: {
try {
await this.options.taskStore.recordRunAuditEvent?.({
taskId: task.id,
ownerNodeId,
ownerNodeHealth: normalizedOwnerNodeHealth,
localNodeId,
handoffPolicy,
decisionReason: handoffReason,
source: "mesh-lease.recover",
epoch: nextEpoch,
recoveryReason: `${reason} (${stale.reason ?? "stale"})`,
},
});
agentId: "mesh-lease-manager",
runId: generateSyntheticRunId("mesh-lease", task.id),
domain: "database",
mutationType: "node:lease:recovered",
target: task.id,
metadata: {
taskId: task.id,
ownerNodeId,
ownerNodeHealth: normalizedOwnerNodeHealth,
localNodeId,
handoffPolicy,
decisionReason: handoffReason,
source: "mesh-lease.recover",
epoch: nextEpoch,
recoveryReason: `${reason} (${stale.reason ?? "stale"})`,
},
});
} catch (error) {
meshLeaseManagerLog.warn(`mesh-lease: failed to emit node:lease:recovered for taskId=${task.id}: ${error instanceof Error ? error.message : String(error)}`);
}
if (isUnreachableOwnerReason) {
await emitNodeUnreachableRecovery({

View File

@@ -1074,51 +1074,59 @@ export class Scheduler {
const reason = `Owning-node handoff parked dispatch: ${handoffDecision.reason}`;
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,
metadata: {
try {
await this.store.recordRunAuditEvent?.({
taskId: freshTask.id,
ownerNodeId: freshTask.checkoutNodeId,
ownerNodeHealth:
ownerNodeHealth === "offline" || ownerNodeHealth === "error" || ownerNodeHealth === "online"
? ownerNodeHealth
: "unknown",
localNodeId,
handoffPolicy: settings.owningNodeHandoffPolicy,
decisionReason: handoffDecision.reason,
source: "scheduler.dispatch",
},
});
agentId: "scheduler",
runId: generateSyntheticRunId("scheduler", freshTask.id),
domain: "database",
mutationType: "node:handoff:parked",
target: freshTask.id,
metadata: {
taskId: freshTask.id,
ownerNodeId: freshTask.checkoutNodeId,
ownerNodeHealth:
ownerNodeHealth === "offline" || ownerNodeHealth === "error" || ownerNodeHealth === "online"
? ownerNodeHealth
: "unknown",
localNodeId,
handoffPolicy: settings.owningNodeHandoffPolicy,
decisionReason: handoffDecision.reason,
source: "scheduler.dispatch",
},
});
} catch (error) {
schedulerLog.warn(`Task ${task.id} failed to emit node:handoff:parked audit: ${error instanceof Error ? error.message : String(error)}`);
}
}
continue;
}
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,
metadata: {
try {
await this.store.recordRunAuditEvent?.({
taskId: freshTask.id,
ownerNodeId: freshTask.checkoutNodeId,
ownerNodeHealth:
ownerNodeHealth === "offline" || ownerNodeHealth === "error" || ownerNodeHealth === "online"
? ownerNodeHealth
: "unknown",
localNodeId,
handoffPolicy: settings.owningNodeHandoffPolicy,
decisionReason: handoffDecision.reason,
source: "scheduler.dispatch",
},
});
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,
metadata: {
taskId: freshTask.id,
ownerNodeId: freshTask.checkoutNodeId,
ownerNodeHealth:
ownerNodeHealth === "offline" || ownerNodeHealth === "error" || ownerNodeHealth === "online"
? ownerNodeHealth
: "unknown",
localNodeId,
handoffPolicy: settings.owningNodeHandoffPolicy,
decisionReason: handoffDecision.reason,
source: "scheduler.dispatch",
},
});
} catch (error) {
schedulerLog.warn(`Task ${task.id} failed to emit node:handoff audit: ${error instanceof Error ? error.message : String(error)}`);
}
const dispatchNodeBefore = effectiveNode.nodeId;
if (handoffDecision.action === "reassign-local") {
effectiveNode = { nodeId: undefined, source: "local" };