diff --git a/server/src/__tests__/heartbeat-stale-queue-invalidation.test.ts b/server/src/__tests__/heartbeat-stale-queue-invalidation.test.ts index fc71329a30..882a81e5d2 100644 --- a/server/src/__tests__/heartbeat-stale-queue-invalidation.test.ts +++ b/server/src/__tests__/heartbeat-stale-queue-invalidation.test.ts @@ -9,6 +9,8 @@ import { createDb, documentRevisions, documents, + environmentLeases, + nativeRunFinalizations, heartbeatRuns, issueComments, issueDocuments, @@ -27,6 +29,7 @@ import { } from "../services/heartbeat.ts"; import { runningProcesses } from "../adapters/index.ts"; import { recoveryService } from "../services/recovery/service.ts"; +import { withQueuedCommentIdsInWakePayload } from "../services/issue-queued-comment-queue.ts"; const mockAdapterExecute = vi.hoisted(() => vi.fn(async () => ({ @@ -311,6 +314,105 @@ describeEmbeddedPostgres("heartbeat stale queued-run invalidation", () => { }); } + it.each([ + { name: "assignments", sameIssue: true, first: "assignment", second: "assignment" }, + { name: "unrelated tasks", sameIssue: false, first: "assignment", second: "assignment" }, + { name: "assignment then direct comment", sameIssue: true, first: "assignment", second: "direct" }, + { name: "direct comment then assignment", sameIssue: true, first: "direct", second: "assignment" }, + { name: "assignment then queued comment", sameIssue: true, first: "assignment", second: "queued" }, + { name: "queued comment then assignment", sameIssue: true, first: "queued", second: "assignment" }, + { name: "queued comments", sameIssue: true, first: "queued", second: "queued" }, + ])("serializes issue claims without serializing unrelated work: $name", async ({ sameIssue, first, second }) => { + const { companyId, agentId } = await seedCompanyAndAgent({ maxConcurrentRuns: 2 }); + const firstIssueId = randomUUID(); + const secondIssueId = sameIssue ? firstIssueId : randomUUID(); + for (const id of new Set([firstIssueId, secondIssueId])) { + await db.insert(issues).values({ id, companyId, title: "Concurrent assignment", status: "todo", assigneeAgentId: agentId }); + } + async function queue(issueId: string, kind: string) { + const commentId = kind === "assignment" ? null : randomUUID(); + if (commentId) await db.insert(issueComments).values({ id: commentId, companyId, issueId, authorUserId: "local-board", body: "Continue this task." }); + const queued = await seedQueuedRun({ companyId, agentId, issueId, + wakeReason: commentId ? "issue_commented" : "issue_assigned", + invocationSource: commentId ? "automation" : "assignment", + contextExtras: commentId ? { wakeCommentIds: [commentId], wakeCommentId: commentId } : {}, + }); + if (kind === "queued") await db.update(agentWakeupRequests).set({ + payload: withQueuedCommentIdsInWakePayload({ issueId }, [commentId!]), + }).where(eq(agentWakeupRequests.id, queued.wakeupRequestId)); + } + await queue(firstIssueId, first); + await queue(secondIssueId, second); + let release!: () => void; + const gate = new Promise((resolve) => { release = resolve; }); + mockAdapterExecute.mockImplementation(async () => { + await gate; + return { exitCode: 0, signal: null, timedOut: false, errorMessage: null, summary: "Assignment finished.", provider: "test", model: "test-model" }; + }); + try { + await heartbeat.resumeQueuedRuns(); + expect(await waitForCondition(() => Promise.resolve(mockAdapterExecute.mock.calls.length > 0))).toBe(true); + const runs = await db.select().from(heartbeatRuns); + const running = runs.filter((run) => run.status === "running"); + expect(running).toHaveLength(sameIssue ? 1 : 2); + expect(runs.filter((run) => run.status === "queued")).toHaveLength(sameIssue ? 1 : 0); + for (const run of running) { + const [issue] = await db.select().from(issues).where(eq(issues.id, (run.contextSnapshot as { issueId: string }).issueId)); + expect(issue.executionRunId).toBe(run.id); + } + } finally { + release(); + await heartbeat.drainActiveRunExecutions(); + } + }); + + it.each(["active_lease", "pending_cleanup", "failed_cleanup", "finalizer_lease", "workspace_finalization", + "retained_ready", "retained_missing_receipt", "retained_failed", "retained_wrong_policy"])( + "waits for durable cleanup from another controller: %s", async (pending) => { + const { companyId, agentId } = await seedCompanyAndAgent(); + const issueId = randomUUID(), previousId = randomUUID(); + await db.insert(issues).values({ id: issueId, companyId, title: "Remote cleanup", status: "todo", assigneeAgentId: agentId }); + await db.insert(heartbeatRuns).values({ id: previousId, companyId, agentId, + status: "succeeded", runtimeMode: "native", nativeIssueId: issueId, + invocationSource: "assignment", contextSnapshot: { issueId }, finishedAt: new Date() }); + await db.update(issues).set({ executionRunId: previousId }).where(eq(issues.id, issueId)); + const leaseId = randomUUID(); + const retained = pending.startsWith("retained_"); + const hasLease = retained || ["active_lease", "pending_cleanup", "failed_cleanup"].includes(pending); + if (hasLease) await db.insert(environmentLeases).values({ id: leaseId, companyId, issueId, + heartbeatRunId: previousId, provider: "daytona", + status: retained ? "retained" : pending === "pending_cleanup" ? "pending_cleanup" : "released", + leasePolicy: pending === "retained_wrong_policy" ? "retain_on_failure" : retained ? "reuse_by_environment" : "ephemeral", + releasedAt: retained || pending === "active_lease" ? null : new Date(), + cleanupStatus: ["retained_ready", "retained_wrong_policy"].includes(pending) ? "success" + : ["failed_cleanup", "retained_failed"].includes(pending) ? "failed" : null }); + await db.insert(nativeRunFinalizations).values({ runId: previousId, companyId, issueId, + phase: pending === "workspace_finalization" ? "workspace_finalizing" : "committed", + leaseOwner: pending === "finalizer_lease" ? "another-controller" : null }); + const next = await seedQueuedRun({ companyId, agentId, issueId, wakeReason: "issue_assigned", invocationSource: "assignment" }); + mockAdapterExecute.mockImplementation(async () => { + await db.update(issues).set({ status: "done" }).where(eq(issues.id, issueId)); + return { exitCode: 0, signal: null, timedOut: false, errorMessage: null, + summary: "Task completed", provider: "test", model: "test-model" }; + }); + // No in-memory executor exists in this service. Durable ownership alone + // must fence a second controller until cleanup settles. A verified warm + // retention receipt is already a settled boundary and allows continuation. + await heartbeat.resumeQueuedRuns(); + if (pending !== "retained_ready") { + expect((await heartbeat.getRun(next.runId))?.status).toBe("queued"); + expect(mockAdapterExecute).not.toHaveBeenCalled(); + expect((await db.select().from(issues).where(eq(issues.id, issueId)))[0].executionRunId).toBe(previousId); + if (hasLease) await db.update(environmentLeases).set({ status: "released", releasedAt: new Date(), cleanupStatus: "success" }).where(eq(environmentLeases.id, leaseId)); + await db.update(nativeRunFinalizations).set({ phase: "committed", leaseOwner: null }).where(eq(nativeRunFinalizations.runId, previousId)); + await heartbeat.resumeQueuedRuns(); + } + await heartbeat.drainActiveRunExecutions(); + expect((await heartbeat.getRun(next.runId))?.status).toBe("succeeded"); + expect(mockAdapterExecute).toHaveBeenCalledTimes(1); + }, + ); + it("skips generic timer wakes with no actionable assigned work before adapter execution", async () => { const { agentId } = await seedCompanyAndAgent({ heartbeatConfig: { diff --git a/server/src/__tests__/heartbeat-task-drain-admission-release.test.ts b/server/src/__tests__/heartbeat-task-drain-admission-release.test.ts index a18cc2ad58..32ddc0de94 100644 --- a/server/src/__tests__/heartbeat-task-drain-admission-release.test.ts +++ b/server/src/__tests__/heartbeat-task-drain-admission-release.test.ts @@ -219,29 +219,21 @@ describeEmbeddedPostgres("heartbeat task-drain admission release", () => { expect(finished?.status).toBe("succeeded"); }, 20_000); - // Wraps db.transaction so the callback's tx object throws the moment code - // calls tx.update(table) for a table named in tablesByCall — this makes a - // real Postgres transaction roll back exactly like a genuine write failure - // partway through, without touching any other table's update path. - // tablesByCall maps a 0-based db.transaction() call index (in call order) - // to the table that call should fail on; a call index with no entry runs - // every update for real. For example { 1: issues } lets the atomic stale-run - // validation transaction complete, then fails only the issue-lock write - // inside releaseRunClaimedJustBeforeSuppression's transaction. - function withFailingTransactionalUpdate(realDb: typeof db, tablesByCall: Record) { - let callIndex = 0; + // Inject one rollback only after the test starts task drain. Admission also + // uses transactions, so transaction call numbers do not identify release. + function withFailingTransactionalUpdate(realDb: typeof db, failingTable: unknown, armed: () => boolean) { + let injected = false; return new Proxy(realDb, { get(target, prop, receiver) { if (prop !== "transaction") return Reflect.get(target, prop, receiver); return (fn: (tx: unknown) => Promise) => { - const failingTable = tablesByCall[callIndex]; - callIndex += 1; return target.transaction((tx) => { const txProxy = new Proxy(tx as object, { get(txTarget, txProp, txReceiver) { if (txProp === "update") { return (table: unknown) => { - if (failingTable !== undefined && table === failingTable) { + if (!injected && armed() && table === failingTable) { + injected = true; throw new Error("simulated transactional write failure"); } return (txTarget as any).update(table); @@ -262,13 +254,15 @@ describeEmbeddedPostgres("heartbeat task-drain admission release", () => { // Fault the release transaction on the issue-lock write, so executeRun's // suppression branch catches the failure, logs it, and returns instead // of throwing. There is no in-process fallback or retry for this path. - const failingDb = withFailingTransactionalUpdate(db, { 1: issues }); + let releaseStarted = false; + const failingDb = withFailingTransactionalUpdate(db, issues, () => releaseStarted); const heartbeat = heartbeatService(failingDb); const unsubscribe = subscribeCompanyLiveEvents(companyId, (event) => { const payload = event.payload as { runId?: string; status?: string }; if (event.type === "heartbeat.run.status" && payload.runId === runId && payload.status === "running") { startTaskDrain({}); + releaseStarted = true; } }); diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index dca86e4abc..abf7e8a342 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -17296,6 +17296,69 @@ export function heartbeatService( responsibleUserId: null, }, }); + // All ordinary and comment claims use the same company-scoped issue + // lock. A batch may claim several runs before executeRun tracks any owner. + async function lockIssueExecutionClaim(tx: Db) { + const [owner] = issueId ? await tx.select({ + assigneeAgentId: issues.assigneeAgentId, + executionRunId: issues.executionRunId, + checkoutRunId: issues.checkoutRunId, + }).from(issues).where(and( + eq(issues.id, issueId), eq(issues.companyId, run.companyId), + )).for("update") : []; + const ownsIssue = owner?.assigneeAgentId === run.agentId && + context.wakeReason !== "source_scoped_recovery_action"; + if (ownsIssue && run.scheduledRetryReason === "native_safe_replacement" && + owner.checkoutRunId && owner.checkoutRunId !== run.id) { + return { ownsIssue, blocked: true }; + } + if (ownsIssue && owner.executionRunId && owner.executionRunId !== run.id) { + const [previous] = await tx.select({ status: heartbeatRuns.status }) + .from(heartbeatRuns).where(and( + eq(heartbeatRuns.id, owner.executionRunId), + eq(heartbeatRuns.companyId, run.companyId), + )); + // A terminal result can precede workspace/lease cleanup on this or + // another controller. Local absence alone is not a release receipt. + if (!isHeartbeatRunTerminalStatus(previous?.status) || + liveRunExecutions.has(owner.executionRunId)) { + return { ownsIssue, blocked: true }; + } + const [pendingLease] = await tx.select({ id: environmentLeases.id }) + .from(environmentLeases).where(and( + eq(environmentLeases.companyId, run.companyId), + eq(environmentLeases.heartbeatRunId, owner.executionRunId), + or(and(isNull(environmentLeases.releasedAt), + // Warm release deliberately retains the sandbox. Its successful + // receipt settles the old run without destroying the resource. + sql`not coalesce(${environmentLeases.status} = 'retained' + and ${environmentLeases.leasePolicy} = 'reuse_by_environment' + and ${environmentLeases.cleanupStatus} = 'success', false)`), + eq(environmentLeases.status, "pending_cleanup"), + eq(environmentLeases.cleanupStatus, "failed")), + )).limit(1); + const [finalization] = await tx.select({ + phase: nativeRunFinalizations.phase, leaseOwner: nativeRunFinalizations.leaseOwner, + }).from(nativeRunFinalizations).where(and( + eq(nativeRunFinalizations.companyId, run.companyId), + eq(nativeRunFinalizations.runId, owner.executionRunId), + )); + if (pendingLease || (finalization && (finalization.leaseOwner || + !["committed", "applied", "terminal_failure"].includes(finalization.phase)))) { + return { ownsIssue, blocked: true }; + } + } + return { ownsIssue, blocked: false }; + } + async function bindClaimedIssueExecution(tx: Db, ownsIssue: boolean, claimedRun: typeof heartbeatRuns.$inferSelect | null | undefined) { + if (!claimedRun || !issueId || !ownsIssue) return; + await tx.update(issues).set({ + executionRunId: claimedRun.id, + executionAgentNameKey: normalizeAgentNameKey(agent.name), + executionLockedAt: claimedAt, + updatedAt: claimedAt, + }).where(and(eq(issues.id, issueId), eq(issues.companyId, run.companyId))); + } const nativeReviewContext = readNativeReviewAssignmentContext(context); const queuedCommentIds = queuedCommentIdsFromRunContext(context); if ( @@ -17316,16 +17379,8 @@ export function heartbeatService( // run becomes running, a concurrent discard must observe the // claimed wake and return an explicit conflict; if discard wins, // this claim observes the cancelled queue and does no work. - await tx - .select({ id: issues.id }) - .from(issues) - .where( - and( - eq(issues.id, issueId), - eq(issues.companyId, run.companyId), - ), - ) - .for("update"); + const issueClaim = await lockIssueExecutionClaim(tx as unknown as Db); + if (issueClaim.blocked) return { kind: "stale" as const, run: null }; const wake = await tx .select() .from(agentWakeupRequests) @@ -17476,6 +17531,7 @@ export function heartbeatService( ), ) .returning(); + await bindClaimedIssueExecution(tx as unknown as Db, issueClaim.ownsIssue, claimedRun); return claimedRun ? { kind: "claimed" as const, run: claimedRun } : { kind: "stale" as const, run: null }; @@ -17578,6 +17634,7 @@ export function heartbeatService( ), ) .returning(); + await bindClaimedIssueExecution(tx as unknown as Db, issueClaim.ownsIssue, claimedRun); return claimedRun ? { kind: "claimed" as const, run: claimedRun } : { kind: "stale" as const, run: null }; @@ -17639,9 +17696,15 @@ export function heartbeatService( agentNameKey: normalizeAgentNameKey(agent.name), }); } - return tx.update(heartbeatRuns).set(claimValues).where(and( - eq(heartbeatRuns.id, run.id), eq(heartbeatRuns.status, "queued"), - )).returning().then((rows) => rows[0] ?? null); + return tx.transaction(async (claimTx) => { + const issueClaim = await lockIssueExecutionClaim(claimTx as unknown as Db); + if (issueClaim.blocked) return null; + const claimedRun = await claimTx.update(heartbeatRuns).set(claimValues).where(and( + eq(heartbeatRuns.id, run.id), eq(heartbeatRuns.status, "queued"), + )).returning().then((rows) => rows[0] ?? null); + await bindClaimedIssueExecution(claimTx as unknown as Db, issueClaim.ownsIssue, claimedRun); + return claimedRun; + }); }); if (!claimed) return null;