From 15b21dead108c0a240ad503587c45223ff06070b Mon Sep 17 00:00:00 2001 From: Phil Larson Date: Wed, 29 Jul 2026 23:27:32 -0700 Subject: [PATCH] fix(dashboard): reconcile task state through live API (#2595) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 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 ## 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). --- .changeset/task-step-live-api.md | 7 + .../store-wedge-resolution.pg.test.ts | 311 ++++++++++++++++++ packages/core/src/store.ts | 5 +- .../core/src/task-store/async-persistence.ts | 27 ++ packages/core/src/task-store/persistence.ts | 18 +- .../core/src/task-store/project-store-ops.ts | 3 +- .../core/src/task-store/task-mutation-ops.ts | 59 +++- .../task-store/workflow-task-create-ops.ts | 2 + ...r-task-workflow-routes.step-update.test.ts | 201 +++++++++++ .../routes/register-task-workflow-routes.ts | 56 ++++ .../dashboard/src/routes/task-lookup-error.ts | 9 +- 11 files changed, 691 insertions(+), 7 deletions(-) create mode 100644 .changeset/task-step-live-api.md create mode 100644 packages/core/src/__tests__/postgres/store-wedge-resolution.pg.test.ts create mode 100644 packages/dashboard/src/routes/__tests__/register-task-workflow-routes.step-update.test.ts diff --git a/.changeset/task-step-live-api.md b/.changeset/task-step-live-api.md new file mode 100644 index 0000000000..10a4252657 --- /dev/null +++ b/.changeset/task-step-live-api.md @@ -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. diff --git a/packages/core/src/__tests__/postgres/store-wedge-resolution.pg.test.ts b/packages/core/src/__tests__/postgres/store-wedge-resolution.pg.test.ts new file mode 100644 index 0000000000..f220330cf1 --- /dev/null +++ b/packages/core/src/__tests__/postgres/store-wedge-resolution.pg.test.ts @@ -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((resolve) => { + reportProjectionPending = resolve; + }); + let releaseProjection!: () => void; + const projectionMayStart = new Promise((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((resolve) => { + reportResolutionProjection = resolve; + }); + let releaseResolutionProjection!: () => void; + const resolutionProjectionMayFinish = new Promise((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((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((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((resolve) => { + releaseReplacement = resolve; + }); + let reportLocked!: () => void; + const rowLocked = new Promise((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((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; + let reportStaleSnapshot!: () => void; + const staleSnapshotReady = new Promise((resolve) => { + reportStaleSnapshot = resolve; + }); + let releaseStaleWrite!: () => void; + const staleWriteMayPersist = new Promise((resolve) => { + releaseStaleWrite = resolve; + }); + const persistSpy = vi.spyOn(store, persistMethod) as unknown as { + mockImplementationOnce: (implementation: (...args: unknown[]) => Promise) => 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", + }); + }); +}); diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index 337788dd76..a585147129 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -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 { async updateTaskAtomic( id: string, updater: ( current: Task, ) => Parameters[1] | null | undefined | Promise[1] | null | undefined>, runContext?: RunMutationContext, ): Promise { 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 | undefined, patch: Record, ): Record { return mergeCustomFieldPatchImpl(this, current, patch); } diff --git a/packages/core/src/task-store/async-persistence.ts b/packages/core/src/task-store/async-persistence.ts index 4024f0f5df..6f75516edc 100644 --- a/packages/core/src/task-store/async-persistence.ts +++ b/packages/core/src/task-store/async-persistence.ts @@ -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 | 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 | undefined; +} + /** * FNXC:TaskStoreArchiveLineage 2026-06-24-15:00: * Soft-delete a task INSIDE a shared transaction handle. This is the diff --git a/packages/core/src/task-store/persistence.ts b/packages/core/src/task-store/persistence.ts index de7e289ae3..39bfda67cd 100644 --- a/packages/core/src/task-store/persistence.ts +++ b/packages/core/src/task-store/persistence.ts @@ -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, task: Task): void { + const durable = fromJson(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. diff --git a/packages/core/src/task-store/project-store-ops.ts b/packages/core/src/task-store/project-store-ops.ts index ba705deac0..2a22488069 100644 --- a/packages/core/src/task-store/project-store-ops.ts +++ b/packages/core/src/task-store/project-store-ops.ts @@ -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); diff --git a/packages/core/src/task-store/task-mutation-ops.ts b/packages/core/src/task-store/task-mutation-ops.ts index 8483a3463f..3a255e2a1d 100644 --- a/packages/core/src/task-store/task-mutation-ops.ts +++ b/packages/core/src/task-store/task-mutation-ops.ts @@ -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 { /* * FNXC:SqliteFinalRemoval 2026-06-26: diff --git a/packages/core/src/task-store/workflow-task-create-ops.ts b/packages/core/src/task-store/workflow-task-create-ops.ts index bd80251170..57fdaf3117 100644 --- a/packages/core/src/task-store/workflow-task-create-ops.ts +++ b/packages/core/src/task-store/workflow-task-create-ops.ts @@ -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"); diff --git a/packages/dashboard/src/routes/__tests__/register-task-workflow-routes.step-update.test.ts b/packages/dashboard/src/routes/__tests__/register-task-workflow-routes.step-update.test.ts new file mode 100644 index 0000000000..12381c9dd3 --- /dev/null +++ b/packages/dashboard/src/routes/__tests__/register-task-workflow-routes.step-update.test.ts @@ -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" })); + }); +}); diff --git a/packages/dashboard/src/routes/register-task-workflow-routes.ts b/packages/dashboard/src/routes/register-task-workflow-routes.ts index 0355a523ab..e12b3b0095 100644 --- a/packages/dashboard/src/routes/register-task-workflow-routes.ts +++ b/packages/dashboard/src/routes/register-task-workflow-routes.ts @@ -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 { diff --git a/packages/dashboard/src/routes/task-lookup-error.ts b/packages/dashboard/src/routes/task-lookup-error.ts index 5ce5ca22e2..cfdba91446 100644 --- a/packages/dashboard/src/routes/task-lookup-error.ts +++ b/packages/dashboard/src/routes/task-lookup-error.ts @@ -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"; }