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 () => {
const layer = h.layer();
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
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:
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:
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.
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 {
if (event === "task:updated" && args.length === 1) {
if (event === "task:updated") {
const task = args[0] as Task;
const metadata = args[1] as TaskStoreEvents["task:updated"][1] | undefined;
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);
}

View File

@@ -318,9 +318,14 @@ export function setupActivityLogListenersImpl(store: TaskStore): void {
);
});
// Task updated (check for failures)
store.on("task:updated", (task) => {
if (task.status === "failed") {
/*
FNXC:ActivityLogFailureDedup 2026-08-19-19:33:
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(
{
type: "task:failed",

View File

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