mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-06 10:48:12 +02:00
fix: allow concurrent agent runs on one Anthropic subscription connection (#13445)
## 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 <noreply@paperclip.ing>
This commit is contained in:
1 parent
52d120f68d
commit
ee81cee76d
4 files changed
+73
-19
No files matched your search
@@ -108,17 +108,22 @@ grant's credentials. Inherited credential variables are cleared. Conflicting
|
|||||||
project authentication and provider-routing overrides are rejected. Managed
|
project authentication and provider-routing overrides are rejected. Managed
|
||||||
failure cannot reactivate host or legacy credentials.
|
failure cannot reactivate host or legacy credentials.
|
||||||
|
|
||||||
Subscription invocations take a grant-scoped transaction advisory lease. The
|
A subscription invocation takes a grant-scoped transaction advisory lease only
|
||||||
reserved database client keeps one transaction open until cleanup, including on
|
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
|
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
|
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
|
cleanup and leave the original lock behind. The lease transaction disables its
|
||||||
idle timeout and contains no application data writes; cleanup rolls it back.
|
idle timeout and contains no application data writes; cleanup rolls it back.
|
||||||
Two
|
Two
|
||||||
different users' grants can run concurrently; a second invocation of the same
|
different users' grants can run concurrently; a second invocation of the same
|
||||||
subscription receives a retryable busy response while it is in use. Refreshes
|
file-backed subscription receives a retryable busy response while it is in
|
||||||
are merged only into the originating active grant, with reconnect/revocation
|
use. Refreshes are merged only into the originating active grant, with
|
||||||
version checks. Temporary homes are removed on normal completion or failure.
|
reconnect/revocation version checks. Temporary homes are removed on normal
|
||||||
|
completion or failure.
|
||||||
|
|
||||||
For a fresh task execution, subscription contention creates a durable scheduled
|
For a fresh task execution, subscription contention creates a durable scheduled
|
||||||
retry checked every 60–120 seconds. The task shows “Waiting for AI subscription”
|
retry checked every 60–120 seconds. The task shows “Waiting for AI subscription”
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ import os from "node:os";
|
|||||||
import path from "node:path";
|
import path from "node:path";
|
||||||
import express from "express";
|
import express from "express";
|
||||||
import request from "supertest";
|
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 { afterAll, beforeAll, describe, expect, it, vi } from "vitest";
|
||||||
import { agents, companies, companyMemberships, createDb, heartbeatRuns, issues, issueThreadInteractions, principalPermissionGrants, toolConnectionInstalls } from "@paperclipai/db";
|
import { agents, companies, companyMemberships, createDb, heartbeatRuns, issues, issueThreadInteractions, principalPermissionGrants, toolConnectionInstalls } from "@paperclipai/db";
|
||||||
import { type AiConnectionBinding } from "@paperclipai/shared";
|
import { type AiConnectionBinding } from "@paperclipai/shared";
|
||||||
@@ -202,7 +202,8 @@ describe("agent-created hires use managed AI connections", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
describe("hired agents sharing a subscription", () => {
|
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 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 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();
|
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);
|
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);
|
||||||
|
}
|
||||||
|
});
|
||||||
});
|
});
|
||||||
@@ -282,12 +282,16 @@ describe("managed AI connections", () => {
|
|||||||
const [agent] = await db.select().from(agents).where(eq(agents.id, agentId));
|
const [agent] = await db.select().from(agents).where(eq(agents.id, agentId));
|
||||||
expect(agent.runtimeConfig.aiConnection).toBeUndefined();
|
expect(agent.runtimeConfig.aiConnection).toBeUndefined();
|
||||||
});
|
});
|
||||||
it("serializes subscription refresh and releases the lease after execution", async () => {
|
it("serializes OpenAI subscription refresh and releases the lease after execution", async () => {
|
||||||
const subscription = { ...input, binding: { ...binding, method: "subscription" as const }, responsibleUserId: "alice", config: { model: "same-model" } };
|
// OpenAI writes its rotated refresh token back to a shared auth file, so
|
||||||
const account = (await service.list(companyId, "alice")).find(account => account.provider === "anthropic" && account.method === "subscription")!;
|
// two runs against the same grant must not overlap.
|
||||||
await service.setDefault(companyId, "alice", account.grantId);
|
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 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 lockKey = `ai-runtime:${selected.grant.id}`;
|
||||||
const held = await db.execute(sql`
|
const held = await db.execute(sql`
|
||||||
select activity.state, activity.xact_start
|
select activity.state, activity.xact_start
|
||||||
@@ -308,6 +312,19 @@ describe("managed AI connections", () => {
|
|||||||
expect(next.identity).toBe(first.identity);
|
expect(next.identity).toBe(first.identity);
|
||||||
await next.cleanup();
|
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 () => {
|
it("reproduces same-agent OpenAI subscription contention and resumes without reconnecting", async () => {
|
||||||
const userId = "subscription-contention-user";
|
const userId = "subscription-contention-user";
|
||||||
await db.insert(companyMemberships).values({ companyId, principalId: userId, principalType: "user", status: "active", membershipRole: "member" });
|
await db.insert(companyMemberships).values({ companyId, principalId: userId, principalType: "user", status: "active", membershipRole: "member" });
|
||||||
|
|||||||
@@ -243,10 +243,15 @@ export async function prepareManagedAiRuntime(
|
|||||||
runnerProvider: input.config.provider,
|
runnerProvider: input.config.provider,
|
||||||
acpxAgent: input.config.acpxAgent,
|
acpxAgent: input.config.acpxAgent,
|
||||||
});
|
});
|
||||||
const release =
|
const subscriptionFile =
|
||||||
selection.attribution.method === "subscription"
|
selection.attribution.method === "subscription" &&
|
||||||
? await acquireCredentialLease(db, selection.grant.id)
|
input.binding.provider !== "anthropic";
|
||||||
: async () => {};
|
// 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;
|
let home: string | undefined;
|
||||||
try {
|
try {
|
||||||
const selectedGrantId = selection.grant.id;
|
const selectedGrantId = selection.grant.id;
|
||||||
@@ -291,9 +296,6 @@ export async function prepareManagedAiRuntime(
|
|||||||
'cli_auth_credentials_store = "file"\n',
|
'cli_auth_credentials_store = "file"\n',
|
||||||
{ mode: 0o600 },
|
{ mode: 0o600 },
|
||||||
);
|
);
|
||||||
const subscriptionFile =
|
|
||||||
selection.attribution.method === "subscription" &&
|
|
||||||
input.binding.provider !== "anthropic";
|
|
||||||
if (subscriptionFile) await writeFile(authFile, value, { mode: 0o600 });
|
if (subscriptionFile) await writeFile(authFile, value, { mode: 0o600 });
|
||||||
else env[capability.envKey] = value;
|
else env[capability.envKey] = value;
|
||||||
if (
|
if (
|
||||||
|
|||||||
Reference in new issue
Block a user