FN-059: deduplicate failure activity logging

Record task failure activity only when a task enters the failed state, while preserving lifecycle metadata compatibility.

- Emit failure-transition metadata only for non-failed to failed updates.
- Preserve warm-cache lane decoration and cold-cache absent metadata.
- Add PostgreSQL parity coverage and a published-package changeset.

Files changed:
 .changeset/fn-059-activity-log-failure-dedup.md    |  7 +++
 .../postgres/activity-log-parity.pg.test.ts        | 50 ++++++++++++++++++++++
 packages/core/src/store.ts                         | 13 ++++--
 packages/core/src/task-store/lifecycle-ops.ts      | 11 +++--
 packages/core/src/task-store/task-update.ts        |  6 ++-

 5 files changed, 80 insertions(+), 7 deletions(-)

Fusion-Task-Id: FN-059

Fusion-Task-Lineage: 28c28c0d-d325-4fa6-844d-33a05a18e5ed

Co-authored-by: Fusion <noreply@runfusion.ai>
This commit is contained in:
Fusion Agent
2026-08-19 20:03:27 +00:00
parent b533220716
commit 47a8b53ad2
5 changed files with 80 additions and 7 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Prevent duplicate Task Failed entries after subsequent task updates.
category: fix
dev: Records failure activity only on a non-failed-to-failed task transition.

View File

@@ -63,6 +63,56 @@ pgDescribe("activity log parity (PostgreSQL)", () => {
}); });
}); });
it("records one failure activity row per failed episode", async () => {
const store = h.store();
const task = await store.createTask({
title: "Deduplicate failed activity",
description: "Failure activity must only record episode transitions",
});
const error = "BLOCKED: prerequisite unavailable";
await store.updateTask(task.id, { status: "failed", error });
await store.updateTask(task.id, { title: "Deduplicate failed activity (updated)" });
await store.updateTask(task.id, { priority: "high" });
await vi.waitFor(async () => {
const failures = await store.getActivityLog({ type: "task:failed" });
expect(failures.filter((entry) => entry.taskId === task.id && entry.metadata?.error === error)).toHaveLength(1);
});
// An identical activity row in another project must not influence this project's feed.
await h.adminDb().insert(schema.project.activityLog).values({
projectId: "other-project",
id: "other-project-failure",
timestamp: new Date().toISOString(),
type: "task:failed",
taskId: task.id,
details: `Task ${task.id} failed: ${error}`,
metadata: { error },
});
expect((await store.getActivityLog({ type: "task:failed" }))
.filter((entry) => entry.taskId === task.id && entry.metadata?.error === error)).toHaveLength(1);
await store.updateTask(task.id, { status: null });
await store.updateTask(task.id, { status: "failed", error });
await store.updateTask(task.id, { description: "Still in the second failure episode" });
await vi.waitFor(async () => {
const failures = await store.getActivityLog({ type: "task:failed" });
expect(failures.filter((entry) => entry.taskId === task.id && entry.metadata?.error === error)).toHaveLength(2);
});
const errorless = await store.createTask({ description: "Errorless failed activity" });
await store.updateTask(errorless.id, { status: "failed" });
await store.updateTask(errorless.id, { title: "Errorless failed activity (updated)" });
await vi.waitFor(async () => {
const failures = await store.getActivityLog({ type: "task:failed" });
expect(failures.filter((entry) => entry.taskId === errorless.id)).toEqual([
expect.objectContaining({ metadata: undefined }),
]);
});
});
it("keeps failed writes best-effort so audit storage cannot break product operations", async () => { it("keeps failed writes best-effort so audit storage cannot break product operations", async () => {
const layer = h.layer(); const layer = h.layer();
const insert = vi.spyOn(layer.db, "insert").mockImplementation(() => { const insert = vi.spyOn(layer.db, "insert").mockImplementation(() => {

View File

@@ -205,7 +205,7 @@ export interface TaskStoreEvents {
second argument. The first argument remains the Task so existing one-argument subscribers are second argument. The first argument remains the Task so existing one-argument subscribers are
unchanged; absent metadata is unknown, never a legacy-lane claim. unchanged; absent metadata is unknown, never a legacy-lane claim.
*/ */
"task:updated": [task: Task, meta?: { lanes?: TaskMoveLanes }]; "task:updated": [task: Task, meta?: { lanes?: TaskMoveLanes; failedTransition?: boolean }];
/* /*
FNXC:CrossProcessDeleteObservation 2026-08-01-11:39: FNXC:CrossProcessDeleteObservation 2026-08-01-11:39:
Observed outbox delivery is at-least-once, including a crash-window duplicate. The explicit Observed outbox delivery is at-least-once, including a crash-window duplicate. The explicit
@@ -608,12 +608,19 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
FNXC:WorkflowEvents 2026-08-01-06:28: FNXC:WorkflowEvents 2026-08-01-06:28:
Safe lifecycle emission invokes listeners directly to isolate listener failures, so it bypasses Safe lifecycle emission invokes listeners directly to isolate listener failures, so it bypasses
EventEmitter.emit. Decorate this path too; otherwise the hot update surfaces silently miss lanes. EventEmitter.emit. Decorate this path too; otherwise the hot update surfaces silently miss lanes.
FNXC:ActivityLogFailureDedup 2026-08-19-20:03:
Update metadata now carries the failure-transition marker, but warm-cache updates must retain the
existing resolved lanes decoration and cold-cache updates must preserve absent metadata.
*/ */
public emitTaskLifecycleEventSafely( event: "task:created" | "task:updated" | "task:deleted", args: TaskStoreEvents["task:created"] | TaskStoreEvents["task:updated"] | TaskStoreEvents["task:deleted"], ): boolean { public emitTaskLifecycleEventSafely( event: "task:created" | "task:updated" | "task:deleted", args: TaskStoreEvents["task:created"] | TaskStoreEvents["task:updated"] | TaskStoreEvents["task:deleted"], ): boolean {
if (event === "task:updated" && args.length === 1) { if (event === "task:updated") {
const task = args[0] as Task; const task = args[0] as Task;
const metadata = args[1] as TaskStoreEvents["task:updated"][1] | undefined;
const lanes = this.laneCache.get(task.id); const lanes = this.laneCache.get(task.id);
if (lanes !== undefined) return emitTaskLifecycleEventSafelyImpl(this, event, [task, { lanes }]); if (lanes !== undefined && metadata?.lanes === undefined) {
return emitTaskLifecycleEventSafelyImpl(this, event, [task, { ...metadata, lanes }]);
}
} }
return emitTaskLifecycleEventSafelyImpl(this, event, args); return emitTaskLifecycleEventSafelyImpl(this, event, args);
} }

View File

@@ -318,9 +318,14 @@ export function setupActivityLogListenersImpl(store: TaskStore): void {
); );
}); });
// Task updated (check for failures) /*
store.on("task:updated", (task) => { FNXC:ActivityLogFailureDedup 2026-08-19-19:33:
if (task.status === "failed") { A failed-state bookkeeping update is not a new durable failure episode. Only the serialized
non-failed → failed transition from updateTask may append task:failed activity, so listeners
cannot turn subsequent task mutations into an unbounded project.activity_log write path.
*/
store.on("task:updated", (task, meta) => {
if (meta?.failedTransition) {
store.recordActivityFromListener( store.recordActivityFromListener(
{ {
type: "task:failed", type: "task:failed",

View File

@@ -1245,7 +1245,11 @@ export async function updateTaskUnlockedImpl(store: TaskStore, id: string, updat
lanes: respecifyMoveLanes, lanes: respecifyMoveLanes,
}); });
} }
store.emitTaskLifecycleEventSafely("task:updated", [task]); const failedTransition = !wasFailed && task.status === "failed";
store.emitTaskLifecycleEventSafely("task:updated", [
task,
failedTransition ? { failedTransition: true } : undefined,
]);
return task; return task;
} }
} }