Merge PR #6744: perf(ai): lazily construct auth broker wire schemas (@usr-bin-roygbiv)

This commit is contained in:
can1357
2026-07-27 04:58:25 +02:00
8 changed files with 937 additions and 413 deletions
+31 -32
View File
@@ -28,22 +28,21 @@ import type {
UsageResponse,
UsageStaleResponse,
} from "./types";
import {
clientUsageReportResponseSchema,
clientUsageSummaryResponseSchema,
credentialBlockResponseSchema,
credentialBlocksDeleteResponseSchema,
credentialDisableResponseSchema,
credentialRefreshResponseSchema,
credentialUploadResponseSchema,
disabledCredentialsResponseSchema,
healthzResponseSchema,
snapshotResponseSchema,
snapshotStreamEventSchema,
usageHistoryResponseSchema,
usageResponseSchema,
usageStaleResponseSchema,
} from "./wire-schemas";
import { getAuthBrokerWireSchemas } from "./wire-schema-resource";
type AuthBrokerResponseSchemaName =
| "clientUsageReportResponseSchema"
| "clientUsageSummaryResponseSchema"
| "credentialBlockResponseSchema"
| "credentialBlocksDeleteResponseSchema"
| "credentialDisableResponseSchema"
| "credentialRefreshResponseSchema"
| "credentialUploadResponseSchema"
| "disabledCredentialsResponseSchema"
| "healthzResponseSchema"
| "usageHistoryResponseSchema"
| "usageResponseSchema"
| "usageStaleResponseSchema";
export interface AuthBrokerClientOptions {
/** Base URL (e.g. `https://broker.tailnet:8765`). Trailing slashes are trimmed. */
@@ -123,7 +122,7 @@ export class AuthBrokerClient {
healthz(signal?: AbortSignal): Promise<HealthzResponse> {
return this.#request<HealthzResponse>("GET", "/v1/healthz", {
schema: healthzResponseSchema,
schema: "healthzResponseSchema",
auth: false,
signal,
});
@@ -152,7 +151,7 @@ export class AuthBrokerClient {
}
const text = await response.text();
const raw = this.#parseJson(text, response.status);
const validated = snapshotResponseSchema(raw);
const validated = getAuthBrokerWireSchemas().snapshotResponseSchema(raw);
if (validated instanceof type.errors) {
throw new AuthBrokerError("Auth broker response failed schema validation", {
status: response.status,
@@ -220,7 +219,7 @@ export class AuthBrokerClient {
cause: err,
});
}
const validated = snapshotStreamEventSchema(parsed);
const validated = getAuthBrokerWireSchemas().snapshotStreamEventSchema(parsed);
if (validated instanceof type.errors) {
throw new AuthBrokerError("Auth broker stream event failed schema validation", {
body: validated.summary,
@@ -250,7 +249,7 @@ export class AuthBrokerClient {
// `metadata`) but leaves provider-specific extension fields permissive so
// the broker can ship new shapes ahead of the client. `raw` is accepted
// but normally stripped by the broker before send.
return this.#request<UsageResponse>("GET", "/v1/usage", { schema: usageResponseSchema, signal });
return this.#request<UsageResponse>("GET", "/v1/usage", { schema: "usageResponseSchema", signal });
}
/** Recorded usage-limit snapshots from the broker host, oldest first. */
@@ -262,14 +261,14 @@ export class AuthBrokerClient {
if (query?.sinceMs !== undefined) params.set("sinceMs", String(query.sinceMs));
if (query?.provider) params.set("provider", query.provider);
const path = `/v1/usage/history${params.size > 0 ? `?${params.toString()}` : ""}`;
return this.#request<UsageHistoryResponse>("GET", path, { schema: usageHistoryResponseSchema, signal });
return this.#request<UsageHistoryResponse>("GET", path, { schema: "usageHistoryResponseSchema", signal });
}
/** Report this client's batched observed request usage for per-install burn tracking. */
reportClientUsage(report: ClientUsageReportRequest, signal?: AbortSignal): Promise<ClientUsageReportResponse> {
return this.#request<ClientUsageReportResponse>("POST", "/v1/usage/observed", {
body: report,
schema: clientUsageReportResponseSchema,
schema: "clientUsageReportResponseSchema",
signal,
});
}
@@ -280,21 +279,21 @@ export class AuthBrokerClient {
if (query?.sinceMs !== undefined) params.set("sinceMs", String(query.sinceMs));
const path = `/v1/usage/clients${params.size > 0 ? `?${params.toString()}` : ""}`;
return this.#request<ClientUsageSummaryResponse>("GET", path, {
schema: clientUsageSummaryResponseSchema,
schema: "clientUsageSummaryResponseSchema",
signal,
});
}
notifyUsageStale(signal?: AbortSignal): Promise<UsageStaleResponse> {
return this.#request<UsageStaleResponse>("POST", "/v1/usage/stale", {
schema: usageStaleResponseSchema,
schema: "usageStaleResponseSchema",
signal,
});
}
async refreshCredential(id: number, signal?: AbortSignal): Promise<CredentialRefreshResponse> {
return this.#request<CredentialRefreshResponse>("POST", `/v1/credential/${id}/refresh`, {
schema: credentialRefreshResponseSchema,
schema: "credentialRefreshResponseSchema",
signal,
});
}
@@ -303,7 +302,7 @@ export class AuthBrokerClient {
const body: CredentialDisableRequest = { cause };
return this.#request<CredentialDisableResponse>("POST", `/v1/credential/${id}/disable`, {
body,
schema: credentialDisableResponseSchema,
schema: "credentialDisableResponseSchema",
signal,
});
}
@@ -319,7 +318,7 @@ export class AuthBrokerClient {
const path = `/v1/credentials/disabled${params.size > 0 ? `?${params.toString()}` : ""}`;
try {
const response = await this.#request<DisabledCredentialsResponse>("GET", path, {
schema: disabledCredentialsResponseSchema,
schema: "disabledCredentialsResponseSchema",
signal,
});
return response.disabled;
@@ -337,7 +336,7 @@ export class AuthBrokerClient {
const body: CredentialUploadRequest = { provider, credential };
return this.#request<CredentialUploadResponse>("POST", "/v1/credential", {
body,
schema: credentialUploadResponseSchema,
schema: "credentialUploadResponseSchema",
signal,
});
}
@@ -350,14 +349,14 @@ export class AuthBrokerClient {
const body: CredentialBlockRequest = block;
return this.#request<CredentialBlockResponse>("POST", `/v1/credential/${id}/block`, {
body,
schema: credentialBlockResponseSchema,
schema: "credentialBlockResponseSchema",
signal,
});
}
async deleteCredentialBlocks(id: number, signal?: AbortSignal): Promise<CredentialBlocksDeleteResponse> {
return this.#request<CredentialBlocksDeleteResponse>("DELETE", `/v1/credential/${id}/blocks`, {
schema: credentialBlocksDeleteResponseSchema,
schema: "credentialBlocksDeleteResponseSchema",
signal,
});
}
@@ -365,12 +364,12 @@ export class AuthBrokerClient {
async #request<t>(
method: "GET" | "POST" | "DELETE",
path: string,
opts: { schema: (input: unknown) => unknown; auth?: boolean; body?: unknown; signal?: AbortSignal },
opts: { schema: AuthBrokerResponseSchemaName; auth?: boolean; body?: unknown; signal?: AbortSignal },
): Promise<t> {
const response = await this.#fetchRaw(method, path, opts);
const text = await response.text();
const raw = this.#parseJson(text, response.status);
const validated = opts.schema(raw);
const validated = getAuthBrokerWireSchemas()[opts.schema](raw);
if (validated instanceof type.errors) {
throw new AuthBrokerError("Auth broker response failed schema validation", {
status: response.status,
+7 -10
View File
@@ -38,12 +38,7 @@ import {
DEFAULT_SERVER_IDLE_TIMEOUT_S,
DEFAULT_STREAM_KEEPALIVE_MS,
} from "./types";
import {
clientUsageReportRequestSchema,
credentialBlockRequestSchema,
credentialDisableRequestSchema,
credentialUploadRequestSchema,
} from "./wire-schemas";
import { getAuthBrokerWireSchemas } from "./wire-schema-resource";
export interface AuthBrokerServerOptions {
/** Underlying credential storage (wraps the local SQLite store on the broker). */
@@ -622,7 +617,7 @@ export function startAuthBroker(opts: AuthBrokerServerOptions): AuthBrokerServer
return json(200, { generatedAt: Date.now(), entries });
}
if (req.method === "POST" && pathname === "/v1/usage/observed") {
const parsed = await parseBody(req, clientUsageReportRequestSchema);
const parsed = await parseBody(req, getAuthBrokerWireSchemas().clientUsageReportRequestSchema);
if (!parsed.ok) return parsed.response;
// Arktype's inferred union collides the `entries` field with
// Array.prototype.entries; the schema already validated the shape.
@@ -689,7 +684,9 @@ export function startAuthBroker(opts: AuthBrokerServerOptions): AuthBrokerServer
const disableMatch = req.method === "POST" ? pathname.match(DISABLE_ROUTE) : null;
if (disableMatch) {
const id = Number.parseInt(disableMatch[1], 10);
const parsed = await parseBody(req, credentialDisableRequestSchema, { allowEmpty: true });
const parsed = await parseBody(req, getAuthBrokerWireSchemas().credentialDisableRequestSchema, {
allowEmpty: true,
});
if (!parsed.ok) return parsed.response;
const cause =
parsed.data.cause && parsed.data.cause.length > 0 ? parsed.data.cause : "disabled via auth-broker";
@@ -705,7 +702,7 @@ export function startAuthBroker(opts: AuthBrokerServerOptions): AuthBrokerServer
const blockMatch = req.method === "POST" ? pathname.match(BLOCK_ROUTE) : null;
if (blockMatch) {
const id = Number.parseInt(blockMatch[1], 10);
const parsed = await parseBody(req, credentialBlockRequestSchema);
const parsed = await parseBody(req, getAuthBrokerWireSchemas().credentialBlockRequestSchema);
if (!parsed.ok) return parsed.response;
const block: StoredCredentialBlock = {
credentialId: id,
@@ -755,7 +752,7 @@ export function startAuthBroker(opts: AuthBrokerServerOptions): AuthBrokerServer
}
}
if (req.method === "POST" && pathname === "/v1/credential") {
const parsed = await parseBody(req, credentialUploadRequestSchema);
const parsed = await parseBody(req, getAuthBrokerWireSchemas().credentialUploadRequestSchema);
if (!parsed.ok) return parsed.response;
const { provider, credential } = parsed.data;
try {
@@ -0,0 +1,487 @@
/**
* ArkType schemas for the auth-broker wire protocol.
*
* Shared between the server (validates inbound request bodies) and the client
* (validates responses from the broker). Schemas mirror the TypeScript types
* in `./types.ts` 1:1; the types remain the source of truth for static typing,
* and `Type` is asserted-compatible with them where possible.
*
* Envelope and fixed-shape schemas use `"+": "reject"` so unknown keys are
* rejected — the previous implementation used a hand-rolled `hasOnlyFields`
* allowlist for the same effect. The OAuth credential schema is the deliberate
* exception (standard type keeps extra keys): it preserves provider-specific extension fields so
* they round-trip through the broker instead of being dropped (see below).
*/
import { type Type, type } from "arktype";
import {
type ApiKeyCredential,
type AuthCredential,
type AuthCredentialSnapshotEntry,
type DisabledCredentialSummary,
type OAuthCredential,
REMOTE_REFRESH_SENTINEL,
type RemoteOAuthCredential,
type SnapshotCredential,
} from "../auth-storage";
import type {
ClientUsageReportRequest,
ClientUsageReportResponse,
ClientUsageSummaryResponse,
CredentialBlockRequest,
CredentialBlockResponse,
CredentialBlockSnapshot,
CredentialBlocksDeleteResponse,
CredentialDisableResponse,
CredentialRefreshResponse,
CredentialUploadRequest,
CredentialUploadResponse,
DisabledCredentialsResponse,
HealthzResponse,
RefresherSchedule,
SnapshotEntry,
SnapshotResponse,
SnapshotStreamEntryEvent,
SnapshotStreamEvent,
SnapshotStreamRemovedEvent,
SnapshotStreamSnapshotEvent,
UsageHistoryResponse,
UsageResponse,
UsageStaleResponse,
} from "./types";
export interface AuthBrokerWireSchemas {
readonly oauthCredentialSchema: Type<OAuthCredential>;
readonly remoteOauthCredentialSchema: Type<RemoteOAuthCredential>;
readonly apiKeyCredentialSchema: Type<ApiKeyCredential>;
readonly writableAuthCredentialSchema: Type<AuthCredential>;
readonly snapshotCredentialSchema: Type<SnapshotCredential>;
readonly credentialSnapshotEntrySchema: Type<AuthCredentialSnapshotEntry>;
readonly credentialBlockSnapshotSchema: Type<CredentialBlockSnapshot>;
readonly snapshotEntrySchema: Type<SnapshotEntry>;
readonly refresherScheduleSchema: Type<RefresherSchedule>;
readonly snapshotResponseSchema: Type<SnapshotResponse>;
readonly snapshotStreamSnapshotEventSchema: Type<SnapshotStreamSnapshotEvent>;
readonly snapshotStreamEntryEventSchema: Type<SnapshotStreamEntryEvent>;
readonly snapshotStreamRemovedEventSchema: Type<SnapshotStreamRemovedEvent>;
readonly snapshotStreamEventSchema: Type<SnapshotStreamEvent>;
readonly healthzResponseSchema: Type<HealthzResponse>;
readonly usageResponseSchema: Type<UsageResponse>;
readonly usageHistoryResponseSchema: Type<UsageHistoryResponse>;
readonly clientUsageReportRequestSchema: Type<ClientUsageReportRequest>;
readonly clientUsageReportResponseSchema: Type<ClientUsageReportResponse>;
readonly clientUsageSummaryResponseSchema: Type<ClientUsageSummaryResponse>;
readonly credentialRefreshResponseSchema: Type<CredentialRefreshResponse>;
readonly credentialDisableRequestSchema: Type<{ cause?: string }>;
readonly credentialDisableResponseSchema: Type<CredentialDisableResponse>;
readonly disabledCredentialSummarySchema: Type<DisabledCredentialSummary>;
readonly disabledCredentialsResponseSchema: Type<DisabledCredentialsResponse>;
readonly credentialBlockRequestSchema: Type<CredentialBlockRequest>;
readonly credentialBlockResponseSchema: Type<CredentialBlockResponse>;
readonly credentialBlocksDeleteResponseSchema: Type<CredentialBlocksDeleteResponse>;
readonly usageStaleResponseSchema: Type<UsageStaleResponse>;
readonly credentialUploadRequestSchema: Type<CredentialUploadRequest>;
readonly credentialUploadResponseSchema: Type<CredentialUploadResponse>;
}
function buildAuthBrokerWireSchemas(): AuthBrokerWireSchemas {
// ─── Credential payloads ───────────────────────────────────────────────────
/** Real OAuth credential (broker-side) — refresh token is the actual upstream value. */
const oauthCredentialSchema = type({
"apiEndpoint?": "string",
type: "'oauth'",
refresh: type("string").narrow(
(value, ctx) =>
value !== REMOTE_REFRESH_SENTINEL ||
ctx.mustBe(`not equal to the remote sentinel (${REMOTE_REFRESH_SENTINEL})`),
),
access: type("string").atLeastLength(1),
expires: "number",
"enterpriseUrl?": "string",
"projectId?": "string",
"email?": "string",
"accountId?": "string",
"orgId?": "string",
"orgName?": "string",
"authorizedAt?": "number",
});
/** OAuth credential as it appears in broker snapshots — refresh replaced with sentinel. */
const remoteOauthCredentialSchema = type({
"apiEndpoint?": "string",
type: "'oauth'",
refresh: type.enumerated(REMOTE_REFRESH_SENTINEL),
access: type("string").atLeastLength(1),
expires: "number",
"enterpriseUrl?": "string",
"projectId?": "string",
"email?": "string",
"accountId?": "string",
"orgId?": "string",
"orgName?": "string",
"authorizedAt?": "number",
});
const apiKeyCredentialSchema = type({
"+": "reject",
type: "'api_key'",
key: type("string").atLeastLength(1),
"source?": "'login'",
});
/** Discriminated union accepted on POST /v1/credential (writes). */
const writableAuthCredentialSchema = oauthCredentialSchema.or(apiKeyCredentialSchema);
/** Discriminated union returned in snapshots (refresh is sentinel for OAuth). */
const snapshotCredentialSchema = remoteOauthCredentialSchema.or(apiKeyCredentialSchema);
// ─── Snapshot ──────────────────────────────────────────────────────────────
const credentialSnapshotEntrySchema = type({
"+": "reject",
id: "number.integer",
provider: type("string").atLeastLength(1),
credential: snapshotCredentialSchema,
identityKey: "string | null",
});
const credentialBlockSnapshotSchema = type({
"+": "reject",
providerKey: type("string").atLeastLength(1),
blockScope: "string",
blockedUntilMs: "number",
"updatedAtMs?": "number",
});
const snapshotEntrySchema = type({
"+": "reject",
id: "number.integer",
provider: type("string").atLeastLength(1),
credential: snapshotCredentialSchema,
identityKey: "string | null",
rotatesInMs: "number | null",
"blocks?": credentialBlockSnapshotSchema.array(),
});
const refresherScheduleSchema = type({
"+": "reject",
enabled: "boolean",
intervalMs: "number",
skewMs: "number",
nextSweepInMs: "number",
});
const snapshotResponseSchema = type({
"+": "reject",
generation: "number.integer",
generatedAt: "number",
serverNowMs: "number",
refresher: refresherScheduleSchema,
credentials: snapshotEntrySchema.array(),
});
// ─── Snapshot stream (SSE) ────────────────────────────────────────────────
/** First frame on connect — full snapshot embedded inline with a `kind` tag. */
const snapshotStreamSnapshotEventSchema = type({
"+": "reject",
generation: "number.integer",
generatedAt: "number",
serverNowMs: "number",
refresher: refresherScheduleSchema,
credentials: snapshotEntrySchema.array(),
kind: "'snapshot'",
});
/** Per-credential upsert/refresh delta. */
const snapshotStreamEntryEventSchema = type({
"+": "reject",
kind: "'entry'",
generation: "number.integer",
serverNowMs: "number",
refresher: refresherScheduleSchema,
entry: snapshotEntrySchema,
});
/** Per-credential delete delta. */
const snapshotStreamRemovedEventSchema = type({
"+": "reject",
kind: "'removed'",
generation: "number.integer",
serverNowMs: "number",
refresher: refresherScheduleSchema,
id: "number.integer",
});
/** Discriminated union over every event frame the snapshot stream emits. */
const snapshotStreamEventSchema = snapshotStreamSnapshotEventSchema
.or(snapshotStreamEntryEventSchema)
.or(snapshotStreamRemovedEventSchema);
// ─── Healthz ────────────────────────────────────────────────────────────────
const healthzResponseSchema = type({
"+": "reject",
ok: "boolean",
"version?": "string",
});
// ─── Usage ─────────────────────────────────────────────────────────────────
const usageUnitSchema = type("'percent' | 'tokens' | 'requests' | 'usd' | 'minutes' | 'bytes' | 'unknown'");
const usageStatusSchema = type("'ok' | 'warning' | 'exhausted' | 'unknown'");
const usageWindowSchema = type({
id: "string",
label: "string",
"durationMs?": "number",
"resetsAt?": "number",
});
const usageAmountSchema = type({
"used?": "number",
"limit?": "number",
"remaining?": "number",
"usedFraction?": "number",
"remainingFraction?": "number",
unit: usageUnitSchema,
});
const usageScopeSchema = type({
provider: "string",
"accountId?": "string",
"projectId?": "string",
"orgId?": "string",
"modelId?": "string",
"tier?": "string",
"windowId?": "string",
"shared?": "boolean",
});
const usageLimitSchema = type({
id: "string",
label: "string",
scope: usageScopeSchema,
"window?": usageWindowSchema,
amount: usageAmountSchema,
"status?": usageStatusSchema,
"notes?": "string[]",
});
const usageResetCreditDetailSchema = type({
"grantedAt?": "string",
"expiresAt?": "string",
"status?": "string",
});
const usageResetCreditsSchema = type({
availableCount: "number",
"credits?": usageResetCreditDetailSchema.array(),
});
const arkUsageReportSchema = type({
provider: "string",
fetchedAt: "number",
limits: usageLimitSchema.array(),
"resetCredits?": usageResetCreditsSchema,
"notes?": "string[]",
"metadata?": { "[string]": "unknown" },
"raw?": "unknown",
});
/**
* Broker `/v1/usage` response. Reports are full {@link UsageReport}s minus the
* heavy provider-specific `raw` field (the server strips it before send) — we
* keep `raw` optional in the underlying schema so a misconfigured broker that
* forgot to strip still validates.
*/
const usageResponseSchema = type({
"+": "reject",
generatedAt: "number",
reports: arkUsageReportSchema.array(),
});
const usageHistoryEntrySchema = type({
recordedAt: "number",
provider: "string",
accountKey: "string",
"email?": "string",
"accountId?": "string",
limitId: "string",
label: "string",
"windowLabel?": "string",
"usedFraction?": "number",
"status?": "'ok' | 'warning' | 'exhausted' | 'unknown'",
"resetsAt?": "number",
});
/** Broker `/v1/usage/history` response — recorded usage-limit snapshots, oldest first. */
const usageHistoryResponseSchema = type({
"+": "reject",
generatedAt: "number",
entries: usageHistoryEntrySchema.array(),
});
const observedUsageEntrySchema = type({
at: "number",
provider: "string",
model: "string",
requests: "number",
inputTokens: "number",
outputTokens: "number",
cacheReadTokens: "number",
cacheWriteTokens: "number",
costUsd: "number",
});
/** Broker `POST /v1/usage/observed` request — one client's batched observed usage. */
const clientUsageReportRequestSchema = type({
"+": "reject",
installId: "string",
"hostname?": "string",
entries: observedUsageEntrySchema.array(),
});
const clientUsageReportResponseSchema = type({
"+": "reject",
ok: "boolean",
});
const clientProviderUsageSchema = type({
provider: "string",
requests: "number",
inputTokens: "number",
outputTokens: "number",
cacheReadTokens: "number",
cacheWriteTokens: "number",
costUsd: "number",
});
const clientUsageClientSummarySchema = type({
installId: "string",
"hostname?": "string",
firstSeen: "number",
lastSeen: "number",
providers: clientProviderUsageSchema.array(),
});
/** Broker `GET /v1/usage/clients` response — per-client token burn aggregates. */
const clientUsageSummaryResponseSchema = type({
"+": "reject",
generatedAt: "number",
clients: clientUsageClientSummarySchema.array(),
});
// ─── Refresh ───────────────────────────────────────────────────────────────
const credentialRefreshResponseSchema = type({
"+": "reject",
entry: credentialSnapshotEntrySchema,
});
// ─── Disable ───────────────────────────────────────────────────────────────
const credentialDisableRequestSchema = type({
"+": "reject",
"cause?": "string",
});
const credentialDisableResponseSchema = type({
"+": "reject",
ok: "boolean",
});
/** One disabled-credential tombstone — identity + cause, never token material. */
const disabledCredentialSummarySchema = type({
"+": "reject",
id: "number.integer",
provider: type("string").atLeastLength(1),
type: "'oauth' | 'api_key'",
"email?": "string",
"accountId?": "string",
"orgId?": "string",
"orgName?": "string",
cause: "string",
"disabledAtMs?": "number",
});
/** Broker `GET /v1/credentials/disabled` response. */
const disabledCredentialsResponseSchema = type({
"+": "reject",
generatedAt: "number",
disabled: disabledCredentialSummarySchema.array(),
});
// ─── Credential blocks ──────────────────────────────────────────────────────
const credentialBlockRequestSchema = credentialBlockSnapshotSchema;
const credentialBlockResponseSchema = type({
"+": "reject",
ok: "boolean",
});
const credentialBlocksDeleteResponseSchema = type({
"+": "reject",
ok: "boolean",
});
const usageStaleResponseSchema = type({
"+": "reject",
ok: "boolean",
});
// ─── Upload ────────────────────────────────────────────────────────────────
const credentialUploadRequestSchema = type({
"+": "reject",
provider: type("string").atLeastLength(1),
credential: writableAuthCredentialSchema,
});
const credentialUploadResponseSchema = type({
"+": "reject",
entries: credentialSnapshotEntrySchema.array(),
});
return {
oauthCredentialSchema,
remoteOauthCredentialSchema,
apiKeyCredentialSchema,
writableAuthCredentialSchema,
snapshotCredentialSchema,
credentialSnapshotEntrySchema,
credentialBlockSnapshotSchema,
snapshotEntrySchema,
refresherScheduleSchema,
snapshotResponseSchema,
snapshotStreamSnapshotEventSchema,
snapshotStreamEntryEventSchema,
snapshotStreamRemovedEventSchema,
snapshotStreamEventSchema,
healthzResponseSchema,
usageResponseSchema,
usageHistoryResponseSchema,
clientUsageReportRequestSchema,
clientUsageReportResponseSchema,
clientUsageSummaryResponseSchema,
credentialRefreshResponseSchema,
credentialDisableRequestSchema,
credentialDisableResponseSchema,
disabledCredentialSummarySchema,
disabledCredentialsResponseSchema,
credentialBlockRequestSchema,
credentialBlockResponseSchema,
credentialBlocksDeleteResponseSchema,
usageStaleResponseSchema,
credentialUploadRequestSchema,
credentialUploadResponseSchema,
};
}
let cachedAuthBrokerWireSchemas: AuthBrokerWireSchemas | undefined;
export function getAuthBrokerWireSchemas(): AuthBrokerWireSchemas {
if (!cachedAuthBrokerWireSchemas) cachedAuthBrokerWireSchemas = buildAuthBrokerWireSchemas();
return cachedAuthBrokerWireSchemas;
}
+38 -371
View File
@@ -1,376 +1,43 @@
/**
* ArkType schemas for the auth-broker wire protocol.
* Eager public facade for the auth-broker wire schemas.
*
* Shared between the server (validates inbound request bodies) and the client
* (validates responses from the broker). Schemas mirror the TypeScript types
* in `./types.ts` 1:1; the types remain the source of truth for static typing,
* and `Type` is asserted-compatible with them where possible.
*
* Envelope and fixed-shape schemas use `"+": "reject"` so unknown keys are
* rejected — the previous implementation used a hand-rolled `hasOnlyFields`
* allowlist for the same effect. The OAuth credential schema is the deliberate
* exception (standard type keeps extra keys): it preserves provider-specific extension fields so
* they round-trip through the broker instead of being dropped (see below).
* Direct consumers of this module receive the shared real ArkType objects.
* Internal client/server startup imports the lazy resource instead.
*/
import { type } from "arktype";
import { REMOTE_REFRESH_SENTINEL } from "../auth-storage";
import { getAuthBrokerWireSchemas } from "./wire-schema-resource";
// ─── Credential payloads ───────────────────────────────────────────────────
const wireSchemas = getAuthBrokerWireSchemas();
/** Real OAuth credential (broker-side) — refresh token is the actual upstream value. */
export const oauthCredentialSchema = type({
"apiEndpoint?": "string",
type: "'oauth'",
refresh: type("string").narrow(
(value, ctx) =>
value !== REMOTE_REFRESH_SENTINEL ||
ctx.mustBe(`not equal to the remote sentinel (${REMOTE_REFRESH_SENTINEL})`),
),
access: type("string").atLeastLength(1),
expires: "number",
"enterpriseUrl?": "string",
"projectId?": "string",
"email?": "string",
"accountId?": "string",
"orgId?": "string",
"orgName?": "string",
"authorizedAt?": "number",
});
/** OAuth credential as it appears in broker snapshots — refresh replaced with sentinel. */
export const remoteOauthCredentialSchema = type({
"apiEndpoint?": "string",
type: "'oauth'",
refresh: type.enumerated(REMOTE_REFRESH_SENTINEL),
access: type("string").atLeastLength(1),
expires: "number",
"enterpriseUrl?": "string",
"projectId?": "string",
"email?": "string",
"accountId?": "string",
"orgId?": "string",
"orgName?": "string",
"authorizedAt?": "number",
});
export const apiKeyCredentialSchema = type({
"+": "reject",
type: "'api_key'",
key: type("string").atLeastLength(1),
"source?": "'login'",
});
/** Discriminated union accepted on POST /v1/credential (writes). */
export const writableAuthCredentialSchema = oauthCredentialSchema.or(apiKeyCredentialSchema);
/** Discriminated union returned in snapshots (refresh is sentinel for OAuth). */
export const snapshotCredentialSchema = remoteOauthCredentialSchema.or(apiKeyCredentialSchema);
// ─── Snapshot ──────────────────────────────────────────────────────────────
export const credentialSnapshotEntrySchema = type({
"+": "reject",
id: "number.integer",
provider: type("string").atLeastLength(1),
credential: snapshotCredentialSchema,
identityKey: "string | null",
});
export const credentialBlockSnapshotSchema = type({
"+": "reject",
providerKey: type("string").atLeastLength(1),
blockScope: "string",
blockedUntilMs: "number",
"updatedAtMs?": "number",
});
export const snapshotEntrySchema = type({
"+": "reject",
id: "number.integer",
provider: type("string").atLeastLength(1),
credential: snapshotCredentialSchema,
identityKey: "string | null",
rotatesInMs: "number | null",
"blocks?": credentialBlockSnapshotSchema.array(),
});
export const refresherScheduleSchema = type({
"+": "reject",
enabled: "boolean",
intervalMs: "number",
skewMs: "number",
nextSweepInMs: "number",
});
export const snapshotResponseSchema = type({
"+": "reject",
generation: "number.integer",
generatedAt: "number",
serverNowMs: "number",
refresher: refresherScheduleSchema,
credentials: snapshotEntrySchema.array(),
});
// ─── Snapshot stream (SSE) ────────────────────────────────────────────────
/** First frame on connect — full snapshot embedded inline with a `kind` tag. */
export const snapshotStreamSnapshotEventSchema = type({
"+": "reject",
generation: "number.integer",
generatedAt: "number",
serverNowMs: "number",
refresher: refresherScheduleSchema,
credentials: snapshotEntrySchema.array(),
kind: "'snapshot'",
});
/** Per-credential upsert/refresh delta. */
export const snapshotStreamEntryEventSchema = type({
"+": "reject",
kind: "'entry'",
generation: "number.integer",
serverNowMs: "number",
refresher: refresherScheduleSchema,
entry: snapshotEntrySchema,
});
/** Per-credential delete delta. */
export const snapshotStreamRemovedEventSchema = type({
"+": "reject",
kind: "'removed'",
generation: "number.integer",
serverNowMs: "number",
refresher: refresherScheduleSchema,
id: "number.integer",
});
/** Discriminated union over every event frame the snapshot stream emits. */
export const snapshotStreamEventSchema = snapshotStreamSnapshotEventSchema
.or(snapshotStreamEntryEventSchema)
.or(snapshotStreamRemovedEventSchema);
// ─── Healthz ────────────────────────────────────────────────────────────────
export const healthzResponseSchema = type({
"+": "reject",
ok: "boolean",
"version?": "string",
});
// ─── Usage ─────────────────────────────────────────────────────────────────
const usageUnitSchema = type("'percent' | 'tokens' | 'requests' | 'usd' | 'minutes' | 'bytes' | 'unknown'");
const usageStatusSchema = type("'ok' | 'warning' | 'exhausted' | 'unknown'");
const usageWindowSchema = type({
id: "string",
label: "string",
"durationMs?": "number",
"resetsAt?": "number",
});
const usageAmountSchema = type({
"used?": "number",
"limit?": "number",
"remaining?": "number",
"usedFraction?": "number",
"remainingFraction?": "number",
unit: usageUnitSchema,
});
const usageScopeSchema = type({
provider: "string",
"accountId?": "string",
"projectId?": "string",
"orgId?": "string",
"modelId?": "string",
"tier?": "string",
"windowId?": "string",
"shared?": "boolean",
});
const usageLimitSchema = type({
id: "string",
label: "string",
scope: usageScopeSchema,
"window?": usageWindowSchema,
amount: usageAmountSchema,
"status?": usageStatusSchema,
"notes?": "string[]",
});
const usageResetCreditDetailSchema = type({
"grantedAt?": "string",
"expiresAt?": "string",
"status?": "string",
});
const usageResetCreditsSchema = type({
availableCount: "number",
"credits?": usageResetCreditDetailSchema.array(),
});
const arkUsageReportSchema = type({
provider: "string",
fetchedAt: "number",
limits: usageLimitSchema.array(),
"resetCredits?": usageResetCreditsSchema,
"notes?": "string[]",
"metadata?": { "[string]": "unknown" },
"raw?": "unknown",
});
/**
* Broker `/v1/usage` response. Reports are full {@link UsageReport}s minus the
* heavy provider-specific `raw` field (the server strips it before send) — we
* keep `raw` optional in the underlying schema so a misconfigured broker that
* forgot to strip still validates.
*/
export const usageResponseSchema = type({
"+": "reject",
generatedAt: "number",
reports: arkUsageReportSchema.array(),
});
const usageHistoryEntrySchema = type({
recordedAt: "number",
provider: "string",
accountKey: "string",
"email?": "string",
"accountId?": "string",
limitId: "string",
label: "string",
"windowLabel?": "string",
"usedFraction?": "number",
"status?": "'ok' | 'warning' | 'exhausted' | 'unknown'",
"resetsAt?": "number",
});
/** Broker `/v1/usage/history` response — recorded usage-limit snapshots, oldest first. */
export const usageHistoryResponseSchema = type({
"+": "reject",
generatedAt: "number",
entries: usageHistoryEntrySchema.array(),
});
const observedUsageEntrySchema = type({
at: "number",
provider: "string",
model: "string",
requests: "number",
inputTokens: "number",
outputTokens: "number",
cacheReadTokens: "number",
cacheWriteTokens: "number",
costUsd: "number",
});
/** Broker `POST /v1/usage/observed` request — one client's batched observed usage. */
export const clientUsageReportRequestSchema = type({
"+": "reject",
installId: "string",
"hostname?": "string",
entries: observedUsageEntrySchema.array(),
});
export const clientUsageReportResponseSchema = type({
"+": "reject",
ok: "boolean",
});
const clientProviderUsageSchema = type({
provider: "string",
requests: "number",
inputTokens: "number",
outputTokens: "number",
cacheReadTokens: "number",
cacheWriteTokens: "number",
costUsd: "number",
});
const clientUsageClientSummarySchema = type({
installId: "string",
"hostname?": "string",
firstSeen: "number",
lastSeen: "number",
providers: clientProviderUsageSchema.array(),
});
/** Broker `GET /v1/usage/clients` response — per-client token burn aggregates. */
export const clientUsageSummaryResponseSchema = type({
"+": "reject",
generatedAt: "number",
clients: clientUsageClientSummarySchema.array(),
});
// ─── Refresh ───────────────────────────────────────────────────────────────
export const credentialRefreshResponseSchema = type({
"+": "reject",
entry: credentialSnapshotEntrySchema,
});
// ─── Disable ───────────────────────────────────────────────────────────────
export const credentialDisableRequestSchema = type({
"+": "reject",
"cause?": "string",
});
export const credentialDisableResponseSchema = type({
"+": "reject",
ok: "boolean",
});
/** One disabled-credential tombstone — identity + cause, never token material. */
export const disabledCredentialSummarySchema = type({
"+": "reject",
id: "number.integer",
provider: type("string").atLeastLength(1),
type: "'oauth' | 'api_key'",
"email?": "string",
"accountId?": "string",
"orgId?": "string",
"orgName?": "string",
cause: "string",
"disabledAtMs?": "number",
});
/** Broker `GET /v1/credentials/disabled` response. */
export const disabledCredentialsResponseSchema = type({
"+": "reject",
generatedAt: "number",
disabled: disabledCredentialSummarySchema.array(),
});
// ─── Credential blocks ──────────────────────────────────────────────────────
export const credentialBlockRequestSchema = credentialBlockSnapshotSchema;
export const credentialBlockResponseSchema = type({
"+": "reject",
ok: "boolean",
});
export const credentialBlocksDeleteResponseSchema = type({
"+": "reject",
ok: "boolean",
});
export const usageStaleResponseSchema = type({
"+": "reject",
ok: "boolean",
});
// ─── Upload ────────────────────────────────────────────────────────────────
export const credentialUploadRequestSchema = type({
"+": "reject",
provider: type("string").atLeastLength(1),
credential: writableAuthCredentialSchema,
});
export const credentialUploadResponseSchema = type({
"+": "reject",
entries: credentialSnapshotEntrySchema.array(),
});
export const {
oauthCredentialSchema,
remoteOauthCredentialSchema,
apiKeyCredentialSchema,
writableAuthCredentialSchema,
snapshotCredentialSchema,
credentialSnapshotEntrySchema,
credentialBlockSnapshotSchema,
snapshotEntrySchema,
refresherScheduleSchema,
snapshotResponseSchema,
snapshotStreamSnapshotEventSchema,
snapshotStreamEntryEventSchema,
snapshotStreamRemovedEventSchema,
snapshotStreamEventSchema,
healthzResponseSchema,
usageResponseSchema,
usageHistoryResponseSchema,
clientUsageReportRequestSchema,
clientUsageReportResponseSchema,
clientUsageSummaryResponseSchema,
credentialRefreshResponseSchema,
credentialDisableRequestSchema,
credentialDisableResponseSchema,
disabledCredentialSummarySchema,
disabledCredentialsResponseSchema,
credentialBlockRequestSchema,
credentialBlockResponseSchema,
credentialBlocksDeleteResponseSchema,
usageStaleResponseSchema,
credentialUploadRequestSchema,
credentialUploadResponseSchema,
} = wireSchemas;
@@ -0,0 +1,37 @@
import { expect, test } from "bun:test";
import * as path from "node:path";
import { TempDir } from "@oh-my-pi/pi-utils";
const preloadPath = path.join(import.meta.dir, "fixtures", "auth-broker-wire-construction-preload.ts");
const probePath = path.join(import.meta.dir, "fixtures", "auth-broker-wire-construction-probe.ts");
test("auth-broker wire schemas construct only on first validation", async () => {
const tempDir = TempDir.createSync("@auth-broker-wire-");
try {
const proc = Bun.spawn([process.execPath, "--preload", preloadPath, probePath, tempDir.path()], {
cwd: path.join(import.meta.dir, "../../.."),
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([
new Response(proc.stdout).text(),
new Response(proc.stderr).text(),
proc.exited,
]);
expect(exitCode, stderr).toBe(0);
expect(JSON.parse(stdout)).toEqual({
counts: {
afterModuleImport: 0,
afterLocalDiscovery: 0,
afterConstruction: 0,
afterFirstHealth: 1,
afterSecondHealth: 1,
},
firstHealth: { ok: true, version: "wire-lazy-probe" },
secondHealth: { ok: true, version: "wire-lazy-probe" },
});
} finally {
await tempDir.remove().catch(() => {});
}
});
@@ -0,0 +1,267 @@
import { describe, expect, test } from "bun:test";
import * as wireSchemas from "@oh-my-pi/pi-ai/auth-broker/wire-schemas";
import { Type, type } from "arktype";
const REFRESHER = {
enabled: false,
intervalMs: 60_000,
skewMs: 300_000,
nextSweepInMs: Number.MAX_SAFE_INTEGER,
};
const BLOCK = {
providerKey: "anthropic:oauth",
blockScope: "tier:fable",
blockedUntilMs: 4_000,
updatedAtMs: 3_000,
};
const REMOTE_OAUTH = {
type: "oauth",
access: "access",
refresh: "__remote__",
expires: 5_000,
tokenUrl: "https://example.test/token",
clientId: "provider-client",
};
const REAL_OAUTH = { ...REMOTE_OAUTH, refresh: "real-refresh" };
const API_KEY = { type: "api_key", key: "secret", source: "login" };
const CREDENTIAL_ENTRY = {
id: 7,
provider: "anthropic",
credential: REMOTE_OAUTH,
identityKey: "account:test",
};
const SNAPSHOT_ENTRY = {
...CREDENTIAL_ENTRY,
rotatesInMs: null,
blocks: [BLOCK],
};
const SNAPSHOT = {
generation: 2,
generatedAt: 1_000,
serverNowMs: 2_000,
refresher: REFRESHER,
credentials: [SNAPSHOT_ENTRY],
};
const STREAM_SNAPSHOT = { kind: "snapshot", ...SNAPSHOT };
const STREAM_ENTRY = {
kind: "entry",
generation: 3,
serverNowMs: 2_500,
refresher: REFRESHER,
entry: SNAPSHOT_ENTRY,
};
const STREAM_REMOVED = {
kind: "removed",
generation: 4,
serverNowMs: 3_000,
refresher: REFRESHER,
id: 7,
};
const USAGE_REPORT = {
provider: "anthropic",
fetchedAt: 2_000,
limits: [
{
id: "rolling",
label: "Rolling window",
scope: { provider: "anthropic", windowId: "rolling", providerExtension: true },
window: { id: "rolling", label: "5 hour", durationMs: 18_000_000 },
amount: { used: 1, limit: 10, remaining: 9, unit: "tokens", providerExtension: "kept" },
status: "ok",
notes: ["limit note"],
providerExtension: 1,
},
],
notes: ["report note"],
metadata: { plan: "max" },
raw: { providerPayload: true },
providerExtension: "kept",
};
const OBSERVED_USAGE = {
at: 1_000,
provider: "anthropic",
model: "claude",
requests: 1,
inputTokens: 2,
outputTokens: 3,
cacheReadTokens: 4,
cacheWriteTokens: 5,
costUsd: 0.01,
};
const schemaNames = [
"oauthCredentialSchema",
"remoteOauthCredentialSchema",
"apiKeyCredentialSchema",
"writableAuthCredentialSchema",
"snapshotCredentialSchema",
"credentialSnapshotEntrySchema",
"credentialBlockSnapshotSchema",
"snapshotEntrySchema",
"refresherScheduleSchema",
"snapshotResponseSchema",
"snapshotStreamSnapshotEventSchema",
"snapshotStreamEntryEventSchema",
"snapshotStreamRemovedEventSchema",
"snapshotStreamEventSchema",
"healthzResponseSchema",
"usageResponseSchema",
"usageHistoryResponseSchema",
"clientUsageReportRequestSchema",
"clientUsageReportResponseSchema",
"clientUsageSummaryResponseSchema",
"credentialRefreshResponseSchema",
"credentialDisableRequestSchema",
"credentialDisableResponseSchema",
"disabledCredentialSummarySchema",
"disabledCredentialsResponseSchema",
"credentialBlockRequestSchema",
"credentialBlockResponseSchema",
"credentialBlocksDeleteResponseSchema",
"usageStaleResponseSchema",
"credentialUploadRequestSchema",
"credentialUploadResponseSchema",
] as const;
type SchemaName = (typeof schemaNames)[number];
type CallableSchema = (input: unknown) => unknown;
const validSamples: Record<SchemaName, unknown> = {
oauthCredentialSchema: REAL_OAUTH,
remoteOauthCredentialSchema: REMOTE_OAUTH,
apiKeyCredentialSchema: API_KEY,
writableAuthCredentialSchema: REAL_OAUTH,
snapshotCredentialSchema: REMOTE_OAUTH,
credentialSnapshotEntrySchema: CREDENTIAL_ENTRY,
credentialBlockSnapshotSchema: BLOCK,
snapshotEntrySchema: SNAPSHOT_ENTRY,
refresherScheduleSchema: REFRESHER,
snapshotResponseSchema: SNAPSHOT,
snapshotStreamSnapshotEventSchema: STREAM_SNAPSHOT,
snapshotStreamEntryEventSchema: STREAM_ENTRY,
snapshotStreamRemovedEventSchema: STREAM_REMOVED,
snapshotStreamEventSchema: STREAM_ENTRY,
healthzResponseSchema: { ok: true, version: "contract" },
usageResponseSchema: { generatedAt: 2_000, reports: [USAGE_REPORT] },
usageHistoryResponseSchema: {
generatedAt: 2_000,
entries: [
{
recordedAt: 1_000,
provider: "anthropic",
accountKey: "account:test",
limitId: "rolling",
label: "Rolling window",
usedFraction: 0.1,
status: "ok",
},
],
},
clientUsageReportRequestSchema: { installId: "install", hostname: "host", entries: [OBSERVED_USAGE] },
clientUsageReportResponseSchema: { ok: true },
clientUsageSummaryResponseSchema: {
generatedAt: 2_000,
clients: [
{
installId: "install",
hostname: "host",
firstSeen: 1_000,
lastSeen: 2_000,
providers: [{ ...OBSERVED_USAGE, firstSeen: undefined, at: undefined, model: undefined }],
},
],
},
credentialRefreshResponseSchema: { entry: CREDENTIAL_ENTRY },
credentialDisableRequestSchema: {},
credentialDisableResponseSchema: { ok: true },
disabledCredentialSummarySchema: {
id: 7,
provider: "anthropic",
type: "oauth",
email: "user@example.test",
cause: "revoked",
disabledAtMs: 2_000,
},
disabledCredentialsResponseSchema: {
generatedAt: 2_000,
disabled: [{ id: 7, provider: "anthropic", type: "oauth", cause: "revoked" }],
},
credentialBlockRequestSchema: BLOCK,
credentialBlockResponseSchema: { ok: true },
credentialBlocksDeleteResponseSchema: { ok: true },
usageStaleResponseSchema: { ok: true },
credentialUploadRequestSchema: { provider: "anthropic", credential: REAL_OAUTH },
credentialUploadResponseSchema: { entries: [CREDENTIAL_ENTRY] },
};
function run(schema: unknown, input: unknown): unknown {
return (schema as CallableSchema)(input);
}
function accept(schema: unknown, input: unknown): unknown {
const result = run(schema, input);
expect(result).not.toBeInstanceOf(type.errors);
if (result instanceof type.errors) throw new Error(`Expected schema acceptance: ${result.summary}`);
return result;
}
function reject(schema: unknown, input: unknown): void {
expect(run(schema, input)).toBeInstanceOf(type.errors);
}
describe("auth-broker public wire schemas", () => {
test("exports all 31 real callable ArkType values with canonical behavior", () => {
expect(Object.keys(wireSchemas).sort()).toEqual([...schemaNames].sort());
for (const name of schemaNames) {
// biome-ignore lint/performance/noDynamicNamespaceImportAccess: this contract intentionally verifies the public namespace.
const schema = wireSchemas[name];
expect(typeof schema).toBe("function");
expect(schema).toBeInstanceOf(Type);
accept(schema, validSamples[name]);
}
});
test("preserves credential extension and sentinel boundaries", () => {
expect(accept(wireSchemas.oauthCredentialSchema, REAL_OAUTH)).toEqual(REAL_OAUTH);
expect(accept(wireSchemas.remoteOauthCredentialSchema, REMOTE_OAUTH)).toEqual(REMOTE_OAUTH);
reject(wireSchemas.oauthCredentialSchema, REMOTE_OAUTH);
reject(wireSchemas.remoteOauthCredentialSchema, REAL_OAUTH);
reject(wireSchemas.oauthCredentialSchema, { ...REAL_OAUTH, access: "" });
reject(wireSchemas.apiKeyCredentialSchema, { ...API_KEY, extra: true });
reject(wireSchemas.apiKeyCredentialSchema, { ...API_KEY, source: "environment" });
reject(wireSchemas.credentialUploadRequestSchema, { provider: "", credential: REAL_OAUTH });
});
test("preserves fixed envelopes, integer fields, discriminators, and block alias identity", () => {
expect(wireSchemas.credentialBlockRequestSchema).toBe(wireSchemas.credentialBlockSnapshotSchema);
accept(wireSchemas.credentialBlockRequestSchema, {
providerKey: BLOCK.providerKey,
blockScope: "",
blockedUntilMs: BLOCK.blockedUntilMs,
});
accept(wireSchemas.credentialDisableRequestSchema, {});
reject(wireSchemas.credentialDisableRequestSchema, { cause: 1 });
reject(wireSchemas.credentialDisableRequestSchema, { extra: true });
reject(wireSchemas.healthzResponseSchema, { ok: true, extra: true });
reject(wireSchemas.snapshotResponseSchema, { ...SNAPSHOT, generation: 1.5 });
reject(wireSchemas.snapshotResponseSchema, { ...SNAPSHOT, extra: true });
reject(wireSchemas.snapshotEntrySchema, { ...SNAPSHOT_ENTRY, id: 1.5 });
reject(wireSchemas.snapshotStreamEventSchema, { ...STREAM_ENTRY, kind: "snapshot" });
reject(wireSchemas.snapshotStreamEventSchema, { ...STREAM_REMOVED, extra: true });
reject(wireSchemas.snapshotStreamEventSchema, { kind: "unknown" });
});
test("preserves usage extensions while rejecting envelope and enum violations", () => {
const response = { generatedAt: 2_000, reports: [USAGE_REPORT] };
expect(accept(wireSchemas.usageResponseSchema, response)).toEqual(response);
reject(wireSchemas.usageResponseSchema, { ...response, extra: true });
reject(wireSchemas.usageResponseSchema, {
...response,
reports: [{ ...USAGE_REPORT, limits: [{ ...USAGE_REPORT.limits[0], amount: { unit: "seconds" } }] }],
});
reject(wireSchemas.usageResponseSchema, {
...response,
reports: [{ ...USAGE_REPORT, limits: [{ ...USAGE_REPORT.limits[0], status: "critical" }] }],
});
});
});
@@ -0,0 +1,15 @@
import { spyOn } from "bun:test";
import { type } from "arktype";
declare global {
var __authBrokerWireConstructionCount: number;
}
globalThis.__authBrokerWireConstructionCount = 0;
const originalEnumerated = type.enumerated;
spyOn(type, "enumerated").mockImplementation((...values) => {
if (values.length === 1 && values[0] === "__remote__") {
globalThis.__authBrokerWireConstructionCount += 1;
}
return originalEnumerated(...values);
});
@@ -0,0 +1,55 @@
import {
AuthBrokerClient,
type AuthBrokerServerHandle,
discoverAuthStorage,
startAuthBroker,
} from "@oh-my-pi/pi-ai/auth-broker";
import type { AuthStorage } from "@oh-my-pi/pi-ai/auth-storage";
declare global {
var __authBrokerWireConstructionCount: number;
}
const agentDir = process.argv[2];
if (!agentDir) throw new Error("Expected an isolated agent directory");
const count = (): number => globalThis.__authBrokerWireConstructionCount;
const counts = {
afterModuleImport: count(),
afterLocalDiscovery: -1,
afterConstruction: -1,
afterFirstHealth: -1,
afterSecondHealth: -1,
};
let storage: AuthStorage | undefined;
let handle: AuthBrokerServerHandle | undefined;
try {
delete process.env.OMP_AUTH_BROKER_URL;
delete process.env.OMP_AUTH_BROKER_TOKEN;
delete process.env.OMP_AUTH_BROKER_ACCOUNT_POOL_FILE;
storage = await discoverAuthStorage({ agentDir });
counts.afterLocalDiscovery = count();
handle = startAuthBroker({
storage,
bind: "127.0.0.1:0",
bearerTokens: [],
version: "wire-lazy-probe",
disableRefresher: true,
});
const client = new AuthBrokerClient({ url: handle.url, token: "unused", maxRetries: 0 });
counts.afterConstruction = count();
const firstHealth = await client.healthz();
counts.afterFirstHealth = count();
const secondHealth = await client.healthz();
counts.afterSecondHealth = count();
process.stdout.write(JSON.stringify({ counts, firstHealth, secondHealth }));
} finally {
await handle?.close();
storage?.close();
}