mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-02 02:07:25 +08:00
fix(heartbeat): cancel obsolete execution continuations before dispatch (#13761)
## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - A comment can queue execution before the task changes state. > - The task may be completed, cancelled, deleted, or reassigned before setup checks ownership. > - The ownership guard must prevent that obsolete execution from starting. > - This pull request records that guard outcome as cancellation and settles the wake request. > - Missing history and authorization errors remain failures that operators can investigate. ## Linked Issues or Issue Description **What happened?** A comment can queue a run just before a user marks its task Done. The ownership guard stops setup before adapter dispatch, but the run becomes a setup failure and can put the agent into an error state. **Expected behavior** Cancel obsolete work when the typed task ownership guard rejects it. Keep genuine setup errors visible. Do not change the task's terminal state or current owner. **Steps to reproduce** Queue a task continuation, then mark the task Done or Cancelled before continuation setup. The run should settle as cancelled without invoking the adapter. Deleting or reassigning the task must also prevent dispatch. Missing source history or explicit user authorization must still fail. **Paperclip version or commit** Refreshed against master `b2e9e82f053af14777fd68e8834fff6bb64c84c4`. Related: #13546 covers lost issue-lock claims, and #13888 covers persisted continuation decisions and retry budgets. This change handles the final continuation ownership guard during setup. Thanks to @MrBlackTongue for the original typed cancellation implementation and regression coverage. This update preserves the contributor's commits, resolves the master conflict, and narrows cancellation to task ownership invalidation. ## What Changed - Add a typed `StaleExecutionContinuationError` for `continuation_task_ownership_changed`. - Use existing cancellation settlement for that typed guard, including the run, wake request, issue execution ownership, and agent state. Suppress immediate recovery of obsolete work. - Preserve failure classification for missing source context, missing user authorization, and untyped errors, including an untyped error with identical text. - Cover Done and Cancelled tasks with real database checks. Retain company and assignment guard coverage. Verify cancellation performs no adapter dispatch or automatic replay. - Document the cancellation event and its error classification in the run-log guide. ## Verification Current head: `b875486e81`. - `pnpm -r typecheck`: passed. - `pnpm build`: passed. - `pnpm exec vitest run server/src/__tests__/heartbeat-process-recovery.test.ts server/src/services/execution-continuation.test.ts --maxWorkers=1 --no-file-parallelism -t 'authorized continuation context|obsolete continuation setup|untyped continuation setup failures'`: 17 passed; 316 unrelated tests excluded by the name filter. - `pnpm test:run`: local full run is still finishing. It is not a clean pass: the observed failures are three company-skills cache checks (`EACCES` while renaming read-only cache directories on macOS), one tool-access assertion, and one comment-wake timeout. The three cache failures reproduce in isolation in unchanged code. The tool-access check passes in isolation, and the entire comment-wake file passes (29 tests), including explicit feedback after completion. Final full-run totals will be added when available. - [CI](https://github.com/paperclipai/paperclip/actions/runs/36282532652): all gates green on this head. An unchanged Runner workspace-diff test initially returned no diff; an isolated check passed, and the failed job plus its dependent gate passed on retry with no code change. Original failed job: 108517010989. - Greptile: 5/5 on the full current head, with no unresolved review threads. Branch is current with master and conflict-free. - Diff check and local secret/PII scan passed. No customer data or private deployment identifiers are included. ## Risks The classification change is limited to the typed ownership guard. A task with missing history remains a failure rather than being treated as an expected cancellation. Existing company, task-owner, terminal-state, and authorization checks still prevent dispatch. Cancellation retains its specific reason in the run log and does not grant a retry. No schema, dependency, or API changes; no migration is required. ## Model Used OpenAI GPT-6 via Codex assisted investigation, implementation, code review, and test execution. The exact serving model 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 - [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 described the issue in-PR following the bug issue template - [x] I have not referenced internal/instance-local Paperclip issues or links - [x] My branch name describes the change and contains no internal ticket id - [x] I have run the focused regression tests locally and they pass; full-suite results are tracked above - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation - [x] I have considered and documented the risks - [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 merge Co-authored-by: Paperclip <noreply@paperclip.ing> Co-authored-by: Devin Foley <devin@paperclip.ing>
This commit is contained in:
co-authored by
Paperclip
Devin Foley
parent
b2e9e82f05
commit
a6c4e7a8d1
@@ -159,6 +159,15 @@ Provider identity diagnostics remain in the local run log. They record the notif
|
||||
|
||||
Recovery lifecycle events retain the original structured failure code, retry attempt, next retry time, and predecessor/successor identifiers. Durable status delivery uses an idempotency marker; delivery grants no provider authority. Failed publication is retried without repeating provider work. These records are not first-party Telemetry.
|
||||
|
||||
If execution-continuation setup finds that a task no longer exists, is closed,
|
||||
or its owner changed, the existing cancellation settlement records
|
||||
`continuation_task_ownership_changed`. The run and wake request become cancelled
|
||||
before adapter dispatch, with the run-log message
|
||||
`stale execution continuation cancelled before dispatch`. Immediate recovery is
|
||||
suppressed. Missing source context, authorization failures, and other setup
|
||||
errors retain their failure classification. An untyped error with the same
|
||||
message is also still a failure; cancellation requires the typed ownership guard.
|
||||
|
||||
Bounded retry exhaustion writes one lifecycle receipt per run, retry reason,
|
||||
scheduled attempt, and retry limit. Repeated or concurrent recovery checks reuse
|
||||
that receipt, including receipts from earlier builds, without advancing the event
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import * as executionContinuation from "../services/execution-continuation.js";
|
||||
import { legacyDispositionFingerprint, LEGACY_DISPOSITION_REPAIR_INSTRUCTION } from "../services/recovery/legacy-continuation.js";
|
||||
import * as controllerLeases from "../services/legacy-controller-lease.js";
|
||||
import { instanceSettingsService } from "../services/instance-settings.js";
|
||||
@@ -1528,6 +1529,62 @@ describeEmbeddedPostgres("heartbeat orphaned process recovery", () => {
|
||||
expect(missingCommentWakeups).toHaveLength(0);
|
||||
});
|
||||
|
||||
it.each([
|
||||
"continuation_task_ownership_changed",
|
||||
] as const)("cancels obsolete continuation setup: %s", async (code) => {
|
||||
const { agentId, runId, wakeupRequestId } = await seedQueuedIssueRunFixture();
|
||||
const build = vi.spyOn(executionContinuation, "buildExecutionContinuation")
|
||||
.mockRejectedValueOnce(new executionContinuation.StaleExecutionContinuationError(code));
|
||||
try {
|
||||
const heartbeat = heartbeatService(db);
|
||||
await heartbeat.resumeQueuedRuns();
|
||||
await waitForRunToSettle(heartbeat, runId);
|
||||
await heartbeat.waitForRunExecutionDrain(runId);
|
||||
expect(build).toHaveBeenCalled();
|
||||
expect(await heartbeat.getRun(runId)).toMatchObject({ status: "cancelled", errorCode: code });
|
||||
const [wakeup] = await db.select().from(agentWakeupRequests)
|
||||
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
||||
expect(wakeup.status).toBe("cancelled");
|
||||
const [agent] = await db.select().from(agents).where(eq(agents.id, agentId));
|
||||
expect(agent.status).not.toBe("error");
|
||||
expect(mockAdapterExecute.mock.calls.some(
|
||||
([input]) => (input as { runId?: string } | undefined)?.runId === runId,
|
||||
)).toBe(false);
|
||||
const runs = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.agentId, agentId));
|
||||
expect(runs).toHaveLength(1);
|
||||
} finally {
|
||||
build.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
it.each([
|
||||
"continuation_source_context_missing",
|
||||
"continuation_user_authorization_missing",
|
||||
"continuation_task_ownership_changed",
|
||||
])("retains untyped continuation setup failures: %s", async (message) => {
|
||||
const { runId, wakeupRequestId } = await seedQueuedIssueRunFixture();
|
||||
const build = vi.spyOn(executionContinuation, "buildExecutionContinuation")
|
||||
.mockRejectedValueOnce(new Error(message));
|
||||
try {
|
||||
const heartbeat = heartbeatService(db);
|
||||
await heartbeat.resumeQueuedRuns();
|
||||
await waitForRunToSettle(heartbeat, runId);
|
||||
await heartbeat.waitForRunExecutionDrain(runId);
|
||||
expect(build).toHaveBeenCalled();
|
||||
expect(await heartbeat.getRun(runId)).toMatchObject({
|
||||
status: "failed", errorCode: "setup_failed", error: message,
|
||||
});
|
||||
const [wakeup] = await db.select().from(agentWakeupRequests)
|
||||
.where(eq(agentWakeupRequests.id, wakeupRequestId));
|
||||
expect(wakeup.status).toBe("failed");
|
||||
expect(mockAdapterExecute.mock.calls.some(
|
||||
([input]) => (input as { runId?: string } | undefined)?.runId === runId,
|
||||
)).toBe(false);
|
||||
} finally {
|
||||
build.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
it("does not immediately continue a low-trust preflight setup failure", async () => {
|
||||
const { agentId, runId, issueId, companyId } =
|
||||
await seedQueuedIssueRunFixture();
|
||||
|
||||
@@ -15,7 +15,20 @@ import {
|
||||
getEmbeddedPostgresTestSupport,
|
||||
startEmbeddedPostgresTestDatabase,
|
||||
} from "../__tests__/helpers/embedded-postgres.js";
|
||||
import { buildExecutionContinuation, currentContinuationOrigins, projectHumanInteractionResponse } from "./execution-continuation.js";
|
||||
import { StaleExecutionContinuationError, buildExecutionContinuation, currentContinuationOrigins, projectHumanInteractionResponse } from "./execution-continuation.js";
|
||||
|
||||
const expectStaleContinuation = async (
|
||||
run: () => Promise<unknown>,
|
||||
code: StaleExecutionContinuationError["code"],
|
||||
) => {
|
||||
await expect(run()).rejects.toThrow(StaleExecutionContinuationError);
|
||||
await expect(run()).rejects.toMatchObject({ code });
|
||||
};
|
||||
const expectMissingContinuationContext = async (run: () => Promise<unknown>) => {
|
||||
const result = run();
|
||||
await expect(result).rejects.toThrow("continuation_source_context_missing");
|
||||
await expect(result).rejects.not.toBeInstanceOf(StaleExecutionContinuationError);
|
||||
};
|
||||
const support = await getEmbeddedPostgresTestSupport();
|
||||
(support.supported ? describe : describe.skip)(
|
||||
"authorized continuation context",
|
||||
@@ -174,9 +187,10 @@ const support = await getEmbeddedPostgresTestSupport();
|
||||
const [source] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.id, runId));
|
||||
await db.update(heartbeatRuns).set({ contextSnapshot: { issueId: randomUUID() } }).where(eq(heartbeatRuns.id, runId));
|
||||
try {
|
||||
await expect(buildExecutionContinuation({ db, companyId, issueId, agentId,
|
||||
context: { interruptedRunId: runId }, summary: null, exposeLowTrustRaw: false }))
|
||||
.rejects.toThrow("continuation_source_context_missing");
|
||||
await expectMissingContinuationContext(
|
||||
() => buildExecutionContinuation({ db, companyId, issueId, agentId,
|
||||
context: { interruptedRunId: runId }, summary: null, exposeLowTrustRaw: false }),
|
||||
);
|
||||
} finally {
|
||||
await db.update(heartbeatRuns).set({ contextSnapshot: source.contextSnapshot }).where(eq(heartbeatRuns.id, runId));
|
||||
}
|
||||
@@ -336,41 +350,63 @@ const support = await getEmbeddedPostgresTestSupport();
|
||||
expect(freshPrompt).not.toContain('"resumeDelta"');
|
||||
});
|
||||
it("fails closed when required originating context is missing", async () => {
|
||||
await expect(
|
||||
buildExecutionContinuation({
|
||||
db,
|
||||
companyId,
|
||||
issueId,
|
||||
agentId,
|
||||
context: { commentId: randomUUID() },
|
||||
summary: null,
|
||||
exposeLowTrustRaw: false,
|
||||
}),
|
||||
).rejects.toThrow("continuation_source_context_missing");
|
||||
await expectMissingContinuationContext(
|
||||
() =>
|
||||
buildExecutionContinuation({
|
||||
db,
|
||||
companyId,
|
||||
issueId,
|
||||
agentId,
|
||||
context: { commentId: randomUUID() },
|
||||
summary: null,
|
||||
exposeLowTrustRaw: false,
|
||||
}),
|
||||
);
|
||||
});
|
||||
it.each(["done", "cancelled"])("rejects continuation after the task becomes %s", async (status) => {
|
||||
await db.update(issues).set({ status }).where(eq(issues.id, issueId));
|
||||
try {
|
||||
await expectStaleContinuation(
|
||||
() => buildExecutionContinuation({
|
||||
db, companyId, issueId, agentId,
|
||||
context: { wakeReason: "issue_commented", commentId: gmailId },
|
||||
summary: null, exposeLowTrustRaw: false,
|
||||
}),
|
||||
"continuation_task_ownership_changed",
|
||||
);
|
||||
const [issue] = await db.select().from(issues).where(eq(issues.id, issueId));
|
||||
expect(issue).toMatchObject({ status, assigneeAgentId: agentId });
|
||||
} finally {
|
||||
await db.update(issues).set({ status: "in_progress" }).where(eq(issues.id, issueId));
|
||||
}
|
||||
});
|
||||
it("rejects another company and an invalidated task owner", async () => {
|
||||
await expect(
|
||||
buildExecutionContinuation({
|
||||
db,
|
||||
companyId: randomUUID(),
|
||||
issueId,
|
||||
agentId,
|
||||
context: {},
|
||||
summary: null,
|
||||
exposeLowTrustRaw: false,
|
||||
}),
|
||||
).rejects.toThrow("continuation_task_ownership_changed");
|
||||
await expect(
|
||||
buildExecutionContinuation({
|
||||
db,
|
||||
companyId,
|
||||
issueId,
|
||||
agentId: randomUUID(),
|
||||
context: {},
|
||||
summary: null,
|
||||
exposeLowTrustRaw: false,
|
||||
}),
|
||||
).rejects.toThrow("continuation_task_ownership_changed");
|
||||
await expectStaleContinuation(
|
||||
() =>
|
||||
buildExecutionContinuation({
|
||||
db,
|
||||
companyId: randomUUID(),
|
||||
issueId,
|
||||
agentId,
|
||||
context: {},
|
||||
summary: null,
|
||||
exposeLowTrustRaw: false,
|
||||
}),
|
||||
"continuation_task_ownership_changed",
|
||||
);
|
||||
await expectStaleContinuation(
|
||||
() =>
|
||||
buildExecutionContinuation({
|
||||
db,
|
||||
companyId,
|
||||
issueId,
|
||||
agentId: randomUUID(),
|
||||
context: {},
|
||||
summary: null,
|
||||
exposeLowTrustRaw: false,
|
||||
}),
|
||||
"continuation_task_ownership_changed",
|
||||
);
|
||||
});
|
||||
},
|
||||
);
|
||||
|
||||
@@ -15,6 +15,13 @@ import { hasConversationContinuationPolicy } from "./conversation-continuation.j
|
||||
import { queuedCommentIdsFromWakePayload } from "./issue-queued-comment-queue.js";
|
||||
import { childReviewOutcomes } from "./native-runtime/child-review-outcomes.js";
|
||||
|
||||
export class StaleExecutionContinuationError extends Error {
|
||||
constructor(readonly code: "continuation_task_ownership_changed") {
|
||||
super(code);
|
||||
this.name = "StaleExecutionContinuationError";
|
||||
}
|
||||
}
|
||||
|
||||
const object = (v: unknown): Record<string, unknown> =>
|
||||
v && typeof v === "object" && !Array.isArray(v)
|
||||
? (v as Record<string, unknown>)
|
||||
@@ -127,7 +134,7 @@ export async function buildExecutionContinuation(input: {
|
||||
issue.assigneeAgentId !== input.agentId ||
|
||||
["done", "cancelled"].includes(issue.status)
|
||||
)
|
||||
throw new Error("continuation_task_ownership_changed");
|
||||
throw new StaleExecutionContinuationError("continuation_task_ownership_changed");
|
||||
const rows = await db
|
||||
.select()
|
||||
.from(issueComments)
|
||||
|
||||
@@ -37,7 +37,7 @@ import {
|
||||
import { executionFailureRetryCount, executionRetryAttemptCount, accountingForScheduledRetry } from "./execution-recovery-attempt.js";
|
||||
import { buildHeartbeatRunStatusLiveEventPayload } from "./heartbeat-run-status-payload.js";
|
||||
export { buildHeartbeatRunStatusLiveEventPayload } from "./heartbeat-run-status-payload.js";
|
||||
import { buildExecutionContinuation } from "./execution-continuation.js";
|
||||
import { buildExecutionContinuation, StaleExecutionContinuationError } from "./execution-continuation.js";
|
||||
import { renderPaperclipWakePrompt } from "@paperclipai/adapter-utils/server-utils";
|
||||
import { PROJECT_REPOSITORIES_DIR, readGitWorkspaceSnapshot } from "@paperclipai/adapter-utils/git-workspace-sync";
|
||||
import { isWorkspaceGitScanError, WorkspaceGitScanError, WORKSPACE_GIT_SCAN_ERROR_CODES } from "./workspace-git-operation-scheduler.js";
|
||||
@@ -25679,6 +25679,15 @@ export function heartbeatService(
|
||||
? outerErr.reason
|
||||
: "adopted_runner_authentication_timeout",
|
||||
}).catch(() => undefined);
|
||||
} else if (outerErr instanceof StaleExecutionContinuationError) {
|
||||
// The queued continuation became obsolete before adapter dispatch.
|
||||
// Use cancellation settlement so wakeup, issue ownership, agent state,
|
||||
// and notifications agree; do not retry work for the previous owner.
|
||||
await cancelRunInternal(run.id, outerErr.code, {
|
||||
errorCode: outerErr.code,
|
||||
eventMessage: "stale execution continuation cancelled before dispatch",
|
||||
suppressImmediateRecovery: true,
|
||||
});
|
||||
} else if (isWorkspaceBusyDeferral(outerErr)) {
|
||||
// Expected contention on a shared project workspace, not a
|
||||
// failure: park the run as a bounded scheduled retry and leave the
|
||||
|
||||
Reference in New Issue
Block a user