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 (