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 { 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 { 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 { 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";
@@ -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> {
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> {
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:
* 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 { 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). */
export interface TaskRow {
@@ -183,6 +183,22 @@ export type TaskColumnDescriptor = {
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:
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 {generateTaskLineageId} from "../task-lineage.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 {isWorkflowDefinitionIdPrimaryKeyCollision, nextWorkflowDefinitionIdAsyncImpl} from "../task-store/workflow-definitions.js";
import {upsertTaskRowInTransaction, buildTaskInsertValues} from "../task-store/async-persistence.js";
@@ -136,6 +136,7 @@ export async function atomicWriteTaskJsonWithAuditImpl(store: TaskStore, dir: st
*/
if (row) {
const existing = store.pgRowToTaskRow(row);
preserveResolvedTaskWedgeEpisode(existing, task);
const changedColumns = store.getChangedTaskColumns(existing, task);
if (changedColumns.size > 0) {
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.
*/
import {TaskStore, storeLog} from "../store.js";
import {TaskDeletedError} from "./errors.js";
import {TaskDeletedError, TaskNotFoundError} from "./errors.js";
import type {LegacyAutoMergeStampReconcileResult} from "../store.js";
import {randomUUID} from "node:crypto";
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 {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.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 {purgeTaskWorkflowSelectionRowsAsyncImpl} from "./workflow-definitions.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> {
/*
* 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 {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.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 {listDueWorkflowWorkItems as listDueWorkflowWorkItemsAsync, withTaskWorkflowSerialization} from "../task-store/async-workflow-workitems.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;
}
const existingRow = store.pgRowToTaskRow(pgRow);
preserveResolvedTaskWedgeEpisode(existingRow, task);
const deletedAt = store.getSoftDeletedWriteConflict(id, task, existingRow);
if (deletedAt) {
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
router.patch("/tasks/:id", async (req, res) => {
try {

View File

@@ -22,7 +22,7 @@
* signal so genuinely file-backed reads (attachments, session files, worktree
* files) keep their 404 behavior.
*/
import { isTaskNotFoundError } from "@fusion/core";
import { isTaskNotFoundError, TaskDeletedError } from "@fusion/core";
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.
*/
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";
}