diff --git a/.changeset/fn-059-activity-log-failure-dedup.md b/.changeset/fn-059-activity-log-failure-dedup.md new file mode 100644 index 0000000000..2e6af4e669 --- /dev/null +++ b/.changeset/fn-059-activity-log-failure-dedup.md @@ -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. diff --git a/packages/core/src/__tests__/postgres/activity-log-parity.pg.test.ts b/packages/core/src/__tests__/postgres/activity-log-parity.pg.test.ts index f656b71806..d44bc83d44 100644 --- a/packages/core/src/__tests__/postgres/activity-log-parity.pg.test.ts +++ b/packages/core/src/__tests__/postgres/activity-log-parity.pg.test.ts @@ -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(() => { diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index e79ac7f37e..29df47da64 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -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 { 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); } diff --git a/packages/core/src/task-store/lifecycle-ops.ts b/packages/core/src/task-store/lifecycle-ops.ts index be504967e1..499af8381e 100644 --- a/packages/core/src/task-store/lifecycle-ops.ts +++ b/packages/core/src/task-store/lifecycle-ops.ts @@ -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", diff --git a/packages/core/src/task-store/task-update.ts b/packages/core/src/task-store/task-update.ts index bd68029fbe..32e48db046 100644 --- a/packages/core/src/task-store/task-update.ts +++ b/packages/core/src/task-store/task-update.ts @@ -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; } }