diff --git a/docs/storage.md b/docs/storage.md index bd5feaf312..c979bf25f6 100644 --- a/docs/storage.md +++ b/docs/storage.md @@ -139,7 +139,10 @@ Important execution nuance: ## FTS5 task-index maintenance (FN-5943 / FN-5976) - Live task search uses the `tasks_fts` external-content FTS5 table in `fusion.db`; the archive log uses a separate `archived_tasks_fts` table in `archive.db`. -- `tasks_fts_au` is value-aware: even though hot task writes still upsert full rows, the trigger only fires when indexed text actually changes (`id`, `title`, `description`, `comments`, `deletedAt`). Status/step/worktree churn no longer rewrites the FTS row on every update. +- `tasks_fts_au` is value-aware and column-scoped. Hot task mutations (`atomicWriteTaskJson` / `atomicWriteTaskJsonWithAudit`) now diff the current row against the incoming task and issue `UPDATE tasks SET , updatedAt = ? WHERE id = ?` instead of rewriting the full task row. Non-text churn (status, steps, leases, scheduler stamps) therefore skips the FTS trigger entirely because those UPDATEs omit the indexed text columns. +- Full-row task persistence is still intentional for create/restore/replication-class paths: `insertTask` / `atomicCreateTaskJson` remain plain `INSERT`, and direct replication-style upserts (`upsertTaskWithFtsRecovery`, for example task-metadata snapshot application) still use the generated full-row `INSERT ... ON CONFLICT DO UPDATE` form. +- After a partial SQLite update, Fusion rewrites compatibility `task.json` from a fresh DB read so the disk mirror stays byte-aligned with the authoritative row even on narrow SQL patches. +- Checkout lease renewal has its own targeted path (`renewCheckoutLease`), updating only `checkoutRunId`, `checkoutLeaseRenewedAt`, and `updatedAt` instead of routing through the broad `updateTask(...)` mutator. - `Database.getFtsIndexBytes()` measures index size via `SELECT SUM(LENGTH(block)) FROM tasks_fts_data`. Fusion intentionally does **not** rely on `dbstat`, because node:sqlite builds do not guarantee `SQLITE_ENABLE_DBSTAT_VTAB`. - `SelfHealingManager` Batch 1 now runs `fts-maintenance` when `fts5Available === true`: - every maintenance tick: incremental `merge` compaction diff --git a/packages/core/src/__tests__/task-partial-update.test.ts b/packages/core/src/__tests__/task-partial-update.test.ts new file mode 100644 index 0000000000..32acef630c --- /dev/null +++ b/packages/core/src/__tests__/task-partial-update.test.ts @@ -0,0 +1,191 @@ +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import { join } from "node:path"; +import { readFile } from "node:fs/promises"; + +import { TaskDeletedError, type TaskStore } from "../store.js"; +import type { Task } from "../types.js"; +import { createSharedTaskStoreTestHarness } from "./store-test-helpers.js"; + +describe("TaskStore partial task updates", () => { + const harness = createSharedTaskStoreTestHarness(); + let store: TaskStore; + let rootDir: string; + + beforeAll(harness.beforeAll); + + beforeEach(async () => { + await harness.beforeEach(); + store = harness.store(); + rootDir = harness.rootDir(); + }); + + afterEach(async () => { + await harness.afterEach(); + }); + + afterAll(harness.afterAll); + + async function captureTaskSql(action: () => Promise): Promise { + const db = store.getDatabase(); + const originalPrepare = db.prepare.bind(db); + const statements: string[] = []; + const spy = vi.spyOn(db, "prepare").mockImplementation(((sql: string) => { + if (/UPDATE tasks\s+SET|INSERT INTO tasks/i.test(sql)) { + statements.push(sql.replace(/\s+/g, " ").trim()); + } + return originalPrepare(sql); + }) as typeof db.prepare); + try { + await action(); + } finally { + spy.mockRestore(); + } + return statements; + } + + function expectLatestUpdate(statements: string[]): string { + const sql = [...statements].reverse().find((statement) => statement.startsWith("UPDATE tasks SET") || statement.startsWith("UPDATE tasks SET ")); + expect(sql).toBeTruthy(); + return sql!; + } + + it("updates only changed fields plus log and updatedAt for a hot updateTask path", async () => { + const task = await store.createTask({ title: "Hot path", description: "status flip" }); + + const statements = await captureTaskSql(() => store.updateTask(task.id, { status: "working" })); + const sql = expectLatestUpdate(statements); + + expect(sql).toContain("status = ?"); + expect(sql).toContain("updatedAt = ?"); + expect(sql).not.toContain("description = ?"); + expect(sql).not.toContain("steps = ?"); + expect(sql).not.toContain("tokenUsageTotalTokens = ?"); + + const diskTask = JSON.parse(await readFile(join(rootDir, ".fusion", "tasks", task.id, "task.json"), "utf8")) as Task; + expect(diskTask.status).toBe("working"); + }); + + it("keeps no-op updates narrow and excludes unchanged fields from the SET list", async () => { + const task = await store.createTask({ title: "Same title", description: "unchanged" }); + + const statements = await captureTaskSql(() => store.updateTask(task.id, { title: "Same title" })); + const sql = expectLatestUpdate(statements); + + expect(sql).toContain("updatedAt = ?"); + expect(sql).not.toContain("title = ?"); + expect(sql).not.toContain("description = ?"); + expect(sql).not.toContain("steps = ?"); + }); + + it("persists field clears via the partial path", async () => { + const task = await store.createTask({ title: "Clear me", description: "field clear" }); + await store.updateTask(task.id, { error: "boom" }); + + const statements = await captureTaskSql(() => store.updateTask(task.id, { error: null })); + const sql = expectLatestUpdate(statements); + expect(sql).toContain("error = ?"); + expect(sql).not.toContain("description = ?"); + + const refreshed = await store.getTask(task.id); + expect(refreshed?.error).toBeUndefined(); + }); + + it("preserves FTS behavior for text changes and skips FTS rewrites for non-text changes", async () => { + const task = await store.createTask({ title: "Alpha", description: "Bravo" }); + const db = store.getDatabase(); + const ftsRowQuery = ` + SELECT fts.rowid as rowid, tasks.id as id, fts.title as title, fts.description as description, fts.comments as comments + FROM tasks_fts fts + JOIN tasks ON tasks.rowid = fts.rowid + WHERE tasks.id = ? + `; + + const before = db.prepare(ftsRowQuery).get(task.id) as Record | undefined; + expect(before).toBeTruthy(); + + await store.updateTask(task.id, { status: "queued" }); + const afterStatus = db.prepare(ftsRowQuery).get(task.id) as Record | undefined; + expect(afterStatus).toEqual(before); + + await store.updateTask(task.id, { title: "Needle title" }); + const results = await store.searchTasks("Needle"); + expect(results.map((entry) => entry.id)).toContain(task.id); + }); + + it("preserves soft-delete guard parity and records resurrection-blocked audit events", async () => { + const task = await store.createTask({ title: "Deleted", description: "guard" }); + const dir = join(rootDir, ".fusion", "tasks", task.id); + await store.deleteTask(task.id); + + await expect((store as any).atomicWriteTaskJson(dir, { ...task, title: "after delete" })).rejects.toBeInstanceOf(TaskDeletedError); + + const events = (store as any).db.prepare( + "SELECT mutationType, metadata FROM runAuditEvents WHERE taskId = ? AND mutationType = ? ORDER BY timestamp ASC" + ).all(task.id, "task:resurrection-blocked") as Array<{ mutationType: string; metadata: string | null }>; + expect(events).toHaveLength(1); + expect(events[0]?.mutationType).toBe("task:resurrection-blocked"); + expect(events[0]?.metadata ?? "").toContain("atomicWriteTaskJson"); + }); + + it("keeps run-audit metadata parity on updateTask with runContext", async () => { + const task = await store.createTask({ title: "Audit me", description: "audit" }); + + await store.updateTask(task.id, { title: "Audited" }, { runId: "run-partial-update", agentId: "agent-1" }); + + const events = store.getRunAuditEvents({ runId: "run-partial-update" }); + const updateEvent = events.find((event) => event.mutationType === "task:update"); + expect(updateEvent?.metadata).toEqual({ updatedFields: ["title"] }); + }); + + it("bumps lastModified exactly once per converted update mutation", async () => { + const task = await store.createTask({ title: "Single bump", description: "counter" }); + const db = store.getDatabase(); + const bumpSpy = vi.spyOn(db, "bumpLastModified"); + + await store.updateTask(task.id, { status: "queued" }); + + expect(bumpSpy).toHaveBeenCalledTimes(1); + }); + + it("renews checkout leases with a targeted checkout UPDATE", async () => { + const task = await store.createTask({ title: "Lease", description: "renew" }); + await store.updateTask(task.id, { + checkedOutBy: "agent-1", + checkedOutAt: "2026-01-01T00:00:00.000Z", + checkoutNodeId: "node-1", + checkoutRunId: "run-old", + checkoutLeaseRenewedAt: "2026-01-01T00:00:00.000Z", + checkoutLeaseEpoch: 1, + }); + + const renewedAt = "2026-01-01T00:01:00.000Z"; + const statements = await captureTaskSql(() => store.renewCheckoutLease(task.id, { + checkoutRunId: "run-new", + checkoutLeaseRenewedAt: renewedAt, + })); + const sql = expectLatestUpdate(statements); + + expect(sql).toContain("checkoutRunId = ?"); + expect(sql).toContain("checkoutLeaseRenewedAt = ?"); + expect(sql).toContain("updatedAt = ?"); + expect(sql).not.toContain("title = ?"); + expect(sql).not.toContain("description = ?"); + expect(sql).not.toContain("steps = ?"); + + const refreshed = await store.getTask(task.id); + expect(refreshed?.checkoutRunId).toBe("run-new"); + expect(refreshed?.checkoutLeaseRenewedAt).toBe(renewedAt); + }); + + it("keeps create and direct upsert/replication paths on full-row SQL", async () => { + const createStatements = await captureTaskSql(() => store.createTask({ title: "Create path", description: "full insert" })); + expect(createStatements.some((statement) => statement.startsWith("INSERT INTO tasks (") && !statement.startsWith("UPDATE tasks SET"))).toBe(true); + + const task = await store.createTask({ title: "Replicate me", description: "full upsert" }); + const replicatedTask = { ...task, title: "Replicated title", updatedAt: new Date(Date.now() + 1_000).toISOString() }; + const upsertStatements = await captureTaskSql(async () => { + (store as any).upsertTaskWithFtsRecovery(replicatedTask); + }); + expect(upsertStatements.some((statement) => statement.includes("ON CONFLICT(id) DO UPDATE SET"))).toBe(true); + }); +}); diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index a44600b437..923527d36f 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -282,6 +282,160 @@ interface TaskRow { allowResurrection: number | null; } +type TaskPersistSerializationContext = { + lineageId: string; +}; + +type TaskColumnDescriptor = { + column: keyof TaskRow; + sqlIdentifier: string; + serialize: (task: Task, context: TaskPersistSerializationContext) => unknown; +}; + +function defineTaskColumn( + column: keyof TaskRow, + serialize: TaskColumnDescriptor["serialize"], + sqlIdentifier: string = column, +): TaskColumnDescriptor { + return { column, sqlIdentifier, serialize }; +} + +const serializeTaskAutoMerge: TaskColumnDescriptor["serialize"] = (task) => task.autoMerge === undefined ? null : (task.autoMerge ? 1 : 0); + +// Keep this descriptor order in lockstep with the named-column INSERT/UPSERT +// clauses we generate below. SQLite binds by the explicit column list we emit, +// so this logical persist order does not need to match the table's physical +// column layout from CREATE TABLE + migrations. +const TASK_COLUMN_DESCRIPTORS: TaskColumnDescriptor[] = [ + defineTaskColumn("id", (task) => task.id), + defineTaskColumn("lineageId", (_task, context) => context.lineageId), + defineTaskColumn("title", (task) => task.title ?? null), + defineTaskColumn("description", (task) => task.description ?? ""), + defineTaskColumn("priority", (task) => normalizeTaskPriority(task.priority)), + defineTaskColumn("column", (task) => task.column, '"column"'), + defineTaskColumn("status", (task) => task.status ?? null), + defineTaskColumn("size", (task) => task.size ?? null), + defineTaskColumn("reviewLevel", (task) => task.reviewLevel ?? null), + defineTaskColumn("currentStep", (task) => task.currentStep || 0), + defineTaskColumn("worktree", (task) => task.worktree ?? null), + defineTaskColumn("blockedBy", (task) => task.blockedBy ?? null), + defineTaskColumn("overlapBlockedBy", (task) => task.overlapBlockedBy ?? null), + defineTaskColumn("paused", (task) => task.paused ? 1 : 0), + defineTaskColumn("pausedReason", (task) => task.pausedReason ?? null), + defineTaskColumn("userPaused", (task) => task.userPaused ? 1 : 0), + defineTaskColumn("baseBranch", (task) => task.baseBranch ?? null), + defineTaskColumn("branch", (task) => task.branch ?? null), + defineTaskColumn("autoMerge", serializeTaskAutoMerge), + defineTaskColumn("executionStartBranch", (task) => task.executionStartBranch ?? null), + defineTaskColumn("baseCommitSha", (task) => task.baseCommitSha ?? null), + defineTaskColumn("modelPresetId", (task) => task.modelPresetId ?? null), + defineTaskColumn("modelProvider", (task) => task.modelProvider ?? null), + defineTaskColumn("modelId", (task) => task.modelId ?? null), + defineTaskColumn("validatorModelProvider", (task) => task.validatorModelProvider ?? null), + defineTaskColumn("validatorModelId", (task) => task.validatorModelId ?? null), + defineTaskColumn("planningModelProvider", (task) => task.planningModelProvider ?? null), + defineTaskColumn("planningModelId", (task) => task.planningModelId ?? null), + defineTaskColumn("mergeRetries", (task) => task.mergeRetries ?? null), + defineTaskColumn("workflowStepRetries", (task) => task.workflowStepRetries ?? null), + defineTaskColumn("stuckKillCount", (task) => task.stuckKillCount ?? 0), + defineTaskColumn("resumeLimboCount", (task) => task.resumeLimboCount ?? 0), + defineTaskColumn("resumeLimboTipSha", (task) => task.resumeLimboTipSha ?? null), + defineTaskColumn("resumeLimboStepSignature", (task) => task.resumeLimboStepSignature ?? null), + defineTaskColumn("postReviewFixCount", (task) => task.postReviewFixCount ?? 0), + defineTaskColumn("recoveryRetryCount", (task) => task.recoveryRetryCount ?? null), + defineTaskColumn("taskDoneRetryCount", (task) => task.taskDoneRetryCount ?? 0), + defineTaskColumn("worktreeSessionRetryCount", (task) => task.worktreeSessionRetryCount ?? 0), + defineTaskColumn("completionHandoffLimboRecoveryCount", (task) => task.completionHandoffLimboRecoveryCount ?? 0), + defineTaskColumn("verificationFailureCount", (task) => task.verificationFailureCount ?? 0), + defineTaskColumn("mergeConflictBounceCount", (task) => task.mergeConflictBounceCount ?? 0), + defineTaskColumn("mergeAuditBounceCount", (task) => task.mergeAuditBounceCount ?? 0), + defineTaskColumn("mergeTransientRetryCount", (task) => task.mergeTransientRetryCount ?? 0), + defineTaskColumn("branchConflictRecoveryCount", (task) => task.branchConflictRecoveryCount ?? 0), + defineTaskColumn("reviewerContextRetryCount", (task) => task.reviewerContextRetryCount ?? 0), + defineTaskColumn("reviewerFallbackRetryCount", (task) => task.reviewerFallbackRetryCount ?? 0), + defineTaskColumn("nextRecoveryAt", (task) => task.nextRecoveryAt ?? null), + defineTaskColumn("error", (task) => task.error ?? null), + defineTaskColumn("summary", (task) => task.summary ?? null), + defineTaskColumn("thinkingLevel", (task) => task.thinkingLevel ?? null), + defineTaskColumn("executionMode", (task) => task.executionMode ?? null), + defineTaskColumn("tokenUsageInputTokens", (task) => task.tokenUsage?.inputTokens ?? null), + defineTaskColumn("tokenUsageOutputTokens", (task) => task.tokenUsage?.outputTokens ?? null), + defineTaskColumn("tokenUsageCachedTokens", (task) => task.tokenUsage?.cachedTokens ?? null), + defineTaskColumn("tokenUsageCacheWriteTokens", (task) => task.tokenUsage?.cacheWriteTokens ?? null), + defineTaskColumn("tokenUsageTotalTokens", (task) => task.tokenUsage?.totalTokens ?? null), + defineTaskColumn("tokenUsageFirstUsedAt", (task) => task.tokenUsage?.firstUsedAt ?? null), + defineTaskColumn("tokenUsageLastUsedAt", (task) => task.tokenUsage?.lastUsedAt ?? null), + defineTaskColumn("tokenBudgetSoftAlertedAt", (task) => task.tokenBudgetSoftAlertedAt ?? null), + defineTaskColumn("tokenBudgetHardAlertedAt", (task) => task.tokenBudgetHardAlertedAt ?? null), + defineTaskColumn("tokenBudgetOverride", (task) => toJsonNullable(task.tokenBudgetOverride)), + defineTaskColumn("createdAt", (task) => task.createdAt), + defineTaskColumn("updatedAt", (task) => task.updatedAt), + defineTaskColumn("columnMovedAt", (task) => task.columnMovedAt ?? null), + defineTaskColumn("firstExecutionAt", (task) => task.firstExecutionAt ?? null), + defineTaskColumn("cumulativeActiveMs", (task) => task.cumulativeActiveMs ?? null), + defineTaskColumn("executionStartedAt", (task) => task.executionStartedAt ?? null), + defineTaskColumn("executionCompletedAt", (task) => task.executionCompletedAt ?? null), + defineTaskColumn("dependencies", (task) => toJson(task.dependencies || [])), + defineTaskColumn("steps", (task) => toJson(task.steps || [])), + defineTaskColumn("customFields", (task) => toJson(task.customFields ?? {})), + defineTaskColumn("log", (task) => toJson(task.log || [])), + defineTaskColumn("attachments", (task) => toJson(task.attachments || [])), + defineTaskColumn("steeringComments", (task) => toJson(task.steeringComments || [])), + defineTaskColumn("comments", (task) => toJson(task.comments || [])), + defineTaskColumn("review", (task) => toJsonNullable(task.review)), + defineTaskColumn("reviewState", (task) => toJsonNullable(task.reviewState)), + defineTaskColumn("workflowStepResults", (task) => toJson(task.workflowStepResults || [])), + defineTaskColumn("prInfo", (task) => toJsonNullable(task.prInfo)), + defineTaskColumn("prInfos", (task) => toJson(task.prInfos || [])), + defineTaskColumn("issueInfo", (task) => toJsonNullable(task.issueInfo)), + defineTaskColumn("githubTracking", (task) => toJsonNullable(task.githubTracking)), + defineTaskColumn("sourceIssueProvider", (task) => task.sourceIssue?.provider ?? null), + defineTaskColumn("sourceIssueRepository", (task) => task.sourceIssue?.repository ?? null), + defineTaskColumn("sourceIssueExternalIssueId", (task) => task.sourceIssue?.externalIssueId ?? null), + defineTaskColumn("sourceIssueNumber", (task) => task.sourceIssue?.issueNumber ?? null), + defineTaskColumn("sourceIssueUrl", (task) => task.sourceIssue?.url ?? null), + defineTaskColumn("mergeDetails", (task) => toJsonNullable(task.mergeDetails)), + defineTaskColumn("breakIntoSubtasks", (task) => task.breakIntoSubtasks ? 1 : 0), + defineTaskColumn("noCommitsExpected", (task) => task.noCommitsExpected ? 1 : 0), + defineTaskColumn("enabledWorkflowSteps", (task) => toJson(task.enabledWorkflowSteps || [])), + defineTaskColumn("modifiedFiles", (task) => toJson(task.modifiedFiles || [])), + defineTaskColumn("missionId", (task) => task.missionId ?? null), + defineTaskColumn("sliceId", (task) => task.sliceId ?? null), + defineTaskColumn("scopeOverride", (task) => task.scopeOverride ? 1 : null), + defineTaskColumn("scopeOverrideReason", (task) => task.scopeOverrideReason ?? null), + defineTaskColumn("scopeAutoWiden", (task) => toJson(task.scopeAutoWiden || [])), + defineTaskColumn("assignedAgentId", (task) => task.assignedAgentId ?? null), + defineTaskColumn("pausedByAgentId", (task) => task.pausedByAgentId ?? null), + defineTaskColumn("assigneeUserId", (task) => task.assigneeUserId ?? null), + defineTaskColumn("nodeId", (task) => task.nodeId ?? null), + defineTaskColumn("effectiveNodeId", (task) => task.effectiveNodeId ?? null), + defineTaskColumn("effectiveNodeSource", (task) => task.effectiveNodeSource ?? null), + defineTaskColumn("sourceType", (task) => task.sourceType ?? null), + defineTaskColumn("sourceAgentId", (task) => task.sourceAgentId ?? null), + defineTaskColumn("sourceRunId", (task) => task.sourceRunId ?? null), + defineTaskColumn("sourceSessionId", (task) => task.sourceSessionId ?? null), + defineTaskColumn("sourceMessageId", (task) => task.sourceMessageId ?? null), + defineTaskColumn("sourceParentTaskId", (task) => task.sourceParentTaskId ?? null), + defineTaskColumn("sourceMetadata", (task) => toJsonNullable(task.sourceMetadata)), + defineTaskColumn("checkedOutBy", (task) => task.checkedOutBy ?? null), + defineTaskColumn("checkedOutAt", (task) => task.checkedOutAt ?? null), + defineTaskColumn("checkoutNodeId", (task) => task.checkoutNodeId ?? null), + defineTaskColumn("checkoutRunId", (task) => task.checkoutRunId ?? null), + defineTaskColumn("checkoutLeaseRenewedAt", (task) => task.checkoutLeaseRenewedAt ?? null), + defineTaskColumn("checkoutLeaseEpoch", (task) => task.checkoutLeaseEpoch ?? 0), + defineTaskColumn("deletedAt", (task) => task.deletedAt ?? null), + defineTaskColumn("allowResurrection", (task) => task.allowResurrection ? 1 : 0), +]; + +const TASK_COLUMN_DESCRIPTOR_BY_COLUMN = new Map( + TASK_COLUMN_DESCRIPTORS.map((descriptor) => [descriptor.column, descriptor]), +); +const TASK_PERSIST_SQL_COLUMNS = TASK_COLUMN_DESCRIPTORS.map((descriptor) => descriptor.sqlIdentifier).join(", "); +const TASK_UPSERT_SQL_ASSIGNMENTS = TASK_COLUMN_DESCRIPTORS + .filter((descriptor) => descriptor.column !== "id") + .map((descriptor) => ` ${descriptor.sqlIdentifier} = excluded.${descriptor.sqlIdentifier}`) + .join(",\n"); + /** Database row shape for the task_documents table. */ const TASK_BRANCH_CONTEXT_METADATA_KEY = "fusionBranchContext"; @@ -2340,128 +2494,23 @@ export class TaskStore extends EventEmitter { return [...columns, limitedLog].join(", "); } - private getTaskPersistValues(task: Task): unknown[] { - return [ - task.id, - task.lineageId ?? generateTaskLineageId(), - task.title ?? null, - task.description ?? "", - normalizeTaskPriority(task.priority), - task.column, - task.status ?? null, - task.size ?? null, - task.reviewLevel ?? null, - task.currentStep || 0, - task.worktree ?? null, - task.blockedBy ?? null, - task.overlapBlockedBy ?? null, - task.paused ? 1 : 0, - task.pausedReason ?? null, - task.userPaused ? 1 : 0, - task.baseBranch ?? null, - task.branch ?? null, - task.autoMerge === undefined ? null : (task.autoMerge ? 1 : 0), - task.executionStartBranch ?? null, - task.baseCommitSha ?? null, - task.modelPresetId ?? null, - task.modelProvider ?? null, - task.modelId ?? null, - task.validatorModelProvider ?? null, - task.validatorModelId ?? null, - task.planningModelProvider ?? null, - task.planningModelId ?? null, - task.mergeRetries ?? null, - task.workflowStepRetries ?? null, - task.stuckKillCount ?? 0, - task.resumeLimboCount ?? 0, - task.resumeLimboTipSha ?? null, - task.resumeLimboStepSignature ?? null, - task.postReviewFixCount ?? 0, - task.recoveryRetryCount ?? null, - task.taskDoneRetryCount ?? 0, - task.worktreeSessionRetryCount ?? 0, - task.completionHandoffLimboRecoveryCount ?? 0, - task.verificationFailureCount ?? 0, - task.mergeConflictBounceCount ?? 0, - task.mergeAuditBounceCount ?? 0, - task.mergeTransientRetryCount ?? 0, - task.branchConflictRecoveryCount ?? 0, - task.reviewerContextRetryCount ?? 0, - task.reviewerFallbackRetryCount ?? 0, - task.nextRecoveryAt ?? null, - task.error ?? null, - task.summary ?? null, - task.thinkingLevel ?? null, - task.executionMode ?? null, - task.tokenUsage?.inputTokens ?? null, - task.tokenUsage?.outputTokens ?? null, - task.tokenUsage?.cachedTokens ?? null, - task.tokenUsage?.cacheWriteTokens ?? null, - task.tokenUsage?.totalTokens ?? null, - task.tokenUsage?.firstUsedAt ?? null, - task.tokenUsage?.lastUsedAt ?? null, - task.tokenBudgetSoftAlertedAt ?? null, - task.tokenBudgetHardAlertedAt ?? null, - toJsonNullable(task.tokenBudgetOverride), - task.createdAt, - task.updatedAt, - task.columnMovedAt ?? null, - task.firstExecutionAt ?? null, - task.cumulativeActiveMs ?? null, - task.executionStartedAt ?? null, - task.executionCompletedAt ?? null, - toJson(task.dependencies || []), - toJson(task.steps || []), - toJson(task.customFields ?? {}), - toJson(task.log || []), - toJson(task.attachments || []), - toJson(task.steeringComments || []), - toJson(task.comments || []), - toJsonNullable(task.review), - toJsonNullable(task.reviewState), - toJson(task.workflowStepResults || []), - toJsonNullable(task.prInfo), - toJson(task.prInfos || []), - toJsonNullable(task.issueInfo), - toJsonNullable(task.githubTracking), - task.sourceIssue?.provider ?? null, - task.sourceIssue?.repository ?? null, - task.sourceIssue?.externalIssueId ?? null, - task.sourceIssue?.issueNumber ?? null, - task.sourceIssue?.url ?? null, - toJsonNullable(task.mergeDetails), - task.breakIntoSubtasks ? 1 : 0, - task.noCommitsExpected ? 1 : 0, - task.autoMerge === undefined ? null : task.autoMerge ? 1 : 0, - toJson(task.enabledWorkflowSteps || []), - toJson(task.modifiedFiles || []), - task.missionId ?? null, - task.sliceId ?? null, - task.scopeOverride ? 1 : null, - task.scopeOverrideReason ?? null, - toJson(task.scopeAutoWiden || []), - task.assignedAgentId ?? null, - task.pausedByAgentId ?? null, - task.assigneeUserId ?? null, - task.nodeId ?? null, - task.effectiveNodeId ?? null, - task.effectiveNodeSource ?? null, - task.sourceType ?? null, - task.sourceAgentId ?? null, - task.sourceRunId ?? null, - task.sourceSessionId ?? null, - task.sourceMessageId ?? null, - task.sourceParentTaskId ?? null, - toJsonNullable(task.sourceMetadata), - task.checkedOutBy ?? null, - task.checkedOutAt ?? null, - task.checkoutNodeId ?? null, - task.checkoutRunId ?? null, - task.checkoutLeaseRenewedAt ?? null, - task.checkoutLeaseEpoch ?? 0, - task.deletedAt ?? null, - task.allowResurrection ? 1 : 0, - ]; + private createTaskPersistSerializationContext( + task: Task, + existingRow?: Pick, + ): TaskPersistSerializationContext { + return { + lineageId: task.lineageId ?? existingRow?.lineageId ?? generateTaskLineageId(), + }; + } + + private getTaskPersistValues(task: Task, existingRow?: Pick): unknown[] { + const context = this.createTaskPersistSerializationContext(task, existingRow); + return TASK_COLUMN_DESCRIPTORS.map((descriptor) => descriptor.serialize(task, context)); + } + + private readTaskRowFromDb(id: string, options?: { includeDeleted?: boolean }): TaskRow | undefined { + const whereClause = options?.includeDeleted ? "id = ?" : `id = ? AND ${TaskStore.ACTIVE_TASKS_WHERE}`; + return this.db.prepare(`SELECT * FROM tasks WHERE ${whereClause}`).get(id) as TaskRow | undefined; } /** @@ -2472,19 +2521,8 @@ export class TaskStore extends EventEmitter { const values = this.getTaskPersistValues(task); const placeholders = values.map(() => "?").join(", "); this.db.prepare(` - INSERT INTO tasks ( - id, lineageId, title, description, priority, "column", status, size, reviewLevel, currentStep, - worktree, blockedBy, overlapBlockedBy, paused, pausedReason, userPaused, baseBranch, branch, autoMerge, executionStartBranch, baseCommitSha, modelPresetId, modelProvider, - modelId, validatorModelProvider, validatorModelId, planningModelProvider, planningModelId, mergeRetries, - workflowStepRetries, stuckKillCount, resumeLimboCount, resumeLimboTipSha, resumeLimboStepSignature, postReviewFixCount, recoveryRetryCount, taskDoneRetryCount, worktreeSessionRetryCount, completionHandoffLimboRecoveryCount, verificationFailureCount, mergeConflictBounceCount, mergeAuditBounceCount, mergeTransientRetryCount, branchConflictRecoveryCount, reviewerContextRetryCount, reviewerFallbackRetryCount, nextRecoveryAt, error, - summary, thinkingLevel, executionMode, tokenUsageInputTokens, tokenUsageOutputTokens, tokenUsageCachedTokens, - tokenUsageCacheWriteTokens, tokenUsageTotalTokens, tokenUsageFirstUsedAt, tokenUsageLastUsedAt, tokenBudgetSoftAlertedAt, tokenBudgetHardAlertedAt, tokenBudgetOverride, createdAt, updatedAt, columnMovedAt, - firstExecutionAt, cumulativeActiveMs, executionStartedAt, executionCompletedAt, - dependencies, steps, customFields, log, attachments, steeringComments, - comments, review, reviewState, workflowStepResults, prInfo, prInfos, issueInfo, githubTracking, - sourceIssueProvider, sourceIssueRepository, sourceIssueExternalIssueId, sourceIssueNumber, sourceIssueUrl, - mergeDetails, breakIntoSubtasks, noCommitsExpected, autoMerge, enabledWorkflowSteps, modifiedFiles, missionId, sliceId, scopeOverride, scopeOverrideReason, scopeAutoWiden, assignedAgentId, pausedByAgentId, assigneeUserId, nodeId, effectiveNodeId, effectiveNodeSource, sourceType, sourceAgentId, sourceRunId, sourceSessionId, sourceMessageId, sourceParentTaskId, sourceMetadata, checkedOutBy, checkedOutAt, checkoutNodeId, checkoutRunId, checkoutLeaseRenewedAt, checkoutLeaseEpoch, deletedAt, allowResurrection - ) VALUES (${placeholders}) + INSERT INTO tasks (${TASK_PERSIST_SQL_COLUMNS}) + VALUES (${placeholders}) `).run(...values); this.db.bumpLastModified(); } @@ -2499,139 +2537,11 @@ export class TaskStore extends EventEmitter { const values = this.getTaskPersistValues(task); const placeholders = values.map(() => "?").join(", "); this.db.prepare(` - INSERT INTO tasks ( - id, lineageId, title, description, priority, "column", status, size, reviewLevel, currentStep, - worktree, blockedBy, overlapBlockedBy, paused, pausedReason, userPaused, baseBranch, branch, autoMerge, executionStartBranch, baseCommitSha, modelPresetId, modelProvider, - modelId, validatorModelProvider, validatorModelId, planningModelProvider, planningModelId, mergeRetries, - workflowStepRetries, stuckKillCount, resumeLimboCount, resumeLimboTipSha, resumeLimboStepSignature, postReviewFixCount, recoveryRetryCount, taskDoneRetryCount, worktreeSessionRetryCount, completionHandoffLimboRecoveryCount, verificationFailureCount, mergeConflictBounceCount, mergeAuditBounceCount, mergeTransientRetryCount, branchConflictRecoveryCount, reviewerContextRetryCount, reviewerFallbackRetryCount, nextRecoveryAt, error, - summary, thinkingLevel, executionMode, tokenUsageInputTokens, tokenUsageOutputTokens, tokenUsageCachedTokens, - tokenUsageCacheWriteTokens, tokenUsageTotalTokens, tokenUsageFirstUsedAt, tokenUsageLastUsedAt, tokenBudgetSoftAlertedAt, tokenBudgetHardAlertedAt, tokenBudgetOverride, createdAt, updatedAt, columnMovedAt, - firstExecutionAt, cumulativeActiveMs, executionStartedAt, executionCompletedAt, - dependencies, steps, customFields, log, attachments, steeringComments, - comments, review, reviewState, workflowStepResults, prInfo, prInfos, issueInfo, githubTracking, - sourceIssueProvider, sourceIssueRepository, sourceIssueExternalIssueId, sourceIssueNumber, sourceIssueUrl, - mergeDetails, breakIntoSubtasks, noCommitsExpected, autoMerge, enabledWorkflowSteps, modifiedFiles, missionId, sliceId, scopeOverride, scopeOverrideReason, scopeAutoWiden, assignedAgentId, pausedByAgentId, assigneeUserId, nodeId, effectiveNodeId, effectiveNodeSource, sourceType, sourceAgentId, sourceRunId, sourceSessionId, sourceMessageId, sourceParentTaskId, sourceMetadata, checkedOutBy, checkedOutAt, checkoutNodeId, checkoutRunId, checkoutLeaseRenewedAt, checkoutLeaseEpoch, deletedAt, allowResurrection - ) VALUES (${placeholders}) + INSERT INTO tasks (${TASK_PERSIST_SQL_COLUMNS}) + VALUES (${placeholders}) ON CONFLICT(id) DO UPDATE SET - lineageId = excluded.lineageId, - title = excluded.title, - description = excluded.description, - priority = excluded.priority, - "column" = excluded."column", - status = excluded.status, - size = excluded.size, - reviewLevel = excluded.reviewLevel, - currentStep = excluded.currentStep, - worktree = excluded.worktree, - blockedBy = excluded.blockedBy, - overlapBlockedBy = excluded.overlapBlockedBy, - paused = excluded.paused, - pausedReason = excluded.pausedReason, - userPaused = excluded.userPaused, - baseBranch = excluded.baseBranch, - branch = excluded.branch, - autoMerge = excluded.autoMerge, - executionStartBranch = excluded.executionStartBranch, - baseCommitSha = excluded.baseCommitSha, - modelPresetId = excluded.modelPresetId, - modelProvider = excluded.modelProvider, - modelId = excluded.modelId, - validatorModelProvider = excluded.validatorModelProvider, - validatorModelId = excluded.validatorModelId, - planningModelProvider = excluded.planningModelProvider, - planningModelId = excluded.planningModelId, - mergeRetries = excluded.mergeRetries, - workflowStepRetries = excluded.workflowStepRetries, - stuckKillCount = excluded.stuckKillCount, - resumeLimboCount = excluded.resumeLimboCount, - resumeLimboTipSha = excluded.resumeLimboTipSha, - resumeLimboStepSignature = excluded.resumeLimboStepSignature, - postReviewFixCount = excluded.postReviewFixCount, - recoveryRetryCount = excluded.recoveryRetryCount, - taskDoneRetryCount = excluded.taskDoneRetryCount, - worktreeSessionRetryCount = excluded.worktreeSessionRetryCount, - completionHandoffLimboRecoveryCount = excluded.completionHandoffLimboRecoveryCount, - verificationFailureCount = excluded.verificationFailureCount, - mergeConflictBounceCount = excluded.mergeConflictBounceCount, - mergeAuditBounceCount = excluded.mergeAuditBounceCount, - mergeTransientRetryCount = excluded.mergeTransientRetryCount, - branchConflictRecoveryCount = excluded.branchConflictRecoveryCount, - reviewerContextRetryCount = excluded.reviewerContextRetryCount, - reviewerFallbackRetryCount = excluded.reviewerFallbackRetryCount, - nextRecoveryAt = excluded.nextRecoveryAt, - error = excluded.error, - summary = excluded.summary, - thinkingLevel = excluded.thinkingLevel, - executionMode = excluded.executionMode, - tokenUsageInputTokens = excluded.tokenUsageInputTokens, - tokenUsageOutputTokens = excluded.tokenUsageOutputTokens, - tokenUsageCachedTokens = excluded.tokenUsageCachedTokens, - tokenUsageCacheWriteTokens = excluded.tokenUsageCacheWriteTokens, - tokenUsageTotalTokens = excluded.tokenUsageTotalTokens, - tokenUsageFirstUsedAt = excluded.tokenUsageFirstUsedAt, - tokenUsageLastUsedAt = excluded.tokenUsageLastUsedAt, - tokenBudgetSoftAlertedAt = excluded.tokenBudgetSoftAlertedAt, - tokenBudgetHardAlertedAt = excluded.tokenBudgetHardAlertedAt, - tokenBudgetOverride = excluded.tokenBudgetOverride, - createdAt = excluded.createdAt, - updatedAt = excluded.updatedAt, - columnMovedAt = excluded.columnMovedAt, - firstExecutionAt = excluded.firstExecutionAt, - cumulativeActiveMs = excluded.cumulativeActiveMs, - executionStartedAt = excluded.executionStartedAt, - executionCompletedAt = excluded.executionCompletedAt, - dependencies = excluded.dependencies, - steps = excluded.steps, - customFields = excluded.customFields, - log = excluded.log, - attachments = excluded.attachments, - steeringComments = excluded.steeringComments, - comments = excluded.comments, - review = excluded.review, - reviewState = excluded.reviewState, - workflowStepResults = excluded.workflowStepResults, - prInfo = excluded.prInfo, - prInfos = excluded.prInfos, - issueInfo = excluded.issueInfo, - githubTracking = excluded.githubTracking, - sourceIssueProvider = excluded.sourceIssueProvider, - sourceIssueRepository = excluded.sourceIssueRepository, - sourceIssueExternalIssueId = excluded.sourceIssueExternalIssueId, - sourceIssueNumber = excluded.sourceIssueNumber, - sourceIssueUrl = excluded.sourceIssueUrl, - mergeDetails = excluded.mergeDetails, - breakIntoSubtasks = excluded.breakIntoSubtasks, - noCommitsExpected = excluded.noCommitsExpected, - autoMerge = excluded.autoMerge, - enabledWorkflowSteps = excluded.enabledWorkflowSteps, - modifiedFiles = excluded.modifiedFiles, - missionId = excluded.missionId, - sliceId = excluded.sliceId, - scopeOverride = excluded.scopeOverride, - scopeOverrideReason = excluded.scopeOverrideReason, - scopeAutoWiden = excluded.scopeAutoWiden, - assignedAgentId = excluded.assignedAgentId, - pausedByAgentId = excluded.pausedByAgentId, - assigneeUserId = excluded.assigneeUserId, - nodeId = excluded.nodeId, - effectiveNodeId = excluded.effectiveNodeId, - effectiveNodeSource = excluded.effectiveNodeSource, - sourceType = excluded.sourceType, - sourceAgentId = excluded.sourceAgentId, - sourceRunId = excluded.sourceRunId, - sourceSessionId = excluded.sourceSessionId, - sourceMessageId = excluded.sourceMessageId, - sourceParentTaskId = excluded.sourceParentTaskId, - sourceMetadata = excluded.sourceMetadata, - checkedOutBy = excluded.checkedOutBy, - checkedOutAt = excluded.checkedOutAt, - checkoutNodeId = excluded.checkoutNodeId, - checkoutRunId = excluded.checkoutRunId, - checkoutLeaseRenewedAt = excluded.checkoutLeaseRenewedAt, - checkoutLeaseEpoch = excluded.checkoutLeaseEpoch, - deletedAt = excluded.deletedAt, - allowResurrection = excluded.allowResurrection - `).run(...this.getTaskPersistValues(task)); +${TASK_UPSERT_SQL_ASSIGNMENTS} + `).run(...values); this.db.bumpLastModified(); } @@ -2689,33 +2599,120 @@ export class TaskStore extends EventEmitter { } } - private upsertTaskWithFtsRecovery(task: Task): void { + private runTaskFtsWriteWithRecovery(taskId: string, operation: string, write: () => void): void { try { - this.upsertTask(task); + write(); return; } catch (error) { if (!this.db.isFts5CorruptionError(error)) { throw error; } - console.warn(`[fusion:store] FTS5 corruption detected during upsert for task ${task.id}; rebuilding index and retrying once`); + console.warn(`[fusion:store] FTS5 corruption detected during ${operation} for task ${taskId}; rebuilding index and retrying once`); try { this.db.rebuildFts5Index(); } catch (rebuildError) { - console.warn("[fusion:store] FTS5 rebuild failed; propagating original upsert error", rebuildError); + console.warn(`[fusion:store] FTS5 rebuild failed; propagating original ${operation} error`, rebuildError); throw error; } try { - this.upsertTask(task); + write(); } catch (retryError) { - console.warn("[fusion:store] Upsert retry after FTS5 rebuild failed; propagating original upsert error", retryError); + console.warn(`[fusion:store] ${operation} retry after FTS5 rebuild failed; propagating original ${operation} error`, retryError); throw error; } } } + private upsertTaskWithFtsRecovery(task: Task): void { + this.runTaskFtsWriteWithRecovery(task.id, "upsert", () => { + this.upsertTask(task); + }); + } + + private getTaskPatchDescriptors(changedColumns: Iterable): TaskColumnDescriptor[] { + const descriptors: TaskColumnDescriptor[] = []; + for (const column of changedColumns) { + const descriptor = TASK_COLUMN_DESCRIPTOR_BY_COLUMN.get(column); + if (!descriptor) { + throw new Error(`Unknown task column for partial patch: ${String(column)}`); + } + descriptors.push(descriptor); + } + return descriptors; + } + + private getChangedTaskColumns(existingRow: TaskRow, task: Task): Set { + const nextValues = this.getTaskPersistValues(task, existingRow); + const changedColumns = new Set(); + for (const [index, descriptor] of TASK_COLUMN_DESCRIPTORS.entries()) { + if (descriptor.column === "updatedAt") { + continue; + } + if (!Object.is(existingRow[descriptor.column], nextValues[index])) { + changedColumns.add(descriptor.column); + } + } + return changedColumns; + } + + private patchTaskRowInTransaction( + id: string, + task: Task, + changedColumns: Iterable, + existingRow?: TaskRow, + ): { deletedAt?: string; current?: Task } { + const currentRow = existingRow ?? this.readTaskRowFromDb(id, { includeDeleted: true }); + const deletedAt = this.getSoftDeletedWriteConflict(id, task, currentRow); + if (deletedAt) { + return { deletedAt }; + } + if (!currentRow || currentRow.deletedAt != null) { + this.upsertTaskWithFtsRecovery(task); + return { current: this.readTaskFromDb(id) }; + } + + const patchDescriptors = this.getTaskPatchDescriptors(changedColumns); + const context = this.createTaskPersistSerializationContext(task, currentRow); + const assignments = patchDescriptors.map((descriptor) => `${descriptor.sqlIdentifier} = ?`); + assignments.push("updatedAt = ?"); + const values = patchDescriptors.map((descriptor) => descriptor.serialize(task, context)); + values.push(task.updatedAt, id); + + this.runTaskFtsWriteWithRecovery(id, "partial update", () => { + this.db.prepare(` + UPDATE tasks + SET ${assignments.join(", ")} + WHERE id = ? AND ${TaskStore.ACTIVE_TASKS_WHERE} + `).run(...values); + }); + this.db.bumpLastModified(); + return { current: this.readTaskFromDb(id) }; + } + + private async applyTaskPatch( + dir: string, + id: string, + task: Task, + changedColumns: Iterable, + options?: { existingRow?: TaskRow; auditInput?: { agentId?: string; runId?: string; timestamp?: string; operation?: string } }, + ): Promise { + let result: { deletedAt?: string; current?: Task } | undefined; + this.db.transactionImmediate(() => { + result = this.patchTaskRowInTransaction(id, task, changedColumns, options?.existingRow); + }); + if (result?.deletedAt) { + this.throwSoftDeletedWriteBlocked(id, result.deletedAt, options?.auditInput?.operation ?? "applyTaskPatch", { + agentId: options?.auditInput?.agentId, + runId: options?.auditInput?.runId, + timestamp: options?.auditInput?.timestamp, + }); + } + await this.writeTaskJsonFile(dir, result?.current ?? task); + } + /** * Read a task from SQLite by ID. */ @@ -2724,7 +2721,7 @@ export class TaskStore extends EventEmitter { ? this.getTaskSelectClauseWithActivityLogLimit(options.activityLogLimit) : "*"; const whereClause = options?.includeDeleted ? "id = ?" : `id = ? AND ${TaskStore.ACTIVE_TASKS_WHERE}`; - const row = this.db.prepare(`SELECT ${selectClause} FROM tasks WHERE ${whereClause}`).get(id) as unknown as TaskRow | undefined; + const row = this.db.prepare(`SELECT ${selectClause} FROM tasks WHERE ${whereClause}`).get(id) as TaskRow | undefined; if (!row) return undefined; return this.rowToTask(row); } @@ -3049,8 +3046,8 @@ export class TaskStore extends EventEmitter { ); } - private getSoftDeletedWriteConflict(id: string, task: Task): string | undefined { - const existing = this.readTaskFromDb(id, { includeDeleted: true }); + private getSoftDeletedWriteConflict(id: string, task: Task, existingRow?: TaskRow): string | undefined { + const existing = existingRow ?? this.readTaskRowFromDb(id, { includeDeleted: true }); if (!existing?.deletedAt || task.deletedAt !== undefined) { return undefined; } @@ -3169,19 +3166,18 @@ export class TaskStore extends EventEmitter { */ private async atomicWriteTaskJson(dir: string, task: Task): Promise { const id = this.getTaskIdFromDir(dir); - let deletedAt: string | undefined; + let result: { deletedAt?: string; current?: Task } | undefined; this.db.transactionImmediate(() => { - // Soft-delete/restore state is written via direct SQL paths (deleteTask and - // future restore flows), so stale task.json upserts must never clear deletedAt. - deletedAt = this.getSoftDeletedWriteConflict(id, task); - if (deletedAt) return; - this.upsertTaskWithFtsRecovery(task); + const existingRow = this.readTaskRowFromDb(id, { includeDeleted: true }); + const changedColumns = existingRow && existingRow.deletedAt == null + ? this.getChangedTaskColumns(existingRow, task) + : new Set(); + result = this.patchTaskRowInTransaction(id, task, changedColumns, existingRow); }); - if (deletedAt) { - this.throwSoftDeletedWriteBlocked(id, deletedAt, "atomicWriteTaskJson"); + if (result?.deletedAt) { + this.throwSoftDeletedWriteBlocked(id, result.deletedAt, "atomicWriteTaskJson"); } - // Also write to disk for backward compatibility - await this.writeTaskJsonFile(dir, task); + await this.writeTaskJsonFile(dir, result?.current ?? task); } /** @@ -3198,29 +3194,28 @@ export class TaskStore extends EventEmitter { auditInput?: RunAuditEventInput, ): Promise { const id = this.getTaskIdFromDir(dir); - let deletedAt: string | undefined; + let result: { deletedAt?: string; current?: Task } | undefined; this.db.transactionImmediate(() => { - deletedAt = this.getSoftDeletedWriteConflict(id, task); - if (deletedAt) return; + const existingRow = this.readTaskRowFromDb(id, { includeDeleted: true }); + const changedColumns = existingRow && existingRow.deletedAt == null + ? this.getChangedTaskColumns(existingRow, task) + : new Set(); + result = this.patchTaskRowInTransaction(id, task, changedColumns, existingRow); + if (result?.deletedAt) return; - // Upsert the task - this.upsertTaskWithFtsRecovery(task); - - // Optionally record the audit event in the same transaction if (auditInput) { this.insertRunAuditEventRow(auditInput); } }); - if (deletedAt) { - this.throwSoftDeletedWriteBlocked(id, deletedAt, auditInput?.mutationType ?? "atomicWriteTaskJsonWithAudit", { + if (result?.deletedAt) { + this.throwSoftDeletedWriteBlocked(id, result.deletedAt, auditInput?.mutationType ?? "atomicWriteTaskJsonWithAudit", { agentId: auditInput?.agentId, runId: auditInput?.runId, timestamp: auditInput?.timestamp, }); } - // File writes are not part of the SQLite transaction - await this.writeTaskJsonFile(dir, task); + await this.writeTaskJsonFile(dir, result?.current ?? task); } /** @@ -6163,6 +6158,55 @@ export class TaskStore extends EventEmitter { return { ok: true, task: post }; } + async renewCheckoutLease( + taskId: string, + update: { + checkoutRunId: string | null; + checkoutLeaseRenewedAt: string; + }, + ): Promise { + const dir = this.taskDir(taskId); + let deletedAt: string | undefined; + let current: Task | undefined; + this.db.transactionImmediate(() => { + const row = this.readTaskRowFromDb(taskId, { includeDeleted: true }); + if (row?.deletedAt) { + deletedAt = row.deletedAt; + return; + } + + const result = this.db.prepare(` + UPDATE tasks + SET checkoutRunId = ?, checkoutLeaseRenewedAt = ?, updatedAt = ? + WHERE id = ? AND ${TaskStore.ACTIVE_TASKS_WHERE} + `).run(update.checkoutRunId, update.checkoutLeaseRenewedAt, update.checkoutLeaseRenewedAt, taskId) as { changes: number }; + + if (result.changes === 0) { + return; + } + + this.db.bumpLastModified(); + current = this.readTaskFromDb(taskId); + }); + + if (deletedAt) { + this.throwSoftDeletedWriteBlocked(taskId, deletedAt, "renewCheckoutLease", { + timestamp: update.checkoutLeaseRenewedAt, + }); + } + + if (!current) { + throw new Error(`Task ${taskId} not found`); + } + + await this.writeTaskJsonFile(dir, current); + if (this.isWatching) { + this.taskCache.set(taskId, { ...current }); + } + this.emitTaskLifecycleEventSafely("task:updated", [current]); + return current; + } + async selectNextTaskForAgent( agentId: string, agent?: Pick, diff --git a/packages/engine/src/__tests__/executor-lease-renewal.test.ts b/packages/engine/src/__tests__/executor-lease-renewal.test.ts new file mode 100644 index 0000000000..c48ad319e3 --- /dev/null +++ b/packages/engine/src/__tests__/executor-lease-renewal.test.ts @@ -0,0 +1,66 @@ +import "./executor-test-helpers.js"; +import { beforeEach, describe, expect, it, vi } from "vitest"; + +import { TaskExecutor } from "../executor.js"; +import { resetExecutorMocks } from "./executor-test-helpers.js"; + +function createStore() { + const listeners = new Map void)[]>(); + return { + on: vi.fn((event: string, listener: (payload: unknown) => void) => { + const existing = listeners.get(event) ?? []; + existing.push(listener); + listeners.set(event, existing); + }), + off: vi.fn(), + getSettings: vi.fn().mockResolvedValue({ globalPause: false, enginePaused: false }), + listTasks: vi.fn().mockResolvedValue([]), + renewCheckoutLease: vi.fn().mockResolvedValue(undefined), + updateTask: vi.fn().mockResolvedValue(undefined), + } as any; +} + +describe("TaskExecutor lease renewal fallback", () => { + beforeEach(() => { + resetExecutorMocks(); + }); + + it("uses renewCheckoutLease instead of updateTask when agentStore is unavailable", async () => { + const store = createStore(); + const executor = new TaskExecutor(store, "/tmp/test"); + + await (executor as any).renewTaskLease("FN-5945", "agent-1", 3, "node-1", "run-1"); + + expect(store.renewCheckoutLease).toHaveBeenCalledTimes(1); + expect(store.renewCheckoutLease).toHaveBeenCalledWith( + "FN-5945", + expect.objectContaining({ + checkoutRunId: "run-1", + checkoutLeaseRenewedAt: expect.any(String), + }), + ); + expect(store.updateTask).not.toHaveBeenCalled(); + }); + + it("keeps the agentStore checkoutTask renewal path unchanged", async () => { + const store = createStore(); + const agentStore = { checkoutTask: vi.fn().mockResolvedValue(undefined) } as any; + const executor = new TaskExecutor(store, "/tmp/test", { agentStore }); + + await (executor as any).renewTaskLease("FN-5945", "agent-1", 4, "node-2", "run-2"); + + expect(agentStore.checkoutTask).toHaveBeenCalledWith( + "agent-1", + "FN-5945", + expect.objectContaining({ + nodeId: "node-2", + runId: "run-2", + leaseEpoch: 4, + renewedAt: expect.any(String), + }), + undefined, + ); + expect(store.renewCheckoutLease).not.toHaveBeenCalled(); + expect(store.updateTask).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/engine/src/executor.ts b/packages/engine/src/executor.ts index ae5675f7c0..933f2c9847 100644 --- a/packages/engine/src/executor.ts +++ b/packages/engine/src/executor.ts @@ -1403,7 +1403,7 @@ export class TaskExecutor { ); return; } - await this.store.updateTask(taskId, { + await this.store.renewCheckoutLease(taskId, { checkoutRunId: runId ?? null, checkoutLeaseRenewedAt: renewedAt, });