Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions plugins/account-pool/PLUGIN_OVERVIEW.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ The hub runs inside BB and serves an Anthropic Messages endpoint and an OpenAI R

The pool waits once on the same account for short temporary rate limits. Longer holds return Retry-After for pinned conversations while new conversations can advance. A model-family limit detours requests for that family without moving the session’s main pin or the provider cursor. The pool commits a new account after a successful response; a failed attempt across every account retains the previous binding. The current account and session pins survive hub restarts. Session pins expire after 30 idle minutes, with the 4,096 most recently used pins retained.

The pooler owns its upstream HTTP connections and uses HTTP/1.1, so a broken HTTP/2 session in the server's shared fetch dispatcher does not strand pooled requests. The transport honors standard proxy environment variables and is disposed on plugin unload. This does not add request replay; existing account-fallback rules still apply. Pooled request connection failures log a known error code when available, without request bodies, credentials, URLs, or raw exception messages.

## Requirements

Accounts you own and are permitted to use this way.
Expand Down
3 changes: 2 additions & 1 deletion plugins/account-pool/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,8 @@
"@dnd-kit/core": "^6.3.1",
"@dnd-kit/sortable": "^10.0.0",
"@dnd-kit/utilities": "^3.2.2",
"zod": "^4.3.6"
"zod": "^4.3.6",
"undici": "^7.28.0"
},
"devDependencies": {
"@get-bb/plugin-sdk": "workspace:*",
Expand Down
11 changes: 11 additions & 0 deletions plugins/account-pool/src/fixtures/localhost-cert.pem
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
-----BEGIN CERTIFICATE-----
MIIBkDCCATagAwIBAgIUPSOkhyOF0idrJUTyAEsjAYBo4SwwCgYIKoZIzj0EAwIw
FDESMBAGA1UEAwwJbG9jYWxob3N0MCAXDTI2MDkwNzA2NTYyOFoYDzIxMjYwODE0
MDY1NjI4WjAUMRIwEAYDVQQDDAlsb2NhbGhvc3QwWTATBgcqhkjOPQIBBggqhkjO
PQMBBwNCAARijKjFmrSp6Lo5hhV+UtpTroyDeviu4O1g3z3XCxINhzQD3a8jCw3j
sjZmZXfmQ1SDbxur5SVL8jQTAJEi1tbEo2QwYjAdBgNVHQ4EFgQU13ZF79SMLdUJ
bQ4eDMUO9J6FD3MwHwYDVR0jBBgwFoAU13ZF79SMLdUJbQ4eDMUO9J6FD3MwDwYD
VR0TAQH/BAUwAwEB/zAPBgNVHREECDAGhwR/AAABMAoGCCqGSM49BAMCA0gAMEUC
IQDU0voLBRfkdTQckCqThNZdTZ9OLYCQS5SgEuxAKm553AIgJzQAyVF0h3aG/5Tn
AlYMphzKC900nxDQcDcf5DvjrY0=
-----END CERTIFICATE-----
5 changes: 5 additions & 0 deletions plugins/account-pool/src/fixtures/localhost-key.pem
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
-----BEGIN PRIVATE KEY-----
MIGHAgEAMBMGByqGSM49AgEGCCqGSM49AwEHBG0wawIBAQQg24rTCfBWq56U/l3T
1WM+dgVzYV906qZqMy+OtQbABiqhRANCAARijKjFmrSp6Lo5hhV+UtpTroyDeviu
4O1g3z3XCxINhzQD3a8jCw3jsjZmZXfmQ1SDbxur5SVL8jQTAJEi1tbE
-----END PRIVATE KEY-----
11 changes: 9 additions & 2 deletions plugins/account-pool/src/hub.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@ interface HubOptions {
usageRefreshIntervalMs: number;
drainTimeoutMs: number;
onAccountsChanged: () => void;
onUpstreamError: (provider: PoolProvider, error: unknown) => void;
}

interface SelectedAccount {
Expand Down Expand Up @@ -978,8 +979,12 @@ export class AccountPoolHub {
: { body: upstreamBody }),
signal: controller.signal,
})
.catch(() => {
throw new UpstreamConnectionError("Upstream connection failed.");
.catch((cause: unknown) => {
if (!controller.signal.aborted)
this.options.onUpstreamError(adapter.provider, cause);
throw new UpstreamConnectionError("Upstream connection failed.", {
cause,
});
});
return { response, controller, release };
} catch (error) {
Expand Down Expand Up @@ -1135,6 +1140,7 @@ export function createHub(options: {
usageRefreshIntervalMs?: number;
drainTimeoutMs?: number;
onAccountsChanged?: () => void;
onUpstreamError?: (provider: PoolProvider, error: unknown) => void;
}): AccountPoolHub {
const adapters: ReadonlyMap<PoolProvider, ProviderAdapter> = new Map([
[
Expand Down Expand Up @@ -1168,6 +1174,7 @@ export function createHub(options: {
options.usageRefreshIntervalMs ?? DEFAULT_USAGE_REFRESH_INTERVAL_MS,
drainTimeoutMs: options.drainTimeoutMs ?? 60_000,
onAccountsChanged: options.onAccountsChanged ?? (() => {}),
onUpstreamError: options.onUpstreamError ?? (() => {}),
});
}

Expand Down
97 changes: 97 additions & 0 deletions plugins/account-pool/src/server.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,7 @@ async function createFixture(args: {
source?: "api-key" | "import";
apiKey?: string;
priority?: number;
beforePlugin?: (host: Fixture["host"]) => void;
}): Promise<Fixture> {
const dataDir = await mkdtemp(path.join(tmpdir(), "bb-account-pool-"));
const host = createFakePluginHost({
Expand All @@ -199,6 +200,7 @@ async function createFixture(args: {
usageUrl: "data:application/json,{}",
...args.options,
});
args.beforePlugin?.(host);
await plugin(host.bb);
const accountMetadata = accountSchema.parse(
await host.harness.behavior.callRpc("account.add", {
Expand Down Expand Up @@ -5760,3 +5762,98 @@ describe("sequential pool recovery", () => {
).toBe(second.id);
});
});

it("logs a sanitized transport cause when pooled fetch fails", async () => {
const fixture = await createFixture({
upstreamUrl: "https://upstream.example",
options: {
fetch: async () => {
throw new TypeError("fetch failed with private request data", {
cause: Object.assign(
new Error("The session has been destroyed: secret-token"),
{
code: "ERR_HTTP2_INVALID_SESSION",
},
),
});
},
},
});
const response = await fixture.host.harness.behavior.fetchHttp(
"POST",
"/v1/messages",
{
headers: authHeaders(fixture.key),
body: "{}",
},
);
expect(response.status).toBe(502);
expect(await response.text()).not.toContain("secret-token");
expect(fixture.host.harness.inspection.logEntries).toContainEqual({
level: "warn",
message:
"Account Pooler claude transport failed: ERR_HTTP2_INVALID_SESSION.",
});
expect(
JSON.stringify(fixture.host.harness.inspection.logEntries),
).not.toContain("secret-token");
expect(
JSON.stringify(fixture.host.harness.inspection.logEntries),
).not.toContain("private request data");
});

it("drains a streamed response before disposing the owned transport", async () => {
const finish = deferred();
const upstream = await startUpstream(async (request, response) => {
await readRequestBody(request);
response.writeHead(200, { "content-type": "text/event-stream" });
response.write("first");
await finish.promise;
response.end("last");
});
cleanups.push(upstream.close);
const hooks: Array<() => void | Promise<void>> = [];
const fixture = await createFixture({
upstreamUrl: upstream.url,
beforePlugin(host) {
const register = host.bb.onDispose.bind(host.bb);
vi.spyOn(host.bb, "onDispose").mockImplementation((hook) => {
hooks.push(hook);
register(hook);
});
},
});
const response = await fixture.host.harness.behavior.fetchHttp(
"POST",
"/v1/messages",
{
headers: authHeaders(fixture.key),
body: "{}",
},
);
const reader = response.body?.getReader();
if (reader === undefined) throw new Error("Expected a stream");
expect(new TextDecoder().decode((await reader.read()).value)).toBe("first");
const disposeTransport = hooks[0];
if (disposeTransport === undefined)
throw new Error("Expected transport disposal");
const tail = reader.read().then(
(result) => ({
kind: "chunk",
text: new TextDecoder().decode(result.value),
}),
() => ({ kind: "error", text: "" }),
);
const disposing = disposeTransport();
try {
await new Promise<void>((resolve) => setImmediate(resolve));
finish.resolve();
expect(await tail).toEqual({ kind: "chunk", text: "last" });
expect((await reader.read()).done).toBe(true);
await disposing;
} finally {
finish.resolve();
await reader.cancel().catch(() => undefined);
await disposing;
}
});
23 changes: 20 additions & 3 deletions plugins/account-pool/src/server.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
import {
createUpstreamTransport,
transportErrorCode,
} from "./upstream-transport.js";
import path from "node:path";
import type { BbPluginApi } from "@get-bb/plugin-sdk";
import { registerPoolCli } from "./cli.js";
Expand Down Expand Up @@ -93,13 +97,16 @@ export function createAccountPoolPlugin(
const db = bb.storage.database();
bb.storage.migrate(db, QUOTA_MIGRATIONS);
const quotas = new QuotaStore(db);
const transport =
options.fetch === undefined ? createUpstreamTransport() : null;
const upstreamFetch = options.fetch ?? transport?.fetch;
const hub = createHub({
accounts,
quotas,
affinity: new PoolAffinityStore(db),
hubTokens,
getSettings: () => currentSettings,
fetch: options.fetch,
fetch: upstreamFetch,
now,
refreshUrl: options.refreshUrl,
codexRefreshUrl: options.codexRefreshUrl,
Expand All @@ -110,9 +117,19 @@ export function createAccountPoolPlugin(
importCodexCredentials: options.importCodexCredentials,
usageRefreshIntervalMs: options.usageRefreshIntervalMs,
drainTimeoutMs: options.drainTimeoutMs,
onUpstreamError: (provider, error) =>
bb.log.warn(
`Account Pooler ${provider} transport failed: ${transportErrorCode(error)}.`,
),
onAccountsChanged: () =>
bb.realtime.publish(ACCOUNT_POOL_ACCOUNTS_CHANGED, {}),
});
if (transport !== null) {
bb.onDispose(async () => {
await hub.stop();
await transport.destroy();
});
}
const operations = new PoolOperations(
accounts,
quotas,
Expand All @@ -127,15 +144,15 @@ export function createAccountPoolPlugin(
(accountId) => hub.refreshUsage(accountId, true),
);
const login = new ClaudeOAuthLogin({
fetch: options.fetch,
fetch: upstreamFetch,
now,
authorizeUrl: options.oauthAuthorizeUrl,
tokenUrl: options.oauthTokenUrl,
profileUrl: options.oauthProfileUrl,
addAccount: (authenticated) => operations.addOAuth(authenticated),
});
const codexLogin = new CodexDeviceLogin({
fetch: options.fetch,
fetch: upstreamFetch,
now,
authBaseUrl: options.codexAuthBaseUrl,
addAccount: (authenticated) => operations.addCodexOAuth(authenticated),
Expand Down
Loading
Loading