fix(dashboard): reconcile task state through live API (#2595)

## Summary

- add a project-scoped live API route for updating individual task
checklist steps
- add an atomic live API route for resolving stale durable wedge
episodes
- prevent operator repair tooling from opening a second embedded store
that can diverge from the running dashboard backend

## Why

Legacy graph-native workflow runs can retain successful
`workflowStepResults` while their narrative checklist remains at 0/N.
The existing `fn task update` fallback may open a separate embedded
store, producing split-brain writes that do not accumulate in the live
dashboard backend. There was also no API surface for the existing atomic
wedge-episode resolver.

## Verification

- `pnpm exec vitest run
src/routes/__tests__/register-task-workflow-routes.step-update.test.ts`
— 5/5 passing
- `pnpm build` in `packages/dashboard` — passing
- full managed runtime workspace build — passing
- deployed to the managed local runtime and used to reconcile six legacy
review-deadlock tasks
- live board audit: zero `in-review-stall-deadlock` paused reasons
- exact local and Tailscale dashboard roots: HTTP 200 with 16,926-byte
bodies


<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## Summary by CodeRabbit

* **New Features**
* Added live API endpoints to update individual task checklist steps
with validation (step index and allowed status values).
* Added an endpoint to reconcile/resolve stale task “wedge” episodes,
resolving only the matching active episode and returning conflicts on
mismatches.
* **Tests**
* Expanded route tests for step updates and wedge resolution, including
consistent 404 behavior for soft-deleted and missing tasks, plus
conflict and invalid-input cases.
* Expanded PostgreSQL coverage for wedge resolution persistence and
concurrent episode replacement scenarios.
* **Bug Fixes**
* Improved task-lookup error handling so soft-deleted tasks are
consistently treated as “not found” (HTTP 404).
<!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
Phil Larson
2026-07-29 23:27:32 -07:00
committed by GitHub
parent 0c07584d51
commit 15b21dead1
11 changed files with 691 additions and 7 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Reconcile legacy task checklist and wedge state through the live dashboard backend.
category: fix
dev: Adds project-scoped task step-update and wedge-resolution API routes.

View File

@@ -0,0 +1,311 @@
/**
* FNXC:TaskStateReconciliation 2026-07-29-16:10:
* Wedge resolution is a PostgreSQL compare-and-set across dashboard processes. A resolver that waits behind a replacement write must preserve the replacement episode, while an exact active episode resolves normally.
*
* FNXC:TaskStateReconciliation 2026-07-29-22:01:
* Resolution projection must hold the task row lock through task JSON, cache, and event publication so a cross-process replacement cannot commit and publish before an older resolved snapshot.
*
* FNXC:TaskStateReconciliation 2026-07-29-22:17:
* A deletion that commits after the compare-and-set but before projection must retain deleted-task semantics instead of returning a stale successful resolution.
*/
import { afterAll, afterEach, beforeAll, beforeEach, expect, it, vi } from "vitest";
import {
createSharedPgTaskStoreTestHarness,
pgDescribe,
type SharedPgTaskStoreHarness,
} from "../../__test-utils__/pg-test-harness.js";
const pgTest = pgDescribe;
pgTest("TaskStore wedge episode resolution (PostgreSQL)", () => {
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({
prefix: "fusion_wedge_resolution",
poolMax: 3,
});
beforeAll(h.beforeAll);
beforeEach(h.beforeEach);
afterEach(h.afterEach);
afterEach(() => vi.restoreAllMocks());
afterAll(h.afterAll);
it("resolves the exact active episode", async () => {
const store = h.store();
const task = await h.createTestTask();
await store.updateTask(task.id, {
wedgeNotification: {
reasonKey: "failed:stale",
episodeId: "episode-observed",
status: "active",
transitionedAt: "2026-07-29T00:00:00.000Z",
},
});
const result = await store.resolveTaskWedgeNotificationEpisode(task.id, "episode-observed");
expect(result.resolved).toBe(true);
expect(result.task.wedgeNotification).toMatchObject({
episodeId: "episode-observed",
status: "resolved",
});
});
it("reports a committed resolution when the derived task JSON projection fails", async () => {
const store = h.store();
const task = await h.createTestTask();
await store.updateTask(task.id, {
wedgeNotification: {
reasonKey: "failed:stale",
episodeId: "episode-projection-failure",
status: "active",
transitionedAt: "2026-07-29T00:00:00.000Z",
},
});
const writeTaskJsonFile = vi
.spyOn(store, "writeTaskJsonFile")
.mockRejectedValueOnce(new Error("projection unavailable"));
const result = await store.resolveTaskWedgeNotificationEpisode(task.id, "episode-projection-failure");
expect(writeTaskJsonFile).toHaveBeenCalledTimes(1);
expect(result.resolved).toBe(true);
expect(result.task.wedgeNotification).toMatchObject({
episodeId: "episode-projection-failure",
status: "resolved",
});
expect((await store.getTask(task.id)).wedgeNotification).toMatchObject({
episodeId: "episode-projection-failure",
status: "resolved",
});
});
it("maps deletion between resolution and projection to a deleted-task failure", async () => {
const store = h.store();
const layer = h.layer();
const task = await h.createTestTask();
const projectId = layer.projectId ?? "__legacy_unscoped__";
await store.updateTask(task.id, {
wedgeNotification: {
reasonKey: "failed:stale",
episodeId: "episode-deleted-before-projection",
status: "active",
transitionedAt: "2026-07-29T00:00:00.000Z",
},
});
const originalTransactionImmediate = layer.transactionImmediate.bind(layer);
let reportProjectionPending!: () => void;
const projectionPending = new Promise<void>((resolve) => {
reportProjectionPending = resolve;
});
let releaseProjection!: () => void;
const projectionMayStart = new Promise<void>((resolve) => {
releaseProjection = resolve;
});
vi.spyOn(layer, "transactionImmediate").mockImplementationOnce(async (callback) => {
reportProjectionPending();
await projectionMayStart;
return originalTransactionImmediate(callback);
});
const resolution = store.resolveTaskWedgeNotificationEpisode(task.id, "episode-deleted-before-projection");
await projectionPending;
const deletedAt = "2026-07-29T00:02:00.000Z";
await h.adminSql()`
UPDATE project.tasks
SET deleted_at = ${deletedAt}, updated_at = ${deletedAt}
WHERE project_id = ${projectId} AND id = ${task.id}
`;
releaseProjection();
await expect(resolution).rejects.toMatchObject({
name: "TaskDeletedError",
deletedAt,
});
});
it("does not let resolved publication hide a replacement episode", async () => {
const store = h.store();
const task = await h.createTestTask();
const projectId = h.layer().projectId ?? "__legacy_unscoped__";
await store.updateTask(task.id, {
wedgeNotification: {
reasonKey: "failed:stale",
episodeId: "episode-publication-observed",
status: "active",
transitionedAt: "2026-07-29T00:00:00.000Z",
},
});
const publishedEpisodes: string[] = [];
store.on("task:updated", (updated) => {
if (updated.wedgeNotification?.episodeId) publishedEpisodes.push(updated.wedgeNotification.episodeId);
});
const originalWriteTaskJsonFile = store.writeTaskJsonFile.bind(store);
let reportResolutionProjection!: () => void;
const resolutionProjectionStarted = new Promise<void>((resolve) => {
reportResolutionProjection = resolve;
});
let releaseResolutionProjection!: () => void;
const resolutionProjectionMayFinish = new Promise<void>((resolve) => {
releaseResolutionProjection = resolve;
});
vi.spyOn(store, "writeTaskJsonFile").mockImplementationOnce(async (...args) => {
reportResolutionProjection();
await resolutionProjectionMayFinish;
await originalWriteTaskJsonFile(...args);
});
const resolution = store.resolveTaskWedgeNotificationEpisode(task.id, "episode-publication-observed");
await resolutionProjectionStarted;
let reportReplacementAttempt!: () => void;
const replacementAttempted = new Promise<void>((resolve) => {
reportReplacementAttempt = resolve;
});
const replacement = (async () => {
reportReplacementAttempt();
await h.adminSql()`
UPDATE project.tasks
SET wedge_notification = ${JSON.stringify({
reasonKey: "failed:new",
episodeId: "episode-publication-replacement",
status: "active",
transitionedAt: "2026-07-29T00:01:00.000Z",
})}, updated_at = ${"2026-07-29T00:01:00.000Z"}
WHERE project_id = ${projectId} AND id = ${task.id}
`;
const replacementTask = await store.getTask(task.id);
await originalWriteTaskJsonFile(store.taskDir(task.id), replacementTask);
store.emitTaskLifecycleEventSafely("task:updated", [replacementTask]);
return replacementTask;
})();
await replacementAttempted;
await new Promise<void>((resolve) => setImmediate(resolve));
releaseResolutionProjection();
const [resolutionResult, replacementTask] = await Promise.all([resolution, replacement]);
const projectedTask = await store.readTaskJson(store.taskDir(task.id));
expect(resolutionResult.resolved).toBe(true);
expect(replacementTask.wedgeNotification).toMatchObject({
episodeId: "episode-publication-replacement",
status: "active",
});
expect(projectedTask.wedgeNotification).toMatchObject({
episodeId: "episode-publication-replacement",
status: "active",
});
expect(publishedEpisodes.at(-1)).toBe("episode-publication-replacement");
});
it("preserves a replacement episode committed while the resolver waits on the row lock", async () => {
const store = h.store();
const task = await h.createTestTask();
const projectId = h.layer().projectId ?? "__legacy_unscoped__";
await store.updateTask(task.id, {
wedgeNotification: {
reasonKey: "failed:stale",
episodeId: "episode-observed",
status: "active",
transitionedAt: "2026-07-29T00:00:00.000Z",
},
});
let releaseReplacement!: () => void;
const replacementMayCommit = new Promise<void>((resolve) => {
releaseReplacement = resolve;
});
let reportLocked!: () => void;
const rowLocked = new Promise<void>((resolve) => {
reportLocked = resolve;
});
const replacement = h.adminSql().begin(async (sql) => {
await sql`
SELECT id
FROM project.tasks
WHERE project_id = ${projectId} AND id = ${task.id}
FOR UPDATE
`;
reportLocked();
await replacementMayCommit;
await sql`
UPDATE project.tasks
SET wedge_notification = ${JSON.stringify({
reasonKey: "failed:new",
episodeId: "episode-replacement",
status: "active",
transitionedAt: "2026-07-29T00:01:00.000Z",
})}, updated_at = ${"2026-07-29T00:01:00.000Z"}
WHERE project_id = ${projectId} AND id = ${task.id}
`;
});
await rowLocked;
const resolution = store.resolveTaskWedgeNotificationEpisode(task.id, "episode-observed");
await new Promise<void>((resolve) => setImmediate(resolve));
releaseReplacement();
const [, result] = await Promise.all([replacement, resolution]);
expect(result.resolved).toBe(false);
expect((await store.getTask(task.id)).wedgeNotification).toMatchObject({
episodeId: "episode-replacement",
status: "active",
});
});
it.each([
{ label: "without audit", runContext: undefined, persistMethod: "atomicWriteTaskJson" as const },
{ label: "with audit", runContext: { agentId: "wedge-test", runId: "stale-wedge-write" }, persistMethod: "atomicWriteTaskJsonWithAudit" as const },
])("does not let an ordinary $label task write reactivate a resolved episode", async ({ runContext, persistMethod }) => {
const store = h.store();
const task = await h.createTestTask();
await store.updateTask(task.id, {
wedgeNotification: {
reasonKey: "failed:stale",
episodeId: "episode-resolved-before-stale-write",
status: "active",
transitionedAt: "2026-07-29T00:00:00.000Z",
},
});
const originalPersist = store[persistMethod].bind(store) as (...args: unknown[]) => Promise<void>;
let reportStaleSnapshot!: () => void;
const staleSnapshotReady = new Promise<void>((resolve) => {
reportStaleSnapshot = resolve;
});
let releaseStaleWrite!: () => void;
const staleWriteMayPersist = new Promise<void>((resolve) => {
releaseStaleWrite = resolve;
});
const persistSpy = vi.spyOn(store, persistMethod) as unknown as {
mockImplementationOnce: (implementation: (...args: unknown[]) => Promise<void>) => void;
};
persistSpy.mockImplementationOnce(async (...args: unknown[]) => {
reportStaleSnapshot();
await staleWriteMayPersist;
await originalPersist(...args);
});
const staleWrite = store.updateTask(task.id, { title: "ordinary concurrent update" }, runContext);
await staleSnapshotReady;
const resolution = await store.resolveTaskWedgeNotificationEpisode(
task.id,
"episode-resolved-before-stale-write",
);
releaseStaleWrite();
const updated = await staleWrite;
expect(resolution.resolved).toBe(true);
expect(updated.title).toBe("ordinary concurrent update");
expect(updated.wedgeNotification).toMatchObject({
episodeId: "episode-resolved-before-stale-write",
status: "resolved",
});
expect((await store.getTask(task.id)).wedgeNotification).toMatchObject({
episodeId: "episode-resolved-before-stale-write",
status: "resolved",
});
});
});

View File

@@ -104,7 +104,7 @@ import { getTaskCommitAssociationsByLineageIdImpl, replaceLegacyTaskCommitAssoci
import { findRecentTasksBySourceParentTaskIdImpl } from "./task-store/branch-and-pr-entities.js"; import { findRecentTasksBySourceParentTaskIdImpl } from "./task-store/branch-and-pr-entities.js";
import { addTaskCommentImpl, applyBuiltInPromptOverridesAsyncImpl, applyBuiltInPromptOverridesSyncImpl, areAllDependenciesDoneImpl, artifactStoredNameImpl, assertWorkflowIrTraitsValidImpl, clearActivityLogImpl, clearTaskWorkflowSelectionImpl, deleteTaskByIdImpl, getDefaultWorkflowIdImpl, resolveOriginWorkflowOverrideIdImpl, type TaskOriginWorkflowKind, getInsightStoreImpl, getMergeQueuedTaskIdsImpl, getMergeRequestRecordImpl, getMergeRequestRecordAsyncImpl, getResearchStoreImpl, getTaskIdFromDirImpl, getTodoStoreImpl, getWorkflowWorkItemByIdentityImpl, hasActiveTaskImpl, invalidateConfigCacheAfterMigrationImpl, isTaskIdConflictErrorImpl, listLegacyAutoMergeStampCandidatesImpl, readTaskRowFromDbImpl, recordBranchGroupMemberLandedImpl, refreshDatabaseHealthAsyncImpl, refreshDatabaseHealthImpl, resolveEffectiveWorkflowIdSyncImpl, resolveTaskCustomFieldDefsSyncImpl, resolveWorkflowBypassGuardsImpl, serializeConfigForDiskImpl, setPluginWorkflowStepTemplatesImpl, shouldSkipWorkflowMovePoliciesImpl, suppressWatcherImpl, upsertTaskWithFtsRecoveryImpl } from "./task-store/task-store-helpers.js"; import { addTaskCommentImpl, applyBuiltInPromptOverridesAsyncImpl, applyBuiltInPromptOverridesSyncImpl, areAllDependenciesDoneImpl, artifactStoredNameImpl, assertWorkflowIrTraitsValidImpl, clearActivityLogImpl, clearTaskWorkflowSelectionImpl, deleteTaskByIdImpl, getDefaultWorkflowIdImpl, resolveOriginWorkflowOverrideIdImpl, type TaskOriginWorkflowKind, getInsightStoreImpl, getMergeQueuedTaskIdsImpl, getMergeRequestRecordImpl, getMergeRequestRecordAsyncImpl, getResearchStoreImpl, getTaskIdFromDirImpl, getTodoStoreImpl, getWorkflowWorkItemByIdentityImpl, hasActiveTaskImpl, invalidateConfigCacheAfterMigrationImpl, isTaskIdConflictErrorImpl, listLegacyAutoMergeStampCandidatesImpl, readTaskRowFromDbImpl, recordBranchGroupMemberLandedImpl, refreshDatabaseHealthAsyncImpl, refreshDatabaseHealthImpl, resolveEffectiveWorkflowIdSyncImpl, resolveTaskCustomFieldDefsSyncImpl, resolveWorkflowBypassGuardsImpl, serializeConfigForDiskImpl, setPluginWorkflowStepTemplatesImpl, shouldSkipWorkflowMovePoliciesImpl, suppressWatcherImpl, upsertTaskWithFtsRecoveryImpl } from "./task-store/task-store-helpers.js";
import { getTaskSelectClauseImpl2, createTaskPersistSerializationContextImpl, getTaskPersistValuesImpl, getTaskPatchDescriptorsImpl, normalizeTaskFromDiskImpl, writeTaskJsonFileImpl, rowToPrEntityImpl, generatePrEntityIdImpl, readTaskForMoveImpl, rowToMergeQueueEntryImpl, rowToMergeRequestRecordImpl, rowToCompletionHandoffMarkerImpl, rowToWorkflowWorkItemImpl, rowToRunAuditEventImpl } from "./task-store/task-row-mappers.js"; import { getTaskSelectClauseImpl2, createTaskPersistSerializationContextImpl, getTaskPersistValuesImpl, getTaskPatchDescriptorsImpl, normalizeTaskFromDiskImpl, writeTaskJsonFileImpl, rowToPrEntityImpl, generatePrEntityIdImpl, readTaskForMoveImpl, rowToMergeQueueEntryImpl, rowToMergeRequestRecordImpl, rowToCompletionHandoffMarkerImpl, rowToWorkflowWorkItemImpl, rowToRunAuditEventImpl } from "./task-store/task-row-mappers.js";
import { getTaskSelectClauseWithActivityLogLimitImpl, getChangedTaskColumnsImpl, getSoftDeletedWriteConflictImpl, readTaskJsonImpl, writeConfigImpl, _maybeAutoArchiveSameAgentDuplicateBackendImpl, updateBranchGroupImpl, updatePrEntityImpl, listTasksForGithubTrackingReconcileImpl, listTasksForGitlabTrackingReconcileImpl, renewCheckoutLeaseImpl, updateTaskAtomicImpl, getWorkflowPromptOverridesImpl, updateWorkflowSettingValuesImpl, rollbackConfigurationImpl, cancelActiveWorkflowWorkItemsForTaskImpl, setCompletionHandoffAcceptedMarkerImpl, reconcileLegacyAutoMergeStampsImpl, recoverExpiredMergeQueueLeasesImpl, rewriteDependentsForRemovalImpl, cleanupBranchForTaskImpl, addAttachmentImpl, deleteAttachmentImpl, registerArtifactImpl, updatePrInfoImpl, unlinkGithubIssueImpl, cleanupArchivedTasksImpl, generatePromptFromArchiveEntryImpl, listWorkflowOccupantTaskIdsImpl, listApprovedCliAutonomyAdaptersImpl, closeImpl, getActivityLogImpl } from "./task-store/task-mutation-ops.js"; import { getTaskSelectClauseWithActivityLogLimitImpl, getChangedTaskColumnsImpl, getSoftDeletedWriteConflictImpl, readTaskJsonImpl, writeConfigImpl, _maybeAutoArchiveSameAgentDuplicateBackendImpl, updateBranchGroupImpl, updatePrEntityImpl, listTasksForGithubTrackingReconcileImpl, listTasksForGitlabTrackingReconcileImpl, renewCheckoutLeaseImpl, updateTaskAtomicImpl, resolveTaskWedgeNotificationEpisodeImpl, getWorkflowPromptOverridesImpl, updateWorkflowSettingValuesImpl, rollbackConfigurationImpl, cancelActiveWorkflowWorkItemsForTaskImpl, setCompletionHandoffAcceptedMarkerImpl, reconcileLegacyAutoMergeStampsImpl, recoverExpiredMergeQueueLeasesImpl, rewriteDependentsForRemovalImpl, cleanupBranchForTaskImpl, addAttachmentImpl, deleteAttachmentImpl, registerArtifactImpl, updatePrInfoImpl, unlinkGithubIssueImpl, cleanupArchivedTasksImpl, generatePromptFromArchiveEntryImpl, listWorkflowOccupantTaskIdsImpl, listApprovedCliAutonomyAdaptersImpl, closeImpl, getActivityLogImpl } from "./task-store/task-mutation-ops.js";
import { getOrCreateForProjectImpl, listGoalCitationsImpl, atomicWriteTaskJsonWithAuditImpl, duplicateTaskImpl, listStrandedRefinementsImpl, tryClaimCheckoutImpl, evaluateWorkflowMovePoliciesImpl, recordRunAuditEventImpl, getRunAuditEventsImpl, dequeueMergeQueueOnColumnExitImpl, updateIssueInfoImpl, listWorkflowStepsImpl, getWorkflowStepImpl, createWorkflowDefinitionImpl, countActiveInCapacitySlotSyncImpl, countActiveInCapacitySlotAsyncImpl, generateSpecifiedPromptImpl, recordActivityImpl, getEvalStoreImpl } from "./task-store/project-store-ops.js"; import { getOrCreateForProjectImpl, listGoalCitationsImpl, atomicWriteTaskJsonWithAuditImpl, duplicateTaskImpl, listStrandedRefinementsImpl, tryClaimCheckoutImpl, evaluateWorkflowMovePoliciesImpl, recordRunAuditEventImpl, getRunAuditEventsImpl, dequeueMergeQueueOnColumnExitImpl, updateIssueInfoImpl, listWorkflowStepsImpl, getWorkflowStepImpl, createWorkflowDefinitionImpl, countActiveInCapacitySlotSyncImpl, countActiveInCapacitySlotAsyncImpl, generateSpecifiedPromptImpl, recordActivityImpl, getEvalStoreImpl } from "./task-store/project-store-ops.js";
import { markLegacyAutoMergeStampsOnceImpl, appendAgentLogImpl, importLegacyAgentLogsImpl, cleanupNoOpTaskMovedActivityRowsOnceImpl, backfillCommitAssociationDiffStatsImpl } from "./task-store/workflow-integrity.js"; import { markLegacyAutoMergeStampsOnceImpl, appendAgentLogImpl, importLegacyAgentLogsImpl, cleanupNoOpTaskMovedActivityRowsOnceImpl, backfillCommitAssociationDiffStatsImpl } from "./task-store/workflow-integrity.js";
import { saveWorkflowRunBranchImpl, clearNearDuplicateReferencesToImpl, selectNextTaskForAgentImpl, pauseTaskImpl, clearLinkedAgentTaskIdsImpl, listArtifactsImpl, rehomeOccupantImpl, type RehomeOccupantResult } from "./task-store/branch-group-ops.js"; import { saveWorkflowRunBranchImpl, clearNearDuplicateReferencesToImpl, selectNextTaskForAgentImpl, pauseTaskImpl, clearLinkedAgentTaskIdsImpl, listArtifactsImpl, rehomeOccupantImpl, type RehomeOccupantResult } from "./task-store/branch-group-ops.js";
@@ -1286,6 +1286,9 @@ export class TaskStore extends EventEmitter<TaskStoreEvents> {
async updateTaskAtomic( id: string, updater: ( current: Task, ) => Parameters<TaskStore["updateTask"]>[1] | null | undefined | Promise<Parameters<TaskStore["updateTask"]>[1] | null | undefined>, runContext?: RunMutationContext, ): Promise<Task> { async updateTaskAtomic( id: string, updater: ( current: Task, ) => Parameters<TaskStore["updateTask"]>[1] | null | undefined | Promise<Parameters<TaskStore["updateTask"]>[1] | null | undefined>, runContext?: RunMutationContext, ): Promise<Task> {
return updateTaskAtomicImpl(this, id, updater, runContext); return updateTaskAtomicImpl(this, id, updater, runContext);
} }
async resolveTaskWedgeNotificationEpisode(id: string, episodeId: string): Promise<{ task: Task; resolved: boolean }> {
return resolveTaskWedgeNotificationEpisodeImpl(this, id, episodeId);
}
public mergeCustomFieldPatch( current: Record<string, unknown> | undefined, patch: Record<string, unknown>, ): Record<string, unknown> { public mergeCustomFieldPatch( current: Record<string, unknown> | undefined, patch: Record<string, unknown>, ): Record<string, unknown> {
return mergeCustomFieldPatchImpl(this, current, patch); return mergeCustomFieldPatchImpl(this, current, patch);
} }

View File

@@ -198,6 +198,33 @@ export async function softDeleteTaskRow(
)); ));
} }
/**
* FNXC:TaskStateReconciliation 2026-07-29-16:10:
* Resolve only the active wedge episode the caller observed. The PostgreSQL predicate is the cross-process compare-and-set authority, so an update waiting behind a replacement episode rechecks the durable row and cannot clear the replacement.
*/
export async function resolveActiveTaskWedgeEpisodeRow(
layer: AsyncDataLayer,
id: string,
episodeId: string,
transitionedAt: string,
): Promise<Record<string, unknown> | undefined> {
const rows = await layer.db
.update(schema.project.tasks)
.set({
wedgeNotification: sql`(${schema.project.tasks.wedgeNotification}::jsonb || jsonb_build_object('status', 'resolved', 'transitionedAt', ${transitionedAt}))::text`,
updatedAt: transitionedAt,
})
.where(and(
taskProjectScope(layer),
eq(schema.project.tasks.id, id),
isNull(schema.project.tasks.deletedAt),
sql`${schema.project.tasks.wedgeNotification}::jsonb ->> 'episodeId' = ${episodeId}`,
sql`${schema.project.tasks.wedgeNotification}::jsonb ->> 'status' = 'active'`,
))
.returning();
return rows[0] as Record<string, unknown> | undefined;
}
/** /**
* FNXC:TaskStoreArchiveLineage 2026-06-24-15:00: * FNXC:TaskStoreArchiveLineage 2026-06-24-15:00:
* Soft-delete a task INSIDE a shared transaction handle. This is the * Soft-delete a task INSIDE a shared transaction handle. This is the

View File

@@ -9,7 +9,7 @@
*/ */
import type { Task } from "../types.js"; import type { Task } from "../types.js";
import { normalizeTaskPriority } from "../task-priority.js"; import { normalizeTaskPriority } from "../task-priority.js";
import { toJson, toJsonNullable } from "../db.js"; import { fromJson, toJson, toJsonNullable } from "../db.js";
/** Database row shape for the tasks table (all columns). */ /** Database row shape for the tasks table (all columns). */
export interface TaskRow { export interface TaskRow {
@@ -183,6 +183,22 @@ export type TaskColumnDescriptor = {
serialize: (task: Task, context: TaskPersistSerializationContext) => unknown; serialize: (task: Task, context: TaskPersistSerializationContext) => unknown;
}; };
/**
* FNXC:TaskStateReconciliation 2026-07-29-20:53:
* A generic PostgreSQL task write may carry an active wedge snapshot that was read before the live API resolved that episode. Preserve the durable resolution for the same episode so changed-column persistence, task.json projection, and cache publication cannot reactivate an acknowledged operator notification; a genuinely new wedge must use a new episode ID.
*/
export function preserveResolvedTaskWedgeEpisode(existingRow: Pick<TaskRow, "wedgeNotification">, task: Task): void {
const durable = fromJson<Task["wedgeNotification"]>(existingRow.wedgeNotification);
const incoming = task.wedgeNotification;
if (
durable?.status === "resolved"
&& incoming?.status === "active"
&& durable.episodeId === incoming.episodeId
) {
task.wedgeNotification = durable;
}
}
/* /*
FNXC:TaskLifecyclePersistence 2026-07-14-13:27: FNXC:TaskLifecyclePersistence 2026-07-14-13:27:
PostgreSQL task JSONB conversion must use one registry for both descriptor writes and SQLite-shaped row hydration. Separate read/write lists drifted when late lifecycle columns were added, allowing JSON strings or parsed objects to cross the wrong serialization boundary. PostgreSQL task JSONB conversion must use one registry for both descriptor writes and SQLite-shaped row hydration. Separate read/write lists drifted when late lifecycle columns were added, allowing JSON strings or parsed objects to cross the wrong serialization boundary.

View File

@@ -32,7 +32,7 @@ import {CentralCore} from "../central-core.js";
import {extractTaskIdTokens, normalizeTitleForTaskId} from "../task-title-id-drift.js"; import {extractTaskIdTokens, normalizeTitleForTaskId} from "../task-title-id-drift.js";
import {generateTaskLineageId} from "../task-lineage.js"; import {generateTaskLineageId} from "../task-lineage.js";
import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.js"; import {sanitizeFileScopeInPromptContent} from "../task-store/file-scope.js";
import {type TaskRow} from "../task-store/persistence.js"; import {preserveResolvedTaskWedgeEpisode, type TaskRow} from "../task-store/persistence.js";
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js"; import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
import {isWorkflowDefinitionIdPrimaryKeyCollision, nextWorkflowDefinitionIdAsyncImpl} from "../task-store/workflow-definitions.js"; import {isWorkflowDefinitionIdPrimaryKeyCollision, nextWorkflowDefinitionIdAsyncImpl} from "../task-store/workflow-definitions.js";
import {upsertTaskRowInTransaction, buildTaskInsertValues} from "../task-store/async-persistence.js"; import {upsertTaskRowInTransaction, buildTaskInsertValues} from "../task-store/async-persistence.js";
@@ -136,6 +136,7 @@ export async function atomicWriteTaskJsonWithAuditImpl(store: TaskStore, dir: st
*/ */
if (row) { if (row) {
const existing = store.pgRowToTaskRow(row); const existing = store.pgRowToTaskRow(row);
preserveResolvedTaskWedgeEpisode(existing, task);
const changedColumns = store.getChangedTaskColumns(existing, task); const changedColumns = store.getChangedTaskColumns(existing, task);
if (changedColumns.size > 0) { if (changedColumns.size > 0) {
const context = store.createTaskPersistSerializationContext(task, existing); const context = store.createTaskPersistSerializationContext(task, existing);

View File

@@ -12,7 +12,7 @@ const severityAuditLog = createLogger("core-task-mutation-ops");
* instance as its first parameter and performs byte-identical work. * instance as its first parameter and performs byte-identical work.
*/ */
import {TaskStore, storeLog} from "../store.js"; import {TaskStore, storeLog} from "../store.js";
import {TaskDeletedError} from "./errors.js"; import {TaskDeletedError, TaskNotFoundError} from "./errors.js";
import type {LegacyAutoMergeStampReconcileResult} from "../store.js"; import type {LegacyAutoMergeStampReconcileResult} from "../store.js";
import {randomUUID} from "node:crypto"; import {randomUUID} from "node:crypto";
import {mkdir, readFile, writeFile, rename, unlink} from "node:fs/promises"; import {mkdir, readFile, writeFile, rename, unlink} from "node:fs/promises";
@@ -26,7 +26,7 @@ import {resolveSameAgentDuplicateIntake} from "./task-creation.js";
import {type TaskRow, TASK_COLUMN_DESCRIPTORS} from "../task-store/persistence.js"; import {type TaskRow, TASK_COLUMN_DESCRIPTORS} from "../task-store/persistence.js";
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js"; import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
import {assertSafeGitBranchName} from "../task-store/shell-safety.js"; import {assertSafeGitBranchName} from "../task-store/shell-safety.js";
import {readTaskRow as readTaskRowAsync, readTaskRowInTransaction} from "../task-store/async-persistence.js"; import {readTaskRow as readTaskRowAsync, readTaskRowInTransaction, resolveActiveTaskWedgeEpisodeRow} from "../task-store/async-persistence.js";
import {upsertArchivedTaskEntry} from "./async-archive-lineage.js"; import {upsertArchivedTaskEntry} from "./async-archive-lineage.js";
import {purgeTaskWorkflowSelectionRowsAsyncImpl} from "./workflow-definitions.js"; import {purgeTaskWorkflowSelectionRowsAsyncImpl} from "./workflow-definitions.js";
import * as schema from "../postgres/schema/index.js"; import * as schema from "../postgres/schema/index.js";
@@ -372,6 +372,61 @@ export async function updateTaskAtomicImpl(store: TaskStore, id: string, updater
}); });
} }
export async function resolveTaskWedgeNotificationEpisodeImpl(
store: TaskStore,
id: string,
episodeId: string,
): Promise<{ task: Task; resolved: boolean }> {
const layer = store.asyncLayer!;
const row = await resolveActiveTaskWedgeEpisodeRow(layer, id, episodeId, new Date().toISOString());
if (row) {
const task = await layer.transactionImmediate(async (tx) => {
const rows = await tx
.select()
.from(schema.project.tasks)
.where(and(
eq(schema.project.tasks.projectId, layer.projectId?.trim() || "__legacy_unscoped__"),
eq(schema.project.tasks.id, id),
))
.for("update");
const currentRow = rows[0];
/*
FNXC:TaskStateReconciliation 2026-07-29-22:17:
The live API must not return a stale successful resolution when deletion wins between the compare-and-set and projection lock. Reclassify the authoritative locked row through the same not-found/deleted contract used by the non-resolved path.
*/
if (!currentRow) throw new TaskNotFoundError(id);
if (currentRow.deletedAt) throw new TaskDeletedError(id, currentRow.deletedAt as string);
const currentTask = store.rowToTask(store.pgRowToTaskRow(currentRow));
/*
FNXC:TaskStateReconciliation 2026-07-29-22:01:
Derived task JSON, cache, and lifecycle publication must remain ordered with PostgreSQL wedge mutations across processes. Lock and re-read the durable row before publication so a replacement episode either publishes first and is selected here, or waits and publishes after this resolved episode.
FNXC:TaskStateReconciliation 2026-07-29-17:43:
PostgreSQL commits wedge resolution before the task JSON projection runs. A projection failure must not report the committed mutation as failed; keep cache and lifecycle observers current while startup reconciliation repairs the derived file.
*/
await store.writeTaskJsonFile(store.taskDir(id), currentTask).catch((error) => {
severityAuditLog.warn("Failed to project committed wedge resolution to task JSON", {
taskId: id,
error: error instanceof Error ? error.message : String(error),
});
});
if (store.isWatching) store.taskCache.set(id, { ...currentTask });
store.emitTaskLifecycleEventSafely("task:updated", [currentTask]);
return currentTask;
});
return { task, resolved: true };
}
const currentRow = await readTaskRowAsync(layer, id, { includeDeleted: true });
if (!currentRow) throw new TaskNotFoundError(id);
if (currentRow.deletedAt) throw new TaskDeletedError(id, currentRow.deletedAt as string);
return {
task: store.rowToTask(store.pgRowToTaskRow(currentRow)),
resolved: false,
};
}
export function getWorkflowPromptOverridesImpl(_store: TaskStore, _workflowId: string, _projectId: string): Record<string, string> { export function getWorkflowPromptOverridesImpl(_store: TaskStore, _workflowId: string, _projectId: string): Record<string, string> {
/* /*
* FNXC:SqliteFinalRemoval 2026-06-26: * FNXC:SqliteFinalRemoval 2026-06-26:

View File

@@ -32,6 +32,7 @@ import {AsyncGoalStore} from "../async-goal-store.js";
import {normalizeTaskCommitAssociation} from "../task-lineage.js"; import {normalizeTaskCommitAssociation} from "../task-lineage.js";
import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js"; import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js";
import {withTaskBranchContextInSourceMetadata} from "../task-store/branch-context.js"; import {withTaskBranchContextInSourceMetadata} from "../task-store/branch-context.js";
import {preserveResolvedTaskWedgeEpisode} from "../task-store/persistence.js";
import {upsertTaskRowInTransaction, readTaskRowInTransaction, buildTaskInsertValues} from "../task-store/async-persistence.js"; import {upsertTaskRowInTransaction, readTaskRowInTransaction, buildTaskInsertValues} from "../task-store/async-persistence.js";
import {listDueWorkflowWorkItems as listDueWorkflowWorkItemsAsync, withTaskWorkflowSerialization} from "../task-store/async-workflow-workitems.js"; import {listDueWorkflowWorkItems as listDueWorkflowWorkItemsAsync, withTaskWorkflowSerialization} from "../task-store/async-workflow-workitems.js";
import {getTaskMovedCountsByDay as getTaskMovedCountsByDayAsync} from "../task-store/async-audit.js"; import {getTaskMovedCountsByDay as getTaskMovedCountsByDayAsync} from "../task-store/async-audit.js";
@@ -102,6 +103,7 @@ export async function atomicWriteTaskJsonImpl2(store: TaskStore, dir: string, ta
return; return;
} }
const existingRow = store.pgRowToTaskRow(pgRow); const existingRow = store.pgRowToTaskRow(pgRow);
preserveResolvedTaskWedgeEpisode(existingRow, task);
const deletedAt = store.getSoftDeletedWriteConflict(id, task, existingRow); const deletedAt = store.getSoftDeletedWriteConflict(id, task, existingRow);
if (deletedAt) { if (deletedAt) {
store.throwSoftDeletedWriteBlocked(id, deletedAt, "atomicWriteTaskJson"); store.throwSoftDeletedWriteBlocked(id, deletedAt, "atomicWriteTaskJson");

View File

@@ -0,0 +1,201 @@
// @vitest-environment node
import { describe, expect, it, vi } from "vitest";
import express from "express";
import { TaskDeletedError, TaskNotFoundError, type TaskStore } from "@fusion/core";
import { createApiRoutes } from "../../routes.js";
import { request as REQUEST } from "../../test-request.js";
interface HarnessOptions {
missing?: boolean;
deleted?: boolean;
rejectStepTransition?: boolean;
wedgeEpisodeId?: string;
}
function createHarness(options: HarnessOptions = {}) {
const task = {
id: "FN-001",
description: "legacy checklist task",
column: "in-review",
dependencies: [],
steps: [
{ title: "Implement", status: "pending" },
{ title: "Verify", status: "pending" },
],
currentStep: 0,
log: [],
wedgeNotification: {
reasonKey: "failed:stale",
episodeId: options.wedgeEpisodeId ?? "episode-observed",
status: "active",
transitionedAt: "2026-07-29T00:00:00.000Z",
},
createdAt: "2026-07-29T00:00:00.000Z",
updatedAt: "2026-07-29T00:00:00.000Z",
} as any;
const getTask = vi.fn(async () => {
if (options.missing) throw new TaskNotFoundError("FN-001");
return task;
});
const updateStep = vi.fn(async (_id: string, index: number, status: string) => {
if (options.deleted) throw new TaskDeletedError("FN-001", "2026-07-29T00:00:00.000Z");
if (!options.rejectStepTransition) task.steps[index].status = status;
return task;
});
const resolveTaskWedgeNotificationEpisode = vi.fn(async (_id: string, episodeId: string) => {
if (options.missing) throw new TaskNotFoundError("FN-001");
if (options.deleted) throw new TaskDeletedError("FN-001", "2026-07-29T00:00:00.000Z");
if (task.wedgeNotification.episodeId !== episodeId || task.wedgeNotification.status !== "active") {
return { task, resolved: false };
}
task.wedgeNotification = {
...task.wedgeNotification,
status: "resolved",
transitionedAt: "2026-07-29T00:01:00.000Z",
};
return { task, resolved: true };
});
const store = {
getRootDir: vi.fn(() => process.cwd()),
getProjectScopedPluginMcpServers: vi.fn(async () => []),
getTask,
updateStep,
resolveTaskWedgeNotificationEpisode,
} as unknown as TaskStore;
const app = express();
app.use(express.json());
app.use("/api", createApiRoutes(store));
return { app, getTask, updateStep, resolveTaskWedgeNotificationEpisode, task };
}
describe("task checklist step update route", () => {
it("updates one step through the live project store", async () => {
const { app, updateStep } = createHarness();
const response = await REQUEST(
app,
"PATCH",
"/api/tasks/FN-001/steps/1",
JSON.stringify({ status: "done" }),
{ "content-type": "application/json" },
);
expect(response.status).toBe(200);
expect(updateStep).toHaveBeenCalledWith("FN-001", 1, "done");
expect((response.body as { steps: Array<{ status: string }> }).steps[1].status).toBe("done");
});
it.each([
["not-an-index", { status: "done" }],
["-1", { status: "done" }],
["2", { status: "done" }],
["0", { status: "invalid" }],
])("rejects invalid step update %s", async (index, body) => {
const { app, updateStep } = createHarness();
const response = await REQUEST(
app,
"PATCH",
`/api/tasks/FN-001/steps/${index}`,
JSON.stringify(body),
{ "content-type": "application/json" },
);
expect(response.status).toBe(400);
expect(updateStep).not.toHaveBeenCalled();
});
it("returns 409 when the store rejects the requested step transition", async () => {
const { app } = createHarness({ rejectStepTransition: true });
const response = await REQUEST(
app,
"PATCH",
"/api/tasks/FN-001/steps/1",
JSON.stringify({ status: "done" }),
{ "content-type": "application/json" },
);
expect(response.status).toBe(409);
expect(response.body).toEqual(expect.objectContaining({ error: expect.stringContaining("was rejected") }));
});
it("resolves only the observed stale task wedge episode", async () => {
const { app, resolveTaskWedgeNotificationEpisode } = createHarness();
const response = await REQUEST(
app,
"POST",
"/api/tasks/FN-001/wedge/resolve",
JSON.stringify({ episodeId: "episode-observed" }),
{ "content-type": "application/json" },
);
expect(response.status).toBe(200);
expect(resolveTaskWedgeNotificationEpisode).toHaveBeenCalledWith("FN-001", "episode-observed");
expect((response.body as { wedgeNotification: { status: string } }).wedgeNotification.status).toBe("resolved");
});
it("does not clear a replacement wedge episode", async () => {
const { app, task } = createHarness({ wedgeEpisodeId: "episode-new" });
const response = await REQUEST(
app,
"POST",
"/api/tasks/FN-001/wedge/resolve",
JSON.stringify({ episodeId: "episode-observed" }),
{ "content-type": "application/json" },
);
expect(response.status).toBe(409);
expect(task.wedgeNotification).toEqual(expect.objectContaining({ episodeId: "episode-new", status: "active" }));
});
it("requires the observed wedge episode id", async () => {
const { app, resolveTaskWedgeNotificationEpisode } = createHarness();
const response = await REQUEST(
app,
"POST",
"/api/tasks/FN-001/wedge/resolve",
JSON.stringify({}),
{ "content-type": "application/json" },
);
expect(response.status).toBe(400);
expect(resolveTaskWedgeNotificationEpisode).not.toHaveBeenCalled();
});
/**
* FNXC:TaskStateReconciliation 2026-07-29-17:43:
* Both live mutation routes expose the same deleted-task 404 contract; testing only wedge resolution leaves checklist updates free to regress to a server error.
*/
it.each([
["PATCH", "/api/tasks/FN-001/steps/0", { status: "done" }],
["POST", "/api/tasks/FN-001/wedge/resolve", { episodeId: "episode-observed" }],
])("maps a soft-deleted task to 404 for %s %s", async (method, path, body) => {
const { app } = createHarness({ deleted: true });
const response = await REQUEST(
app,
method,
path,
JSON.stringify(body),
{ "content-type": "application/json" },
);
expect(response.status).toBe(404);
expect(response.body).toEqual(expect.objectContaining({ error: expect.stringContaining("Task FN-001 is soft-deleted") }));
});
it.each([
["PATCH", "/api/tasks/FN-001/steps/0", { status: "done" }],
["POST", "/api/tasks/FN-001/wedge/resolve", { episodeId: "episode-observed" }],
])("maps a missing task to 404 for %s %s", async (method, path, body) => {
const { app } = createHarness({ missing: true });
const response = await REQUEST(
app,
method,
path,
JSON.stringify(body),
{ "content-type": "application/json" },
);
expect(response.status).toBe(404);
expect(response.body).toEqual(expect.objectContaining({ error: "Task FN-001 not found" }));
});
});

View File

@@ -4890,6 +4890,62 @@ export function registerTaskWorkflowRoutes(ctx: ApiRoutesContext, deps: TaskWork
} }
}); });
/**
* FNXC:TaskStateReconciliation 2026-07-29-11:40:
* Checklist repair must use the live project-scoped store, map missing tasks to 404, reject out-of-range indices, and report 409 when lifecycle ordering rejects the requested transition instead of returning a false-success 200.
*/
router.patch("/tasks/:id/steps/:stepIndex", async (req, res) => {
try {
const { store: scopedStore } = await getProjectContext(req);
const stepIndex = Number(req.params.stepIndex);
const validStatuses = ["pending", "in-progress", "done", "skipped"] as const;
const status = req.body?.status;
if (!Number.isInteger(stepIndex) || stepIndex < 0) {
throw badRequest("stepIndex must be a non-negative integer");
}
if (!validStatuses.includes(status)) {
throw badRequest(`status must be one of: ${validStatuses.join(", ")}`);
}
const task = await scopedStore.getTask(req.params.id);
if (stepIndex >= (task.steps?.length ?? 0)) {
throw badRequest(`stepIndex ${stepIndex} is out of range`);
}
const updated = await scopedStore.updateStep(req.params.id, stepIndex, status);
if (updated.steps?.[stepIndex]?.status !== status) {
throw conflict(`Step ${stepIndex} transition to ${status} was rejected`);
}
res.json(updated);
} catch (err: unknown) {
rethrowTaskApiError(err, req.params.id);
}
});
/**
* FNXC:TaskStateReconciliation 2026-07-29-11:40:
* Wedge resolution is compare-and-set against the episode the operator observed. A concurrent replacement episode must remain active rather than being cleared by a stale request from another dashboard process.
*/
router.post("/tasks/:id/wedge/resolve", async (req, res) => {
try {
const { store: scopedStore } = await getProjectContext(req);
const { id } = req.params;
const episodeId = req.body?.episodeId;
if (typeof episodeId !== "string" || episodeId.length === 0) {
throw badRequest("episodeId must be a non-empty string");
}
const result = await scopedStore.resolveTaskWedgeNotificationEpisode(id, episodeId);
if (!result.resolved) {
throw conflict(`Wedge episode ${episodeId} is no longer active`);
}
res.json(result.task);
} catch (err: unknown) {
rethrowTaskApiError(err, req.params.id);
}
});
// Update task // Update task
router.patch("/tasks/:id", async (req, res) => { router.patch("/tasks/:id", async (req, res) => {
try { try {

View File

@@ -22,7 +22,7 @@
* signal so genuinely file-backed reads (attachments, session files, worktree * signal so genuinely file-backed reads (attachments, session files, worktree
* files) keep their 404 behavior. * files) keep their 404 behavior.
*/ */
import { isTaskNotFoundError } from "@fusion/core"; import { isTaskNotFoundError, TaskDeletedError } from "@fusion/core";
import { ApiError, notFound, rethrowAsApiError } from "../api-error.js"; import { ApiError, notFound, rethrowAsApiError } from "../api-error.js";
/** /**
@@ -32,7 +32,12 @@ import { ApiError, notFound, rethrowAsApiError } from "../api-error.js";
* file-backed reads on the same handlers. * file-backed reads on the same handlers.
*/ */
export function isTaskLookupMiss(error: unknown): boolean { export function isTaskLookupMiss(error: unknown): boolean {
if (isTaskNotFoundError(error)) return true; /*
* FNXC:TaskLookup404 2026-07-29-16:10:
* Soft-deleted tasks are absent from the live API even when a mutation reports the tombstone with TaskDeletedError. Preserve structural matching for duplicated bundled/workspace core instances.
*/
if (isTaskNotFoundError(error) || error instanceof TaskDeletedError) return true;
if (error && typeof error === "object" && (error as { name?: unknown }).name === "TaskDeletedError") return true;
return (error as NodeJS.ErrnoException | undefined)?.code === "ENOENT"; return (error as NodeJS.ErrnoException | undefined)?.code === "ENOENT";
} }