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:
7
.changeset/fn-059-activity-log-failure-dedup.md
Normal file
7
.changeset/fn-059-activity-log-failure-dedup.md
Normal 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.
|
||||
@@ -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(() => {
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user