From ee81cee76d9a374f603cc3167dfe7712e03dc56c Mon Sep 17 00:00:00 2001 From: Nicky Leach Date: Mon, 14 Sep 2026 20:53:56 -0700 Subject: [PATCH] fix: allow concurrent agent runs on one Anthropic subscription connection (#13445) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## Thinking Path > - Paperclip manages AI agents and their work. > - Agent runs use provider connections and subscription credentials. > - File-backed subscriptions need a lease while a run writes new credentials. > - Anthropic subscriptions pass a token and do not write a credential file. > - The current lease blocks concurrent Anthropic runs without protecting shared state. > - This pull request limits the lease to file-backed subscriptions and tests both paths. > - The result lets Anthropic runs share one connection safely while file-backed credentials remain serialized. ## Linked Issues or Issue Description **What existing behavior does this improve?** This change improves concurrent use of one subscription connection. Anthropic runs no longer wait for a lease that protects no shared file. **Subsystem affected** server/ — REST API and orchestration services, with related server tests and connection documentation. **Current behavior** The server takes a connection lease for every subscription invocation. The lease blocks a second Anthropic run while the first runtime remains open. **Proposed behavior** The server takes the lease only when the subscription writes a provider authentication file back to the grant. Anthropic runs can proceed at the same time. **Reason and benefit** Anthropic passes its token through an environment variable and writes no file. Removing this unnecessary wait improves concurrency and preserves credential safety for OpenAI and xAI. **Breaking changes** None. OpenAI and xAI keep the existing lease and retry behavior. ## What Changed - Limit the AI credential lease to subscriptions that write a provider authentication file. - Add coverage for concurrent Anthropic runs and hired-agent runs. - Keep the existing OpenAI contention assertions as a sensitivity control. - Update the AI connection documentation to describe the lease boundary. ## Verification - Run `server/src/__tests__/ai-connections.test.ts`. - Run `server/src/__tests__/agent-hire-ai-connections.test.ts`. - Run `server/src/__tests__/heartbeat-ai-subscription-contention.test.ts`. - Run the server type check. - Review the pull request checks after GitHub starts continuous integration. ## Risks The main risk is an incorrect provider classification. The write-back condition remains unchanged, and OpenAI contention tests keep the file-backed lease behavior under test. No schema or API contract changes occur. ## Model Used OpenAI GPT-5. Exact runtime model ID and context window are not exposed in this execution. The model used tool calls, shell commands, and GitHub operations. ## 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 --------- Co-authored-by: Paperclip --- doc/connections/AI-CONNECTIONS.md | 15 +++++--- .../agent-hire-ai-connections.test.ts | 34 +++++++++++++++++-- server/src/__tests__/ai-connections.test.ts | 27 ++++++++++++--- server/src/services/ai-connection-runtime.ts | 16 +++++---- 4 files changed, 73 insertions(+), 19 deletions(-) diff --git a/doc/connections/AI-CONNECTIONS.md b/doc/connections/AI-CONNECTIONS.md index baf936fe6b..626baf504b 100644 --- a/doc/connections/AI-CONNECTIONS.md +++ b/doc/connections/AI-CONNECTIONS.md @@ -108,17 +108,22 @@ grant's credentials. Inherited credential variables are cleared. Conflicting project authentication and provider-routing overrides are rejected. Managed failure cannot reactivate host or legacy credentials. -Subscription invocations take a grant-scoped transaction advisory lease. The -reserved database client keeps one transaction open until cleanup, including on +A subscription invocation takes a grant-scoped transaction advisory lease only +when it writes a provider authentication file back to the grant. OpenAI and xAI +subscriptions do this, because their refresh tokens are single use. An +Anthropic subscription invocation writes no file back, so it takes no lease; +two Anthropic invocations of one grant run at the same time. The reserved +database client keeps one transaction open until cleanup, including on transaction-pooling proxies such as PgBouncer. Session-level advisory locks must not be used here: a pooled connection can return to a different backend for cleanup and leave the original lock behind. The lease transaction disables its idle timeout and contains no application data writes; cleanup rolls it back. Two different users' grants can run concurrently; a second invocation of the same -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. +file-backed 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” diff --git a/server/src/__tests__/agent-hire-ai-connections.test.ts b/server/src/__tests__/agent-hire-ai-connections.test.ts index d2586f74d6..c7645acc3c 100644 --- a/server/src/__tests__/agent-hire-ai-connections.test.ts +++ b/server/src/__tests__/agent-hire-ai-connections.test.ts @@ -4,7 +4,7 @@ import os from "node:os"; import path from "node:path"; import express from "express"; import request from "supertest"; -import { eq } from "drizzle-orm"; +import { and, eq } from "drizzle-orm"; import { afterAll, beforeAll, describe, expect, it, vi } from "vitest"; import { agents, companies, companyMemberships, createDb, heartbeatRuns, issues, issueThreadInteractions, principalPermissionGrants, toolConnectionInstalls } from "@paperclipai/db"; import { type AiConnectionBinding } from "@paperclipai/shared"; @@ -202,7 +202,8 @@ describe("agent-created hires use managed AI connections", () => { }); describe("hired agents sharing a subscription", () => { - it.each(["anthropic", "openai"] as const)("waits for the %s parent and resumes without asking for a new connection", async (provider) => { + it("waits for the openai parent and resumes without asking for a new connection", async () => { + const provider = "openai"; const f = await fixture(provider, "subscription"); const agent = hired(await request(f.app).post(`/api/companies/${f.companyId}/agent-hires`).send({ name: "Waiting teammate", role: "engineer", adapterType: f.adapterType, reportsTo: f.agentId, adapterConfig: { cwd: home, engine: "cli" }, runtimeConfig: { heartbeat: { enabled: false } } })); const [issue] = await db.insert(issues).values({ companyId: f.companyId, title: "Subscription child task", status: "todo", assigneeAgentId: agent.id, responsibleUserId: f.userId, createdByUserId: f.userId }).returning(); @@ -239,4 +240,33 @@ describe("hired agents sharing a subscription", () => { unregisterServerAdapter(f.adapterType); } }); + + it("runs the anthropic child alongside a live parent and inherits its connection", async () => { + const provider = "anthropic"; + const f = await fixture(provider, "subscription"); + const agent = hired(await request(f.app).post(`/api/companies/${f.companyId}/agent-hires`).send({ name: "Concurrent teammate", role: "engineer", adapterType: f.adapterType, reportsTo: f.agentId, adapterConfig: { cwd: home, engine: "cli" }, runtimeConfig: { heartbeat: { enabled: false } } })); + const [issue] = await db.insert(issues).values({ companyId: f.companyId, title: "Subscription child task", status: "todo", assigneeAgentId: agent.id, responsibleUserId: f.userId, createdByUserId: f.userId }).returning(); + const parentRuntime = await prepareManagedAiRuntime(db, { companyId: f.companyId, agentId: f.agentId, responsibleUserId: f.userId, adapterType: f.adapterType, binding: f.binding, config: { cwd: home } }); + const execute = vi.fn(async () => { + await db.update(issues).set({ status: "done", completedAt: new Date() }).where(eq(issues.id, issue.id)); + return { exitCode: 0, signal: null, timedOut: false, resultJson: {} }; + }); + registerServerAdapter({ ...getServerAdapter(f.adapterType), execute }); + const heartbeat = heartbeatService(db); + try { + const run = await heartbeat.invoke(agent.id, "assignment", { issueId: issue.id, wakeReason: "issue_assigned", responsibleUserId: f.userId }, "system"); + expect(run).not.toBeNull(); + await expect.poll(async () => (await heartbeat.getRun(run!.id))?.status, { timeout: 20_000 }).toBe("succeeded"); + expect(execute).toHaveBeenCalledTimes(1); + const finished = await heartbeat.getRun(run!.id); + expect(finished?.errorCode).not.toBe("ai_connection_busy"); + expect(finished?.contextSnapshot?.aiConnection).toMatchObject({ connectionId: f.account.connectionId, responsibleUserId: f.userId, method: "subscription" }); + expect(await db.select().from(heartbeatRuns).where(and(eq(heartbeatRuns.agentId, agent.id), eq(heartbeatRuns.scheduledRetryReason, "ai_connection_busy")))).toEqual([]); + } finally { + await parentRuntime.cleanup(); + await db.update(heartbeatRuns).set({ status: "cancelled", finishedAt: new Date() }).where(eq(heartbeatRuns.id, f.runId)); + await heartbeat.drainActiveRunExecutions(); + unregisterServerAdapter(f.adapterType); + } + }); }); diff --git a/server/src/__tests__/ai-connections.test.ts b/server/src/__tests__/ai-connections.test.ts index 6628fedbd5..c36e35ef58 100644 --- a/server/src/__tests__/ai-connections.test.ts +++ b/server/src/__tests__/ai-connections.test.ts @@ -282,12 +282,16 @@ describe("managed AI connections", () => { const [agent] = await db.select().from(agents).where(eq(agents.id, agentId)); expect(agent.runtimeConfig.aiConnection).toBeUndefined(); }); - it("serializes subscription refresh and releases the lease after execution", async () => { - const subscription = { ...input, binding: { ...binding, method: "subscription" as const }, responsibleUserId: "alice", config: { model: "same-model" } }; - const account = (await service.list(companyId, "alice")).find(account => account.provider === "anthropic" && account.method === "subscription")!; - await service.setDefault(companyId, "alice", account.grantId); + it("serializes OpenAI subscription refresh and releases the lease after execution", async () => { + // OpenAI writes its rotated refresh token back to a shared auth file, so + // two runs against the same grant must not overlap. + const userId = "openai-lease-user"; + await db.insert(companyMemberships).values({ companyId, principalId: userId, principalType: "user", status: "active", membershipRole: "member" }); + const token = JSON.stringify({ tokens: { access_token: "fixture-lease-access", refresh_token: "fixture-lease-refresh", id_token: "fixture-lease-id", account_id: "fixture-lease-account" } }); + await service.save(companyId, userId, { provider: "openai", method: "subscription", ownership: "personal", name: "Lease subscription", loginSessionId: "fixture", allAgents: true, agentIds: [] }, token); + const subscription = { ...input, adapterType: "codex_local", binding: { provider: "openai", method: "subscription", mode: "responsible_user" } as const, responsibleUserId: userId, config: { model: "same-model" } }; const first = await prepareManagedAiRuntime(db, subscription); - const selected = await service.select({ ...subscription, userId: "alice" }); + const selected = await service.select({ ...subscription, userId }); const lockKey = `ai-runtime:${selected.grant.id}`; const held = await db.execute(sql` select activity.state, activity.xact_start @@ -308,6 +312,19 @@ describe("managed AI connections", () => { expect(next.identity).toBe(first.identity); await next.cleanup(); }); + it("runs two Claude subscription executions for the same grant at the same time", async () => { + // Claude writes no auth file back to the grant, so two runs share no + // mutable state and must not wait for each other. + const subscription = { ...input, binding: { ...binding, method: "subscription" as const }, responsibleUserId: "alice", config: { model: "same-model" } }; + const account = (await service.list(companyId, "alice")).find(account => account.provider === "anthropic" && account.method === "subscription")!; + await service.setDefault(companyId, "alice", account.grantId); + const [first, second] = await Promise.all([prepareManagedAiRuntime(db, subscription), prepareManagedAiRuntime(db, subscription)]); + try { + expect(second.identity).toBe(first.identity); + } finally { + await Promise.all([first.cleanup(), second.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" }); diff --git a/server/src/services/ai-connection-runtime.ts b/server/src/services/ai-connection-runtime.ts index 8e05a3d79a..0f304a9118 100644 --- a/server/src/services/ai-connection-runtime.ts +++ b/server/src/services/ai-connection-runtime.ts @@ -243,10 +243,15 @@ export async function prepareManagedAiRuntime( runnerProvider: input.config.provider, acpxAgent: input.config.acpxAgent, }); - const release = - selection.attribution.method === "subscription" - ? await acquireCredentialLease(db, selection.grant.id) - : async () => {}; + const subscriptionFile = + selection.attribution.method === "subscription" && + input.binding.provider !== "anthropic"; + // A file-backed subscription runs one credential rotation at a time, so it + // needs the lease. Anthropic writes no file back, so two runs share no + // mutable state and can run at the same time without the lease. + const release = subscriptionFile + ? await acquireCredentialLease(db, selection.grant.id) + : async () => {}; let home: string | undefined; try { const selectedGrantId = selection.grant.id; @@ -291,9 +296,6 @@ export async function prepareManagedAiRuntime( 'cli_auth_credentials_store = "file"\n', { mode: 0o600 }, ); - const subscriptionFile = - selection.attribution.method === "subscription" && - input.binding.provider !== "anthropic"; if (subscriptionFile) await writeFile(authFile, value, { mode: 0o600 }); else env[capability.envKey] = value; if (