fix(heartbeat): claim task ownership with queued runs (#13973)

## Thinking Path

> - Paperclip manages agents and their tasks.
> - The scheduler can start several runs for one agent.
> - Each task still needs one execution owner.
> - Task creation and recovery can queue assignments for the same task.
> - The scheduler previously marked both runs running before starting
either executor.
> - This change claims the run and task ownership in one transaction.
> - Tasks with different IDs retain concurrent execution.

## Linked Issues or Issue Description

Refs #13882. Related queue work: #11830 and #13965; those address
different admission and cancellation paths.

**What happened?**

The Grok subscription Product campaign exposed two assignment runs for
one task. Task creation and periodic recovery started runs within 200
ms. Both acquired Daytona sandboxes. One failed before provider work
with `paperclip_runner_attachment_staging_not_authorized` because the
other owned the task. Fixture cleanup then found an active session. The
original failed result is retained in [campaign
36071063537](https://github.com/paperclipai/paperclip/actions/runs/36071063537).

**Expected behavior**

Only the task execution owner may start provider setup. Competing queued
work must wait. Different tasks may use the agent's available slots.

**Steps to reproduce**

1. Set an agent's concurrent run limit to two.
2. Queue task-creation and recovery assignment runs for the same task.
3. Resume the queue while holding adapter execution open.
4. Before this fix, both runs become running. Only one has the task
execution lock.

**Paperclip version or commit**

Reproduced on master `8781f06a8` and in the Grok campaign at `4196a4cd`.

## What Changed

- Use the same company-scoped task ownership gate for assignments,
direct comments, and queued comments.
- Leave competing work queued while another live run owns the task.
Transfer a terminal pointer only after the tracked executor and durable
environment leases/finalization settle, including cleanup on another
controller.
- Commit the running state and task execution owner together.
- Preserve the dedicated review path and the native replacement checkout
guard.
- Add sixteen database regressions for assignment/comment orderings,
unrelated concurrent tasks, and cross-controller cleanup fences.

## Verification

- Live subscription Product E2E planning passed on its first attempt in
both Daytona (351,197 ms) and local execution (268,883 ms), with 6/6
matchers and cleanup passing in each, on combined source
`1b0551bb7c8de3c54f4bee64dbe2c88328b3645e`
([campaign](https://github.com/paperclipai/paperclip/actions/runs/36080870743)).
The local screenshots and persisted state show the same plan revised,
the first approval rejected, the revised approval accepted, and one
completion. The full campaign remains unqualified because of a separate
startup retry and EC2 Spot interruption.
- The new same-task regression failed before the fix: two running rows
instead of one.
- All seven focused concurrency cases passed. Three comment cases failed
before the shared gate was added. Five cross-controller cleanup cases
failed before the durable gate was added. A warm-retention regression
also failed before its successful release receipt was admitted;
missing/failed receipts and a different retention policy stay blocked.
- `node ../node_modules/vitest/vitest.mjs run
src/__tests__/heartbeat-stale-queue-invalidation.test.ts` from `server`:
48 passed. Native cleanup admission adds one passing test; task-drain
release adds two. Total focused checks: 51 passed.
- `git diff --check`: passed.
- Current-head
[CI](https://github.com/paperclipai/paperclip/actions/runs/36078495577)
passes repository typecheck, tests, Rust checks, build, and browser
suites. The unrelated repository-form browser shard passed one bounded
rerun without assertion or code changes; its original navigation failure
remains retained. No local Docker or Rust build was used.
- The failed test's exact two sandboxes were checked: one was already
absent; the other was deleted after its company, environment, and run
labels matched.

## Risks

The claim transaction adds a task row lock. It follows task-before-run
lock order. A competing run stays queued until the current owner
releases the task. The change does not alter attachment authorization,
provider permissions, schemas, or public APIs. Review is 5/5 with all
threads resolved on `23c0d13b6`. All CI checks pass on that exact head.
Live Grok requalification will use the combined runner and scheduler
source.

## Model Used

OpenAI GPT-6 through Codex, with repository tools and code execution.
The exact serving identifier and context-window size are not exposed in
this session.

## Checklist

- [x] I have included a thinking path that traces from project context
to this change
- [x] I have specified the model used (with version and capability
details)
- [x] I have checked ROADMAP.md and confirmed this PR does not duplicate
planned core work
- [x] I have searched GitHub for duplicate or related PRs and linked
them above
- [x] I have either (a) linked existing issues with `Fixes: #` / `Closes
#` / `Refs #` OR (b) described the issue in-PR following the relevant
issue template
- [x] I have not referenced internal/instance-local Paperclip issues or
links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip`
URLs)
- [x] My branch name describes the change (e.g. `docs/...`, `fix/...`)
and contains no internal Paperclip ticket id or instance-derived details
- [x] I have run tests locally and they pass
- [x] I have added or updated tests where applicable
- [x] I have updated relevant documentation to reflect my changes
- [x] I have considered and documented any risks above
- [x] All Paperclip CI gates are green
- [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups
- [x] I will address all Greptile and reviewer comments before
requesting merge

### Live regression verification

The Daytona plan/revise/accept lifecycle in [Grok campaign
36080870743](https://github.com/paperclipai/paperclip/actions/runs/36080870743)
passed on its first attempt at combined source
`1b0551bb7c8de3c54f4bee64dbe2c88328b3645e`, in 351,197 ms, with all six
outcome matchers and cleanup passing. The earlier competing-run failure
remains retained. This is one live planning result; the complete Grok
matrix and repeated qualification are separate gates, and this campaign
also has separately retained infrastructure failures.

---------

Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
DottaandPaperclip authored and GitHub committed 2026-09-25 10:06:05 -05:00
1 parent efce9356b5
commit c96c751f39
3 files changed
+187 -28

No files matched your search

@@ -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<void>((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: {
@@ -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<number, unknown>) {
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<unknown>) => {
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;
}
});
+76 -13
View File
@@ -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;