diff --git a/doc/connections/AI-CONNECTIONS.md b/doc/connections/AI-CONNECTIONS.md index bb90256852..980b07463f 100644 --- a/doc/connections/AI-CONNECTIONS.md +++ b/doc/connections/AI-CONNECTIONS.md @@ -117,6 +117,18 @@ subscription receives a retryable busy response while it is in use. Refreshes are merged only into the originating active grant, with reconnect/revocation version checks. Temporary homes are removed on normal completion or failure. +For a fresh task execution, subscription contention creates a durable scheduled +retry checked every 60–120 seconds. The task shows “Waiting for AI subscription” +and does not request a reconnect or consume its provider-failure retry allowance. +Each attempt rechecks task eligibility, ownership, budget, and current credential +access. Revocation and other configuration failures still require user action. +Authorized comment wakes that started as non-assignee runs can resume without +claiming the assignee’s execution lock. Admission records this authority while +holding the task and run locks. A reassignment during preflight cannot grant it. +Assignee retries must still own that lock. +Already-started native sessions retain their existing same-run recovery path; +they must not be replaced by a fresh execution with a pre-provider receipt. + Session reuse includes grant identity, responsible user, and credential generation. A changed identity starts a fresh provider session. Managed native executions use per-turn lifecycle cleanup; a suspended native execution whose @@ -275,7 +287,7 @@ Concurrent runs using the same subscription wait through scheduled retries while the credential lease is held. They do not request new credentials or consume the provider-failure retry allowance. Each retry revalidates the account and existing run-dispatch rules still suppress cancelled, reassigned, or otherwise ineligible work. -Task retries must retain execution-lock ownership at scheduling, promotion, and dispatch. +Assignee retries must retain execution-lock ownership at scheduling, promotion, and dispatch. `server/src/__tests__/agent-hire-ai-connections.test.ts` covers both creation routes, both providers and methods, approval gates, native provider mapping, shared access diff --git a/server/src/__tests__/ai-connections.test.ts b/server/src/__tests__/ai-connections.test.ts index 880df5edf0..6628fedbd5 100644 --- a/server/src/__tests__/ai-connections.test.ts +++ b/server/src/__tests__/ai-connections.test.ts @@ -308,6 +308,44 @@ describe("managed AI connections", () => { expect(next.identity).toBe(first.identity); await next.cleanup(); }); + it("reproduces same-agent OpenAI subscription contention and resumes without reconnecting", async () => { + const userId = "subscription-contention-user"; + await db.insert(companyMemberships).values({ companyId, principalId: userId, principalType: "user", status: "active", membershipRole: "member" }); + const credential = JSON.stringify({ tokens: { access_token: "fixture-access", refresh_token: "fixture-refresh", id_token: "fixture-id", account_id: "fixture-account" } }); + const account = await service.save(companyId, userId, { provider: "openai", method: "subscription", ownership: "personal", name: "Contention fixture", loginSessionId: "fixture", allAgents: true, agentIds: [] }, credential); + const runInput = { + companyId, + agentId, + adapterType: "codex_local", + responsibleUserId: userId, + binding: { provider: "openai", method: "subscription", mode: "responsible_user" } as const, + config: {}, + }; + // A parent can create and assign a child before its own execution ends. + // Both runs select the same personal subscription, even on the same agent. + const parent = await prepareManagedAiRuntime(db, runInput); + try { + await expect(prepareManagedAiRuntime(db, runInput)).rejects.toMatchObject({ + status: 422, + message: "This subscription is in use. Retry when its current execution finishes.", + details: { code: "ai_connection_busy" }, + }); + const selected = await service.select({ ...runInput, userId }); + expect(selected.grant).toMatchObject({ id: account.grantId, status: "active" }); + expect(await service.credential(selected)).toBe(credential); + } finally { + await parent.cleanup(); + } + // Releasing the parent alone is sufficient: no reconnect, credential + // rotation, account switch, or agent configuration change is required. + const child = await prepareManagedAiRuntime(db, runInput); + try { + expect(child.identity).toBe(parent.identity); + expect(child.attribution.grantId).toBe(account.grantId); + } finally { + await child.cleanup(); + } + }); it("persists refreshed credentials only to their original grant and fences reconnects", async () => { const auth = (marker: string, hour: number) => JSON.stringify({ tokens: { account_id: "fixture-account", id_token: `id-${marker}`, access_token: `access-${marker}`, refresh_token: `refresh-${marker}` }, last_refresh: `2026-09-10T${hour}:00:00Z` }); const intent = { provider: "openai" as const, method: "subscription" as const, name: "Refresh test", ownership: "personal" as const, agentIds: [], allAgents: true, loginSessionId: "fixture" }; diff --git a/server/src/__tests__/heartbeat-ai-subscription-contention.test.ts b/server/src/__tests__/heartbeat-ai-subscription-contention.test.ts new file mode 100644 index 0000000000..570506083f --- /dev/null +++ b/server/src/__tests__/heartbeat-ai-subscription-contention.test.ts @@ -0,0 +1,250 @@ +import { randomUUID } from "node:crypto"; +import { mkdtemp, rm } from "node:fs/promises"; +import os from "node:os"; +import path from "node:path"; +import { eq } from "drizzle-orm"; +import { afterAll, beforeAll, describe, expect, it, vi } from "vitest"; +import { + agents, companies, companyMemberships, connectionGrants, createDb, heartbeatRuns, + issueComments, issueRecoveryActions, issueThreadInteractions, issues, +} from "@paperclipai/db"; +import { startEmbeddedPostgresTestDatabase } from "@paperclipai/db/test-embedded-postgres"; +import { heartbeatService } from "../services/heartbeat.js"; +import { aiConnectionService } from "../services/ai-connections.js"; +import { prepareManagedAiRuntime } from "../services/ai-connection-runtime.js"; +import { executionFailureRetryCount } from "../services/execution-recovery-attempt.js"; +import { executionProjectionsForRuns } from "../services/execution-projection.js"; +import type { AdapterExecutionContext } from "../adapters/index.js"; + +const execute = vi.hoisted(() => vi.fn(async (_input: AdapterExecutionContext) => ({ + exitCode: 0, signal: null, timedOut: false, resultJson: {}, summary: "Fixture work completed.", +}))); +vi.mock("../adapters/index.js", async () => ({ + ...await vi.importActual("../adapters/index.js"), + getServerAdapter: () => ({ supportsLocalAgentJwt: false, execute }), +})); + +const afterCheckout = vi.hoisted(() => vi.fn(async () => {})); +vi.mock("../services/issues.js", async () => { + const actual = await vi.importActual("../services/issues.js"); + return { ...actual, issueService: (...args: Parameters) => { + const service = actual.issueService(...args); + return { ...service, checkout: async (...checkoutArgs: Parameters) => { + const result = await service.checkout(...checkoutArgs); + await afterCheckout(); + return result; + } }; + } }; +}); + +describe("heartbeat AI subscription contention", () => { + let database: Awaited>; + let db: ReturnType; + let heartbeat: ReturnType; + let home: string; + + beforeAll(async () => { + home = await mkdtemp(path.join(os.tmpdir(), "paperclip-subscription-contention-")); + vi.stubEnv("PAPERCLIP_HOME", home); + vi.stubEnv("PAPERCLIP_INSTANCE_ID", "subscription-contention-tests"); + database = await startEmbeddedPostgresTestDatabase("paperclip-subscription-contention-db-"); + db = createDb(database.connectionString); + heartbeat = heartbeatService(db); + execute.mockImplementation(async (input) => { + await db.update(issues).set({ status: "done" }).where(eq(issues.id, String(input.context.issueId))); + return { exitCode: 0, signal: null, timedOut: false, resultJson: {}, summary: "Fixture work completed." }; + }); + }, 90_000); + afterAll(async () => { + await heartbeat?.drainActiveRunExecutions(); + await database?.cleanup(); + vi.unstubAllEnvs(); + if (home) await rm(home, { recursive: true, force: true }); + }); + + async function fixture() { + execute.mockClear(); + const companyId = randomUUID(); + const agentId = randomUUID(); + const issueId = randomUUID(); + const userId = randomUUID(); + await db.insert(companies).values({ id: companyId, name: "Contention fixture", issuePrefix: `S${companyId.slice(0, 8).toUpperCase()}`, defaultResponsibleUserId: userId }); + await db.insert(companyMemberships).values({ companyId, principalId: userId, principalType: "user", status: "active", membershipRole: "member" }); + const binding = { provider: "openai", method: "subscription", mode: "responsible_user" } as const; + await db.insert(agents).values({ id: agentId, companyId, name: "CEO", adapterType: "codex_local", adapterConfig: { cwd: home }, runtimeConfig: { aiConnection: binding, heartbeat: { wakeOnDemand: true, maxConcurrentRuns: 20 } } }); + await db.insert(issues).values({ id: issueId, companyId, title: "Child task", status: "todo", assigneeAgentId: agentId, responsibleUserId: userId }); + const account = await aiConnectionService(db).save(companyId, userId, { provider: "openai", method: "subscription", ownership: "personal", name: "Fixture subscription", loginSessionId: "fixture", allAgents: true, agentIds: [] }, JSON.stringify({ tokens: { access_token: "fixture-access", refresh_token: "fixture-refresh", id_token: "fixture-id", account_id: "fixture-account" } })); + const parent = await prepareManagedAiRuntime(db, { companyId, agentId, responsibleUserId: userId, adapterType: "codex_local", binding, config: { cwd: home } }); + let released = false; + const release = async () => { + if (released) return; + released = true; + await parent.cleanup(); + }; + return { companyId, agentId, issueId, userId, account, parent, release }; + } + + async function defer(f: Awaited>, context: Record = {}) { + const run = await heartbeat.invoke(f.agentId, "assignment", { + issueId: f.issueId, wakeReason: "issue_assigned", ...context, + }, "system"); + expect(run).not.toBeNull(); + await heartbeat.drainActiveRunExecutions(); + expect(await heartbeat.getRun(run!.id)).toMatchObject({ status: "cancelled", errorCode: "ai_connection_busy" }); + const [retry] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, run!.id)); + expect(retry).toMatchObject({ status: "scheduled_retry", scheduledRetryReason: "ai_connection_busy" }); + return retry!; + } + + async function dispatch(retry: typeof heartbeatRuns.$inferSelect, service = heartbeat) { + await service.promoteDueScheduledRetries(new Date(retry.scheduledRetryAt!.getTime() + 1)); + await service.resumeQueuedRuns(); + await service.drainActiveRunExecutions(); + } + + it("defers a child without a reconnect blocker and executes after its parent's subscription is released", async () => { + const f = await fixture(); + try { + const run = await heartbeat.invoke(f.agentId, "assignment", { issueId: f.issueId, wakeReason: "issue_assigned" }, "system"); + expect(run).not.toBeNull(); + await heartbeat.drainActiveRunExecutions(); + const deferred = await heartbeat.getRun(run!.id); + expect(deferred).toMatchObject({ status: "cancelled", errorCode: "ai_connection_busy" }); + expect(execute).not.toHaveBeenCalled(); + const [retry] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, run!.id)); + expect(retry).toMatchObject({ status: "scheduled_retry", scheduledRetryReason: "ai_connection_busy" }); + expect(executionFailureRetryCount(retry)).toBe(0); + const [issue] = await db.select().from(issues).where(eq(issues.id, f.issueId)); + expect(issue.status).not.toBe("blocked"); + expect(issue.executionRunId).toBe(retry.id); + expect(await db.select().from(issueRecoveryActions).where(eq(issueRecoveryActions.companyId, f.companyId))).toHaveLength(0); + expect(await db.select().from(issueThreadInteractions).where(eq(issueThreadInteractions.companyId, f.companyId))).toHaveLength(0); + expect(await db.select().from(issueComments).where(eq(issueComments.issueId, f.issueId))).toHaveLength(0); + await f.release(); + await heartbeat.promoteDueScheduledRetries(new Date(retry.scheduledRetryAt!.getTime() + 1)); + await heartbeat.resumeQueuedRuns(); + await heartbeat.drainActiveRunExecutions(); + const completed = await heartbeat.getRun(retry.id); + expect(completed).toMatchObject({ status: "succeeded" }); + expect(execute).toHaveBeenCalledTimes(1); + expect(completed?.contextSnapshot?.aiConnection).toMatchObject({ grantId: f.account.grantId, identity: f.parent.identity }); + } finally { + await f.release(); + } + }); + + it("keeps waiting beyond the failure-attempt limit and survives a new service instance", async () => { + const f = await fixture(); + try { + let retry = await defer(f, { + failureRetriesBeforeAiConnectionWait: 99, aiConnectionBusyDeferredWhileAssignee: false, + }); + expect(retry.contextSnapshot?.aiConnectionBusyDeferredWhileAssignee).toBe(true); + for (let count = 0; count < 3; count++) { + await dispatch(retry); + const [successor] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, retry.id)); + expect(successor).toMatchObject({ status: "scheduled_retry", scheduledRetryAttempt: count + 2 }); + expect(executionFailureRetryCount(successor)).toBe(0); + retry = successor!; + const projections = await executionProjectionsForRuns(db, f.companyId, [retry.id]); + expect(projections.get(retry.id)).toMatchObject({ + phase: "retry_scheduled", label: "Waiting for AI subscription", attempt: 1, + }); + } + expect(execute).not.toHaveBeenCalled(); + await f.release(); + await dispatch(retry, heartbeatService(db)); + expect(await heartbeat.getRun(retry.id)).toMatchObject({ status: "succeeded" }); + expect(execute).toHaveBeenCalledTimes(1); + } finally { + await f.release(); + } + }); + + it.each([ + ["cancelled", "issue_cancelled"], + ["done", "issue_terminal_status"], + ["reassigned", "issue_reassigned"], + ["agent_paused", "agent_not_invokable"], + ])("does not dispatch after the waiting task is %s", async (change, errorCode) => { + const f = await fixture(); + try { + const retry = await defer(f); + if (change === "agent_paused") { + await db.update(agents).set({ status: "paused" }).where(eq(agents.id, f.agentId)); + } else { + await db.update(issues).set(change === "reassigned" + ? { assigneeAgentId: null, assigneeUserId: f.userId } + : { status: change }).where(eq(issues.id, f.issueId)); + } + await f.release(); + await dispatch(retry); + expect(await heartbeat.getRun(retry.id)).toMatchObject({ status: "cancelled", errorCode }); + expect(execute).not.toHaveBeenCalled(); + const [issue] = await db.select().from(issues).where(eq(issues.id, f.issueId)); + expect(issue.executionRunId).toBeNull(); + } finally { + await f.release(); + } + }); + + it.each(["issue_assigned", "issue_commented"])("does not turn an assignee %s wake into a non-assignee retry after reassignment during preflight", async (wakeReason) => { + const f = await fixture(); + try { + afterCheckout.mockImplementationOnce(async () => { + await db.update(issues).set({ + assigneeAgentId: null, assigneeUserId: f.userId, executionRunId: null, + }).where(eq(issues.id, f.issueId)); + }); + const commentId = randomUUID(); + await db.insert(issueComments).values({ id: commentId, companyId: f.companyId, issueId: f.issueId, + authorType: "user", authorUserId: f.userId, body: "Continue this task." }); + const run = await heartbeat.invoke(f.agentId, "assignment", { + issueId: f.issueId, wakeReason, commentId, aiConnectionBusyDeferredWhileAssignee: false, + }, "system"); + await heartbeat.drainActiveRunExecutions(); + expect(await heartbeat.getRun(run!.id)).toMatchObject({ status: "cancelled", errorCode: "ai_connection_busy" }); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, run!.id))).toHaveLength(0); + expect(execute).not.toHaveBeenCalled(); + } finally { + afterCheckout.mockReset(); + await f.release(); + } + }); + + it("resumes an authorized comment wake without stealing the assignee's task lock", async () => { + const f = await fixture(); + try { + await db.update(issues).set({ assigneeAgentId: null, assigneeUserId: f.userId }).where(eq(issues.id, f.issueId)); + const retry = await defer(f, { wakeReason: "issue_comment_mentioned", commentId: randomUUID() }); + expect(retry.contextSnapshot?.aiConnectionBusyDeferredWhileAssignee).toBe(false); + await dispatch(retry); + const [successor] = await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, retry.id)); + expect(successor).toMatchObject({ status: "scheduled_retry" }); + expect(successor.contextSnapshot?.aiConnectionBusyDeferredWhileAssignee).toBe(false); + const [issue] = await db.select().from(issues).where(eq(issues.id, f.issueId)); + expect(issue.executionRunId).toBeNull(); + await f.release(); + await dispatch(successor); + expect(await heartbeat.getRun(successor.id)).toMatchObject({ status: "succeeded" }); + expect(execute).toHaveBeenCalledTimes(1); + } finally { + await f.release(); + } + }); + + it("rechecks revoked credentials after waiting and retains the configuration failure", async () => { + const f = await fixture(); + try { + const retry = await defer(f); + await f.release(); + await db.update(connectionGrants).set({ status: "revoked", revokedAt: new Date() }).where(eq(connectionGrants.id, f.account.grantId)); + await dispatch(retry); + expect(await heartbeat.getRun(retry.id)).toMatchObject({ status: "failed", errorCode: "configuration_incomplete" }); + expect(execute).not.toHaveBeenCalled(); + expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, retry.id))).toHaveLength(0); + } finally { + await f.release(); + } + }); +}); diff --git a/server/src/__tests__/heartbeat-retry-scheduling.test.ts b/server/src/__tests__/heartbeat-retry-scheduling.test.ts index 6484e678b0..38bbf2bb49 100644 --- a/server/src/__tests__/heartbeat-retry-scheduling.test.ts +++ b/server/src/__tests__/heartbeat-retry-scheduling.test.ts @@ -265,12 +265,15 @@ describeEmbeddedPostgres("heartbeat bounded retry scheduling", () => { expect(await db.select().from(heartbeatRuns).where(eq(heartbeatRuns.retryOfRunId, runId))).toHaveLength(1); }); - it("retains the failure budget after many pre-provider workspace waits", async () => { + it.each([ + ["workspace_busy", "failureRetriesBeforeWorkspaceWait"], + ["ai_connection_busy", "failureRetriesBeforeAiConnectionWait"], + ])("retains the failure budget after many pre-provider %s waits", async (reason, countKey) => { const runId = randomUUID(), companyId = randomUUID(), agentId = randomUUID(); const now = new Date("2026-04-20T12:00:00.000Z"); await seedRetryFixture({ runId, companyId, agentId, now, errorCode: "overloaded", errorFamily: "transient_upstream" }); - await db.update(heartbeatRuns).set({ scheduledRetryReason: "workspace_busy", scheduledRetryAttempt: 12, - contextSnapshot: { failureRetriesBeforeWorkspaceWait: 1 } }).where(eq(heartbeatRuns.id, runId)); + await db.update(heartbeatRuns).set({ scheduledRetryReason: reason, scheduledRetryAttempt: 12, + contextSnapshot: { [countKey]: 1 } }).where(eq(heartbeatRuns.id, runId)); const scheduled = await heartbeat.scheduleBoundedRetry(runId, { now, random: () => 0 }); expect(scheduled).toMatchObject({ outcome: "scheduled", run: { scheduledRetryAttempt: 2, scheduledRetryReason: "transient_failure" } }); if (scheduled.outcome !== "scheduled") throw new Error("Expected a bounded retry"); diff --git a/server/src/modules/run-dispatch/adapters/postgres.ts b/server/src/modules/run-dispatch/adapters/postgres.ts index b89f2f0731..1190dc3268 100644 --- a/server/src/modules/run-dispatch/adapters/postgres.ts +++ b/server/src/modules/run-dispatch/adapters/postgres.ts @@ -908,9 +908,9 @@ export function createPostgresRunDispatchAdapter( async function decideCurrentRunStaleness(tx: Db, run: HeartbeatRun, now: Date) { const contextSnapshot = parseObject(run.contextSnapshot); const issueId = readNonEmptyString(contextSnapshot.issueId); - if (!issueId) return { issueId: null, decision: { stale: false as const } }; + if (!issueId) return { issueId: null, facts: null, decision: { stale: false as const } }; const recovery = await getExecutionBlocker(tx, run.companyId, issueId, { conversationResetCommentId: deriveCommentId(contextSnapshot) }); - if (recovery) return { issueId, decision: { stale: true as const, + if (recovery) return { issueId, facts: null, decision: { stale: true as const, errorCode: "execution_reconciliation_required" as const, reason: recovery.nextAction, details: { issueId, recoveryActionId: recovery.recoveryActionId }, } }; @@ -926,7 +926,7 @@ export function createPostgresRunDispatchAdapter( now, tx, ); - return { issueId, decision: decideQueuedRunStaleness(facts, now) }; + return { issueId, facts, decision: decideQueuedRunStaleness(facts, now) }; } async function cancelStaleQueuedRun( @@ -934,8 +934,21 @@ export function createPostgresRunDispatchAdapter( ): Promise { const cancelLockedRun = async (tx: Db, run: HeartbeatRun) => { if (run.status !== input.expectedStatus) return { outcome: "lost_race" as const }; - const { issueId, decision } = await decideCurrentRunStaleness(tx, run, input.now); - if (!decision.stale || !issueId) return { outcome: "not_stale" as const }; + const { issueId, facts, decision } = await decideCurrentRunStaleness(tx, run, input.now); + if (!decision.stale || !issueId) { + if (input.expectedStatus === "queued" && facts?.isInteractionWake) { + // Preserve the authority accepted under the issue/run locks. Later + // preflight reads can observe a reassignment; they must not turn an + // assignee comment into a non-assignee subscription-wait exception. + await tx.update(heartbeatRuns).set({ + runnerProfileJson: { + ...parseObject(run.runnerProfileJson), + aiConnectionNonAssigneeCommentWake: facts.issueAssigneeAgentId !== run.agentId, + }, + }).where(and(eq(heartbeatRuns.id, run.id), eq(heartbeatRuns.companyId, run.companyId))); + } + return { outcome: "not_stale" as const }; + } return cancelStaleRunInTx(tx, run, issueId, decision, input.expectedStatus, input.now); }; diff --git a/server/src/modules/run-dispatch/domain/policy.test.ts b/server/src/modules/run-dispatch/domain/policy.test.ts index eeb4035d6b..2efe03f3af 100644 --- a/server/src/modules/run-dispatch/domain/policy.test.ts +++ b/server/src/modules/run-dispatch/domain/policy.test.ts @@ -486,3 +486,23 @@ describe("native replacement execution authority", () => { expect(decideQueuedRunStaleness({ ...baseStalenessFacts(), retryReasonKind: "native_safe_replacement", ...locks }, NOW)).toEqual({ stale: false }); }); }); + + +describe("AI subscription wait ownership", () => { + it.each([null, "another-run"])("preserves assignee lock guards and admits only recorded non-assignee waits: %s", (issueExecutionRunId) => { + const gate = { + ...baseGateFacts(), retryReasonKind: "ai_connection_wait" as const, + enforceIssueExecutionLock: true, issueExecutionRunId, + }; + const queued = { ...baseStalenessFacts(), retryReasonKind: "ai_connection_wait" as const, issueExecutionRunId }; + expect(decideScheduledRetryGate(gate, NOW)).toMatchObject({ allowed: false, errorCode: "issue_execution_lock_changed" }); + expect(decideQueuedRunStaleness(queued, NOW)).toMatchObject({ stale: true, errorCode: "issue_execution_lock_changed" }); + const nonAssignee = { issueAssigneeAgentId: "another-agent", isNonAssigneeWorkspaceBusyRetry: true }; + expect(decideScheduledRetryGate({ ...gate, ...nonAssignee }, NOW)).toEqual({ allowed: true }); + expect(decideQueuedRunStaleness({ ...queued, ...nonAssignee }, NOW)).toEqual({ stale: false }); + expect(decideScheduledRetryGate({ ...gate, ...nonAssignee, retryReasonKind: "max_turn_continuation" }, NOW)) + .toMatchObject({ allowed: false, errorCode: "issue_execution_lock_changed" }); + expect(decideQueuedRunStaleness({ ...queued, ...nonAssignee, retryReasonKind: "max_turn_continuation" }, NOW)) + .toMatchObject({ stale: true, errorCode: "issue_execution_lock_changed" }); + }); +}); diff --git a/server/src/modules/run-dispatch/domain/policy.ts b/server/src/modules/run-dispatch/domain/policy.ts index 3049b2f141..884ca30816 100644 --- a/server/src/modules/run-dispatch/domain/policy.ts +++ b/server/src/modules/run-dispatch/domain/policy.ts @@ -404,7 +404,10 @@ export function decideScheduledRetryGate( } const lockOutcome = decideExecutionLock({ - requiresExecutionLock: (requiresInProgress || facts.retryReasonKind === "ai_connection_wait") && facts.enforceIssueExecutionLock, + requiresExecutionLock: + (requiresInProgress || + (facts.retryReasonKind === "ai_connection_wait" && !facts.isNonAssigneeWorkspaceBusyRetry)) && + facts.enforceIssueExecutionLock, runId: facts.runId, issueExecutionRunId: facts.issueExecutionRunId, }); @@ -630,7 +633,9 @@ export function decideQueuedRunStaleness( } const lockOutcome = decideExecutionLock({ - requiresExecutionLock: requiresInProgress || facts.retryReasonKind === "ai_connection_wait", + // A server-recorded non-assignee wake never held the task execution lock. + requiresExecutionLock: requiresInProgress || + (facts.retryReasonKind === "ai_connection_wait" && !facts.isNonAssigneeWorkspaceBusyRetry), runId: facts.runId, issueExecutionRunId: facts.issueExecutionRunId, }); diff --git a/server/src/modules/run-dispatch/domain/wake-context.test.ts b/server/src/modules/run-dispatch/domain/wake-context.test.ts index 574f9e7ac8..223ee6fb8b 100644 --- a/server/src/modules/run-dispatch/domain/wake-context.test.ts +++ b/server/src/modules/run-dispatch/domain/wake-context.test.ts @@ -9,6 +9,15 @@ import { } from "./wake-context.js"; describe("wake context", () => { + it("requires the subscription-specific non-assignee receipt for a subscription wait", () => { + expect(isNonAssigneeWorkspaceBusyRetry("ai_connection_busy", { + aiConnectionBusyDeferredWhileAssignee: false, + })).toBe(true); + for (const context of [ {}, { aiConnectionBusyDeferredWhileAssignee: true }, { workspaceBusyDeferredWhileAssignee: false } ]) { + expect(isNonAssigneeWorkspaceBusyRetry("ai_connection_busy", context)).toBe(false); + } + expect(isNonAssigneeWorkspaceBusyRetry("transient_failure", { aiConnectionBusyDeferredWhileAssignee: false })).toBe(false); + }); it("recognizes only workspace-busy retries deferred outside assignee-ship", () => { expect(isNonAssigneeWorkspaceBusyRetry(WORKSPACE_BUSY_RETRY_REASON, { workspaceBusyDeferredWhileAssignee: false, diff --git a/server/src/modules/run-dispatch/domain/wake-context.ts b/server/src/modules/run-dispatch/domain/wake-context.ts index 98b9748d63..b0dd653959 100644 --- a/server/src/modules/run-dispatch/domain/wake-context.ts +++ b/server/src/modules/run-dispatch/domain/wake-context.ts @@ -17,6 +17,7 @@ function readNonEmptyString(value: unknown): string | null { export const MAX_TURN_CONTINUATION_RETRY_REASON = "max_turns_continuation"; export const WORKSPACE_BUSY_RETRY_REASON = "workspace_busy"; +export const AI_CONNECTION_BUSY_RETRY_REASON = "ai_connection_busy"; export const INTERACTION_CONTINUATION_INFRA_RETRY_REASON = "interaction_continuation_infra_retry"; export const INTERACTION_CONTINUATION_INFRA_WAKE_REASON = "interaction_continuation_infra_retry"; export const WAKE_COMMENT_IDS_KEY = "wakeCommentIds"; @@ -28,7 +29,7 @@ export const RESOLVED_INTERACTION_CONTINUATION_STATUSES = new Set([ ]); /** - * True for the retry of a workspace-busy deferral whose original run did not + * True for a resource-wait retry whose original run did not * execute under assignee-ship (a comment or review-participant wake). Such a * retry has an expected assignee mismatch, so the scheduled-retry gate and * the queued-run staleness check must not treat it as a reassignment. @@ -38,8 +39,10 @@ export function isNonAssigneeWorkspaceBusyRetry( contextSnapshot: Record, ): boolean { return ( - retryReason === WORKSPACE_BUSY_RETRY_REASON && - contextSnapshot.workspaceBusyDeferredWhileAssignee === false + (retryReason === WORKSPACE_BUSY_RETRY_REASON && + contextSnapshot.workspaceBusyDeferredWhileAssignee === false) || + (retryReason === AI_CONNECTION_BUSY_RETRY_REASON && + contextSnapshot.aiConnectionBusyDeferredWhileAssignee === false) ); } diff --git a/server/src/modules/run-dispatch/index.ts b/server/src/modules/run-dispatch/index.ts index 20f3473bda..d6fda59345 100644 --- a/server/src/modules/run-dispatch/index.ts +++ b/server/src/modules/run-dispatch/index.ts @@ -12,6 +12,7 @@ import type { RunDispatchWriter, ScheduledRetryReader } from "./application/port export { MAX_TURN_CONTINUATION_RETRY_REASON, WORKSPACE_BUSY_RETRY_REASON, + AI_CONNECTION_BUSY_RETRY_REASON, INTERACTION_CONTINUATION_INFRA_RETRY_REASON, INTERACTION_CONTINUATION_INFRA_WAKE_REASON, WAKE_COMMENT_IDS_KEY, diff --git a/server/src/services/ai-connection-runtime.ts b/server/src/services/ai-connection-runtime.ts index 0608286723..8e05a3d79a 100644 --- a/server/src/services/ai-connection-runtime.ts +++ b/server/src/services/ai-connection-runtime.ts @@ -1,5 +1,5 @@ import { createHash } from "node:crypto"; -import { unprocessable } from "../errors.js"; +import { HttpError, unprocessable } from "../errors.js"; import { mkdtemp, mkdir, writeFile, readFile, rm } from "node:fs/promises"; import os from "node:os"; import path from "node:path"; @@ -16,6 +16,11 @@ import type { AdapterExecutionTarget } from "@paperclipai/adapter-utils/executio import { runAdapterExecutionTargetProcess } from "@paperclipai/adapter-utils/execution-target"; import { decideGrokAuthMerge } from "@paperclipai/adapter-grok-local/server"; +export function isAiConnectionBusy(error: unknown): error is HttpError { + return error instanceof HttpError && error.status === 422 && + (error.details as { code?: unknown } | undefined)?.code === "ai_connection_busy"; +} + // Blank values intentionally override inherited credentials in CLI child environments. export const AI_AUTH_ENV_KEYS = [ "ANTHROPIC_API_KEY", diff --git a/server/src/services/execution-projection.test.ts b/server/src/services/execution-projection.test.ts index 2db13acc7a..f0dbddb79c 100644 --- a/server/src/services/execution-projection.test.ts +++ b/server/src/services/execution-projection.test.ts @@ -37,6 +37,11 @@ const project = ( ) => projectExecution(r, c, pending, undefined, now); describe("execution truth projection", () => { + it("shows subscription contention as a resource wait without failed provider attempts", () => { + expect(projectExecution(run({ runtimeMode: "legacy", status: "scheduled_retry", scheduledRetryReason: "ai_connection_busy", + scheduledRetryAttempt: 12, contextSnapshot: { failureRetriesBeforeAiConnectionWait: 0 } }), undefined, [], undefined, now)) + .toMatchObject({ label: "Waiting for AI subscription", phase: "retry_scheduled", attempt: 1, recoveryOwner: null }); + }); it("shows a workspace wait without presenting its deferral count as failed attempts", () => { expect(projectExecution(run({ runtimeMode: "legacy", status: "scheduled_retry", scheduledRetryReason: "workspace_busy", scheduledRetryAttempt: 12, contextSnapshot: { failureRetriesBeforeWorkspaceWait: 1 } }), undefined, [], undefined, now)) diff --git a/server/src/services/execution-projection.ts b/server/src/services/execution-projection.ts index 24e60004a1..65d4501098 100644 --- a/server/src/services/execution-projection.ts +++ b/server/src/services/execution-projection.ts @@ -29,6 +29,7 @@ const executionRunColumns = { status: heartbeatRuns.status, contextSnapshot: sql>`jsonb_build_object( 'issueId', ${heartbeatRuns.contextSnapshot}->'issueId', + 'failureRetriesBeforeAiConnectionWait', ${heartbeatRuns.contextSnapshot}->'failureRetriesBeforeAiConnectionWait', 'failureRetriesBeforeWorkspaceWait', ${heartbeatRuns.contextSnapshot}->'failureRetriesBeforeWorkspaceWait')`, }; type Run = Pick; @@ -230,6 +231,10 @@ export function projectExecution( ) return set("finishing", "Finishing"); if (successorRunId) return set("completed", "Continued in another run"); + if (run.status === "scheduled_retry" && run.scheduledRetryReason === "ai_connection_busy") { + projection.nextAction = "Waiting for the AI subscription's current execution to finish; the scheduled check will revalidate access."; + return set("retry_scheduled", "Waiting for AI subscription"); + } if (run.status === "scheduled_retry" && run.scheduledRetryReason === "workspace_busy") { projection.nextAction = "Waiting for the live workspace holder to finish; the scheduled check will revalidate ownership."; return set("retry_scheduled", "Waiting for workspace"); diff --git a/server/src/services/execution-recovery-attempt.test.ts b/server/src/services/execution-recovery-attempt.test.ts index f1757f5e91..f18baa135a 100644 --- a/server/src/services/execution-recovery-attempt.test.ts +++ b/server/src/services/execution-recovery-attempt.test.ts @@ -2,11 +2,15 @@ import { describe, expect, it } from "vitest"; import { executionFailureRetryCount } from "./execution-recovery-attempt.js"; describe("failure attempts across resource waits", () => { - it("preserves prior failures through repeated AI subscription waits", () => { + it("preserves prior failures through subscription waits without trusting unrelated context", () => { expect(executionFailureRetryCount({ scheduledRetryReason: "ai_connection_busy", scheduledRetryAttempt: 12, contextSnapshot: { failureRetriesBeforeAiConnectionWait: 1 } })).toBe(1); expect(executionFailureRetryCount({ scheduledRetryReason: "transient_failure", scheduledRetryAttempt: 2, contextSnapshot: { failureRetriesBeforeAiConnectionWait: 0 } })).toBe(2); + for (const count of [undefined, -1, 1.5, "0"]) { + expect(executionFailureRetryCount({ scheduledRetryReason: "ai_connection_busy", scheduledRetryAttempt: 4, + contextSnapshot: { failureRetriesBeforeAiConnectionWait: count } })).toBe(4); + } }); it("preserves prior failures through repeated workspace waits", () => { expect(executionFailureRetryCount({ scheduledRetryReason: "workspace_busy", scheduledRetryAttempt: 12, diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 220621e69a..e2bb3fc2be 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -7,7 +7,7 @@ import { hasRemoteTerminationReceipt, remoteExecutionHasStopped, remoteTerminati import { applyConnectorSkills, prepareConnectorSkillDelivery, resolveConnectorAssignments } from "./connector-runtime.js"; import { admitExplicitNativeContinuation, undeliveredLegacyUserCommentIds } from "./explicit-native-continuation.js"; import { connectionIntentService } from "./connection-intents.js"; -import { prepareManagedAiRuntime, assertManagedAiProjectAuth, stripAiAuthBindings, AI_AUTH_ENV_KEYS } from "./ai-connection-runtime.js"; +import { prepareManagedAiRuntime, assertManagedAiProjectAuth, stripAiAuthBindings, isAiConnectionBusy, AI_AUTH_ENV_KEYS } from "./ai-connection-runtime.js"; import { aiConnectionBindingSchema } from "@paperclipai/shared"; import { executionBlockerPredicate, getExecutionBlocker } from "./execution-blocker.js"; import { CONVERSATION_CONTINUATION_POLICY, claimedAdapterType, runUsedConversationAdapter, hasConversationContinuationPolicy, isConversationAdapter } from "./conversation-continuation.js"; @@ -491,6 +491,7 @@ import { type PostCommitEffect, MAX_TURN_CONTINUATION_RETRY_REASON, WORKSPACE_BUSY_RETRY_REASON, + AI_CONNECTION_BUSY_RETRY_REASON, INTERACTION_CONTINUATION_INFRA_RETRY_REASON, INTERACTION_CONTINUATION_INFRA_WAKE_REASON, WAKE_COMMENT_IDS_KEY, @@ -842,10 +843,6 @@ const MAX_TURN_CONTINUATION_LIVE_RUN_STATUSES = [ export { WORKSPACE_BUSY_RETRY_REASON }; export const WORKSPACE_BUSY_RETRY_WAKE_REASON = "workspace_busy_retry"; export const WORKSPACE_BUSY_ERROR_CODE = "workspace_busy"; -const AI_CONNECTION_BUSY_RETRY_REASON = "ai_connection_busy"; -function isAiConnectionBusy(error: unknown): error is HttpError { - return error instanceof HttpError && parseObject(error.details).code === AI_CONNECTION_BUSY_RETRY_REASON; -} export const WORKSPACE_BUSY_RETRY_BASE_DELAY_MS = 60 * 1000; export const WORKSPACE_BUSY_RETRY_JITTER_MS = 60 * 1000; // A running run stops counting as a shared-workspace holder once it has been @@ -15543,7 +15540,10 @@ export function heartbeatService( } } - if (retryReason === AI_CONNECTION_BUSY_RETRY_REASON && issueId) { + if ( + retryReason === AI_CONNECTION_BUSY_RETRY_REASON && issueId && + !isNonAssigneeWorkspaceBusyRetry(retryReason, contextSnapshot) + ) { // The issue row is locked above. Recheck after the preflight gate so // cancellation or recovery cannot leave a successor without its lock. const [lockedIssue] = await tx.select({ executionRunId: issues.executionRunId }) @@ -15943,11 +15943,19 @@ export function heartbeatService( // by another run. Contention is a resource wait, not broken authentication. // The database lease is released on completion/disconnect; retrying remains // safe across processes and each attempt revalidates the selected account. - async function finalizeAiConnectionBusyDeferral(run: typeof heartbeatRuns.$inferSelect, error: HttpError) { + async function finalizeAiConnectionBusyDeferral( + run: typeof heartbeatRuns.$inferSelect, + error: HttpError, + wasIssueAssignee: boolean, + ) { const now = new Date(); const cancelled = await setRunStatusIfRunning(run.id, "cancelled", { error: error.message, errorCode: AI_CONNECTION_BUSY_RETRY_REASON, finishedAt: now, resultJson: { executionRecovery: { kind: "ai_connection_wait", providerWorkStarted: false } }, + contextSnapshot: { + ...parseObject(run.contextSnapshot), + aiConnectionBusyDeferredWhileAssignee: wasIssueAssignee, + }, }); if (!cancelled.updated) return; await setWakeupStatus(run.wakeupRequestId, "cancelled", { finishedAt: now, error: error.message }).catch(() => undefined); @@ -15969,8 +15977,11 @@ export function heartbeatService( }); } } finally { - if (cancelledRun && !scheduled) await releaseIssueExecutionAndPromote(cancelledRun); - await finalizeAgentStatus(run.agentId, "cancelled", null, { wasFirstHeartbeat: timerClaimWasFirstHeartbeat(run) }); + try { + if (cancelledRun && !scheduled) await releaseIssueExecutionAndPromote(cancelledRun); + } finally { + await finalizeAgentStatus(run.agentId, "cancelled", null, { wasFirstHeartbeat: timerClaimWasFirstHeartbeat(run) }); + } } } @@ -20769,7 +20780,18 @@ export function heartbeatService( try { managedAiRuntime = await prepareManagedAiRuntime(db, { companyId: agent.companyId, agentId: agent.id, responsibleUserId, adapterType: agent.adapterType, binding: aiBinding, config: resolvedConfig }); } catch (error) { - if (isAiConnectionBusy(error)) throw error; + // Only fresh executions can receive a pre-provider wait receipt. A + // persisted native input may already have provider effects to recover. + if (isAiConnectionBusy(error) && !persistedNativeExecutionInput) { + // Use the authority recorded by the locked admission gate, never + // the issue's mutable assignee observed during runtime preparation. + const authorizedNonAssigneeWake = + parseObject(run.runnerProfileJson).aiConnectionNonAssigneeCommentWake === true || + (run.scheduledRetryReason === AI_CONNECTION_BUSY_RETRY_REASON && + isNonAssigneeWorkspaceBusyRetry(run.scheduledRetryReason, parseObject(run.contextSnapshot))); + await finalizeAiConnectionBusyDeferral(run, error, !authorizedNonAssigneeWake); + return; + } if (responsibleUserId && issueId && aiBinding.mode === "responsible_user") { await connectionIntentService(db).request({ sub: agent.id, company_id: agent.companyId, run_id: run.id, responsible_user_id: responsibleUserId }, aiBinding.provider, { purpose: "ai" }).catch(() => { logger.warn({ runId: run.id, agentId: agent.id }, "Could not attach AI connection request; runtime configuration action remains available"); @@ -25077,10 +25099,6 @@ export function heartbeatService( ? outerErr.reason : "adopted_runner_authentication_timeout", }).catch(() => undefined); - } else if (isAiConnectionBusy(outerErr)) { - await finalizeAiConnectionBusyDeferral(run, outerErr).catch((error) => { - logger.error({ err: error, runId }, "failed to schedule a retry for the busy AI subscription"); - }); } 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 diff --git a/server/src/services/legacy-execution-recovery.test.ts b/server/src/services/legacy-execution-recovery.test.ts index 62558c6053..e1b81f3cca 100644 --- a/server/src/services/legacy-execution-recovery.test.ts +++ b/server/src/services/legacy-execution-recovery.test.ts @@ -9,6 +9,20 @@ const stopped = { }, }; +it("permits subscription waits only with explicit evidence that provider work never started", () => { + const waiting = { + runtimeMode: "legacy", status: "cancelled", errorCode: "ai_connection_busy", scheduledRetryAttempt: 12, + resultJson: { executionRecovery: { kind: "ai_connection_wait", providerWorkStarted: false } }, + }; + expect(legacyExecutionNeedsReconciliation(waiting)).toBe(false); + expect(legacyExecutionNeedsReconciliation({ ...waiting, status: "failed" })).toBe(true); + expect(legacyExecutionNeedsReconciliation({ ...waiting, errorCode: "cancelled" })).toBe(true); + expect(legacyExecutionNeedsReconciliation({ ...waiting, resultJson: {} })).toBe(true); + expect(legacyExecutionNeedsReconciliation({ ...waiting, resultJson: { + executionRecovery: { kind: "ai_connection_wait", providerWorkStarted: true }, + } })).toBe(true); +}); + it("allows a confirmed interrupted checkpoint without treating ordinary cancellation as replay permission", () => { expect(legacyExecutionNeedsReconciliation(stopped)).toBe(false); expect(legacyExecutionNeedsReconciliation({ ...stopped, resultJson: {} })).toBe(true);