diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index b65995898..07a378850 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -179,6 +179,7 @@ jobs: run: bun run test:e2e env: E2E_REQUIRED: "1" + DATABASE_URL: postgres://postgres:postgres@localhost:5432/workbench - name: Run the two-org isolation suite run: bun test test/isolation @@ -188,9 +189,14 @@ jobs: # Postgres. E2E_REQUIRED=1 turns a would-be skip here into a hard # failure, so a misconfigured invocation can't pass vacuously. - name: Run the hub's database-backed suites - run: bun test apps/hub/test + run: | + set -a + . ./.env + set +a + bun test apps/hub/test env: E2E_REQUIRED: "1" + DATABASE_URL: postgres://postgres:postgres@localhost:5432/workbench # Every DB-gated suite under packages/* and apps/hub/src gates on # DATABASE_URL and describe.skip's without one, so build-test (no @@ -200,6 +206,10 @@ jobs: # would-be skip into a hard failure. - name: Run the database-backed package suites run: | - bun test $(grep -rl DATABASE_URL --include='*.test.ts' packages apps | grep -v apps/hub/test) + set -a + . ./.env + set +a + bun --env-file="${GITHUB_WORKSPACE}/.env" test $(grep -rl DATABASE_URL --include='*.test.ts' packages apps | grep -v apps/hub/test) env: E2E_REQUIRED: "1" + DATABASE_URL: postgres://postgres:postgres@localhost:5432/workbench diff --git a/apps/hub/src/index.ts b/apps/hub/src/index.ts index 049f618ec..ad714a3b7 100644 --- a/apps/hub/src/index.ts +++ b/apps/hub/src/index.ts @@ -97,12 +97,14 @@ import { localPartOf, parseParticipants, postRoomMessage, + recordSourcesDigest, startWorkflowCommand, sendWorkbenchMessage, settleConnectedService, workbenchLaunchPersistExtra, } from "@corbits/chat"; import type { RelaunchNoticePort } from "@corbits/chat"; +import { reportError } from "@corbits/error-sink"; import type { FinalizedTurnToolCall } from "@corbits/turn-artifacts"; import { decodedOrNull } from "@corbits/url-path"; import { @@ -286,6 +288,7 @@ import { } from "@workbench/onboarding"; import { createConnectionRoutes, + isInferenceProvider, createMcpOAuthRoutes, createMcpServerRoutes, createOAuthConnectRoutes, @@ -1987,6 +1990,8 @@ export async function createHub(config: HubConfig) { cryptoProviderCache: foldedRunCryptoProviders, launchMode: AGENT_SECTION_MODE, persistLaunch: workbenchLaunchPersistExtra, + recordLaunchSources: ({ instanceId, sourcesDigest }) => + recordSourcesDigest(db, instanceId, sourcesDigest), }, trigger, payload, @@ -1999,8 +2004,15 @@ export async function createHub(config: HubConfig) { // clears (flipping the in-room connect card via `chat.settings`), and // the host agent is woken via `dispatchTurn` without a forged // signed-in-user timeline row. - const settleServiceConnection: ServiceConnectedHook = (info) => - settleConnectedService( + // + // An inference provider's credential landing also re-checks every + // live participant's deployed inference chain (CL-6687): a rotated + // key only ever reaches an agent at deploy time, so the relaunch has + // to be kicked here, not left for the next message. Not awaited — a + // relaunch is a sidecar deploy round-trip, and the connect response + // must not wait on it. + const settleServiceConnection: ServiceConnectedHook = async (info) => { + await settleConnectedService( { store: chatStore, platform: chatPlatform, @@ -2015,6 +2027,19 @@ export async function createHub(config: HubConfig) { displayName: info.displayName, }, ); + if (!isInferenceProvider(info.connectorId)) return; + void chatPlatform + .reconcileInferenceSources(info.tenantId) + .then(({ scanned, relaunched }) => { + log.info`inference credential ${info.connectorId} changed on tenant ${info.tenantId}: re-checked ${String(scanned)} live agents, relaunched ${String(relaunched)}`; + }) + .catch((cause: unknown) => { + reportError(cause, { + operation: "connections.reconcile-inference-sources", + tenantId: info.tenantId, + }); + }); + }; // Connections: the settings surface's tenant-scoped credential // test-and-store, mounted under the same tenant prefix and reusing // the same grant store/condition registry every other credential- diff --git a/apps/hub/src/routine-launcher.test.ts b/apps/hub/src/routine-launcher.test.ts index a657afcb2..13572235e 100644 --- a/apps/hub/src/routine-launcher.test.ts +++ b/apps/hub/src/routine-launcher.test.ts @@ -25,6 +25,7 @@ const FOLDED_BODY = { const FRAME_HEADER = "Input from this routine's setup:"; let launchFoldedRunCalls: unknown[] = []; +const launchRowUpdates: unknown[] = []; let sendFoldedMailWithRetryCalls: unknown[] = []; let sendFoldedMailWithRetryResult: unknown = { ok: true, @@ -50,7 +51,11 @@ mock.module("@corbits/folded-runs", () => ({ }, launchFoldedRun: async (...args: unknown[]) => { launchFoldedRunCalls.push(args); - return { instancePrincipalId: "prn_run1", sessionId: "ses_run1" }; + return { + instancePrincipalId: "prn_run1", + sessionId: "ses_run1", + sourcesDigest: "digest_run1", + }; }, sendFoldedMailWithRetry: async (...args: unknown[]) => { sendFoldedMailWithRetryCalls.push(args); @@ -118,6 +123,15 @@ function createFakeDb( }), }), }), + // `recordSourcesDigest` writes the deployed inference chain's digest + // onto the launch row once `launchFoldedRun` returns (CL-6687). + update: () => ({ + set: (values: unknown) => ({ + where: async () => { + launchRowUpdates.push(values); + }, + }), + }), }; } diff --git a/apps/hub/src/routine-launcher.ts b/apps/hub/src/routine-launcher.ts index c38cb0662..891d39925 100644 --- a/apps/hub/src/routine-launcher.ts +++ b/apps/hub/src/routine-launcher.ts @@ -50,6 +50,7 @@ import type { AssetService } from "@intx/hub-sessions"; import { AGENT_SECTION_MODE, handleFromName, + recordSourcesDigest, workbenchLaunchPersistExtra, } from "@corbits/chat"; import { renderRoutineInput, type RoutineLauncher } from "@corbits/routines"; @@ -199,6 +200,7 @@ export function createHubRoutineLauncher( foldedBody, }), }); + await recordSourcesDigest(deps.db, instanceId, launched.sourcesDigest); if ( input.deliveryWorkbenchId !== undefined && diff --git a/docs/credential-wiring.md b/docs/credential-wiring.md index 32faf5db1..ab801e887 100644 --- a/docs/credential-wiring.md +++ b/docs/credential-wiring.md @@ -70,3 +70,31 @@ tool package. See the comment on `packageFromToolId` (`apps/sidecar/src/step-agent-tools.ts`) for the same note in code; the loader-side provenance check this implies for a future untrusted-package story is tracked as its own follow-up, not solved here. + +## Rotation reaches live agents (CL-6687) + +Inference sources -- the decrypted provider key included -- are resolved +at deploy time and rendered into the run's bytes; nothing re-reads them +while the run is resident. So a rotated key (Settings > AI providers, +disconnect + reconnect after a 401) used to reach only the _next_ deploy, +and an already-open workbench kept sending the dead key until the stack +restarted. + +Every deploy now records `inferenceSourcesDigest` (a SHA-256 over the +resolved chain, secret included; the secret itself is never stored) on +the run's `workbench_launch` row. Two paths compare it against today's +resolution and relaunch the run on a mismatch, the same relaunch a +definition-content drift triggers (CL-6588): + +- `reconcileDriftedRun`, ahead of every send, so the next message in an + open room uses the new key even if nothing else fired. +- `reconcileInferenceSources(tenantId)`, kicked by the hub the moment an + inference provider's `/complete` stores a credential, so the fix an + operator just applied lands without waiting for a message. + +Re-saving the same key produces the same digest and relaunches nothing. +A run whose row predates the column (`sources_digest IS NULL`) gets +today's chain recorded as its baseline on its first check and is left +alone that once. The per-send check is throttled to one resolution per +participant per 30 seconds; the provider-connect hook is the path that +reaches live rooms immediately. diff --git a/packages/chat/src/agent-binding.ts b/packages/chat/src/agent-binding.ts index d29fdf0de..c1f49aa29 100644 --- a/packages/chat/src/agent-binding.ts +++ b/packages/chat/src/agent-binding.ts @@ -43,6 +43,8 @@ export interface AgentBinding { /** Every run this participant used to be, oldest first. */ readonly priorRunIds: readonly string[]; readonly foldedBody: FoldedBody; + /** See `workbenchLaunch.sourcesDigest`. */ + readonly sourcesDigest: string | null; } /** @@ -97,6 +99,7 @@ function bindingFrom(row: LaunchRow, domain: string): AgentBinding { liveAddress: formatRunAddress(row.currentRunId, domain), priorRunIds: priorRunIdsFrom(row), foldedBody: parsed, + sourcesDigest: row.sourcesDigest ?? null, }; } @@ -288,16 +291,60 @@ export async function repointBinding( db: DB["db"], binding: AgentBinding, newRunId: string, + sourcesDigest: string, ): Promise { const history = [...binding.priorRunIds, binding.currentRunId].slice( -PRIOR_RUN_HISTORY_LIMIT, ); await db .update(workbenchLaunch) - .set({ currentRunId: newRunId, priorRunIds: history }) + .set({ currentRunId: newRunId, priorRunIds: history, sourcesDigest }) .where(eq(workbenchLaunch.instanceId, binding.stableId)); } +/** + * Records the inference chain a deploy just pinned for the participant + * `stableId` names — a wake, or a standalone launch whose mapping row + * `workbenchLaunchPersistExtra` wrote before the deploy resolved — so + * the next send can tell whether the tenant's catalog (a rotated key, a + * moved endpoint) has moved on from it since. + */ +export async function recordSourcesDigest( + db: DB["db"], + stableId: string, + sourcesDigest: string, +): Promise { + await db + .update(workbenchLaunch) + .set({ sourcesDigest }) + .where(eq(workbenchLaunch.instanceId, stableId)); +} + +/** + * Every participant a tenant has launched, with its live run, for the + * pass that re-checks each one's inference chain after a provider + * credential changes (`platform-adapter.ts`'s + * `reconcileInferenceSources`). Bounded like `listLaunchesBeyondWake`. + */ +export async function listLaunchesForTenant( + db: DB["db"], + tenantId: string, + limit: number, +): Promise { + const rows = await db + .select() + .from(workbenchLaunch) + .where(eq(workbenchLaunch.tenantId, tenantId)) + .limit(limit); + const live: LiveAgent[] = []; + for (const row of rows) { + const run = await readRun(db, row.currentRunId); + if (run === undefined || run.address === null) continue; + live.push({ binding: bindingFrom(row, requireDomain(run.address)), run }); + } + return live; +} + /** * Every participant whose current run is beyond waking, for the boot * sweep that relaunches them (`platform-adapter.ts`'s diff --git a/packages/chat/src/index.ts b/packages/chat/src/index.ts index 0253ba4f2..d97b7c446 100644 --- a/packages/chat/src/index.ts +++ b/packages/chat/src/index.ts @@ -133,6 +133,7 @@ export { AGENT_SECTION_MODE, workbenchLaunchPersistExtra, } from "./standalone-launch"; +export { recordSourcesDigest } from "./agent-binding"; export { createWorkbenchTurnQueue, TurnQueuedEvent } from "./turn-queue"; export type { DispatchTurnBatch, diff --git a/packages/chat/src/migrations.ts b/packages/chat/src/migrations.ts index 9e92205ee..ace9b7a3f 100644 --- a/packages/chat/src/migrations.ts +++ b/packages/chat/src/migrations.ts @@ -407,6 +407,13 @@ export const chatMigrations: readonly ChatMigration[] = [ DROP COLUMN IF EXISTS "noop_inference"; `, }, + { + name: "0024_workbench_launch_sources_digest", + sql: ` + ALTER TABLE "chat"."workbench_launch" + ADD COLUMN IF NOT EXISTS "sources_digest" text; + `, + }, ]; /** diff --git a/packages/chat/src/platform-adapter.ts b/packages/chat/src/platform-adapter.ts index ba8ab665f..7338e6198 100644 --- a/packages/chat/src/platform-adapter.ts +++ b/packages/chat/src/platform-adapter.ts @@ -17,10 +17,13 @@ import { createCryptoProviderCache, DefinitionProjectionMissingError, domainOf, + inferenceSourcesDigest, + InferenceResolutionError, launchFoldedRun, mintFoldedRun, readFoldedBody, resolveFoldedRunSessionId, + resolveLaunchSources, resolveNewestProjectedDefinition, sendFoldedMail, wakeFoldedRun, @@ -32,8 +35,10 @@ import { findStandingLaunchByDefinition, isBeyondWake, listLaunchesBeyondWake, + listLaunchesForTenant, readBindingByAddress, readPriorRuns, + recordSourcesDigest, repointBinding, resolveLiveAgent, resolveLiveByStableId, @@ -184,6 +189,19 @@ export type HubChatPlatform = ChatPlatform & { * that run's terminal event. */ sweepTerminalRuns(): Promise<{ scanned: number; relaunched: number }>; + /** + * Re-checks every live participant in `tenantId` against the tenant's + * current inference catalog and relaunches the ones whose deployed + * chain no longer matches it (CL-6687) — a rotated API key, a moved + * Ollama endpoint, a changed default. The host runs this the moment a + * provider credential is stored, so the fix an operator just applied + * in Settings reaches an already-open workbench without waiting for + * its next message. Best-effort per participant: one failed relaunch + * is logged and the pass moves on. + */ + reconcileInferenceSources( + tenantId: string, + ): Promise<{ scanned: number; relaunched: number }>; }; /** @@ -342,7 +360,7 @@ export function createHubChatPlatform( ); wakeLogger.info`relaunching ${binding.roomAddress}: run ${run.id} is terminal (${run.status}); minting fresh run ${newRunId}`; - await launchFoldedRun(foldedRunsDeps, { + const deployed = await launchFoldedRun(foldedRunsDeps, { tenantId: binding.tenantId, instanceId: newRunId, triggerAddress: newAddress, @@ -356,7 +374,7 @@ export function createHubChatPlatform( // outlived a failed launch would leave the room addressing a run // `launchFoldedRun` had already rolled back, and the next message // would resolve nothing at all. - await repointBinding(deps.db, binding, newRunId); + await repointBinding(deps.db, binding, newRunId, deployed.sourcesDigest); lifecycle?.untrack(binding.liveAddress); lifecycle?.track(newAddress); @@ -504,6 +522,60 @@ export function createHubChatPlatform( return freshFoldedBody; } + /** + * How long one participant's inference chain is trusted after a + * check. Resolving the chain decrypts every candidate credential, and + * `reconcileDriftedRun` runs ahead of every send — a provider connect + * reaches live rooms through `reconcileInferenceSources` anyway, so + * the per-send check is the backstop and need not run on every + * message. + */ + const SOURCES_CHECK_INTERVAL_MS = 30_000; + const sourcesCheckedAt = new Map(); + + /** + * Whether the inference chain this run deployed with no longer + * matches what the tenant catalog resolves today (CL-6687). The + * deployed bytes carry the decrypted secret, so a rotated key is + * invisible to `foldedBody` comparison — only the chain's own digest + * can see it. A run whose digest was never recorded (a row from before + * the column existed) gets today's chain recorded as its baseline and + * is left alone this time; a chain that does not resolve at all today + * is equally left alone — its next wake fails loud on the same + * `InferenceResolutionError` either way, and relaunching it now would + * just fail sooner. + */ + async function hasDriftedSources(binding: AgentBinding): Promise { + const checkedAt = sourcesCheckedAt.get(binding.stableId); + const now = Date.now(); + if ( + checkedAt !== undefined && + now - checkedAt < SOURCES_CHECK_INTERVAL_MS + ) { + return false; + } + const { fallbackModel } = await deployShapeFor(binding); + let resolved: Awaited>; + try { + resolved = await resolveLaunchSources(foldedRunsDeps, { + tenantId: binding.tenantId, + foldedBody: binding.foldedBody, + launchLabel: "the inference-source drift check", + ...(fallbackModel !== undefined ? { fallbackModel } : {}), + }); + } catch (cause: unknown) { + if (cause instanceof InferenceResolutionError) return false; + throw cause; + } + sourcesCheckedAt.set(binding.stableId, now); + const digest = inferenceSourcesDigest(resolved); + if (binding.sourcesDigest === null) { + await recordSourcesDigest(deps.db, binding.stableId, digest); + return false; + } + return digest !== binding.sourcesDigest; + } + /** * Brings the run behind `address` back to routable — or, now, current * — whichever kind of "not serving what it should" it is in. @@ -567,10 +639,15 @@ export function createHubChatPlatform( principalId: live.run.principalId, foldedBody: binding.foldedBody, }; - await wakeFoldedRun(foldedRunsDeps, { + const deployed = await wakeFoldedRun(foldedRunsDeps, { ...wakeParams, ...(await deployShapeFor(binding)), }); + await recordSourcesDigest( + deps.db, + binding.stableId, + deployed.sourcesDigest, + ); } /** @@ -602,13 +679,13 @@ export function createHubChatPlatform( * that never needed waking must not gain a new failure mode from a * check that used to never run for it. */ - async function reconcileDriftedRun(address: string): Promise { + async function reconcileDriftedRun(address: string): Promise { const binding = await readBindingByAddress(deps.db, address); - if (binding === undefined) return; + if (binding === undefined) return false; const live = await resolveLiveAgent(deps.db, binding); - if (live === undefined || live.run.address === null) return; - if (await isBeyondWake(deps.db, live.run)) return; - if (!isRoutable(live.run.address)) return; + if (live === undefined || live.run.address === null) return false; + if (await isBeyondWake(deps.db, live.run)) return false; + if (!isRoutable(live.run.address)) return false; // Best-effort: a staleness check that cannot resolve the current // authored projection (e.g. `DefinitionProjectionMissingError` for // a pre-cutover definition with no frozen wire projection at all) @@ -617,24 +694,63 @@ export function createHubChatPlatform( // "nothing to reconcile", the same posture `sweepTerminalRuns` // takes for a relaunch failure. let driftedFoldedBody: Awaited>; + let driftedSources: boolean; try { driftedFoldedBody = await resolveDriftedFoldedBody( binding.tenantId, live.run, binding.foldedBody, ); + driftedSources = await hasDriftedSources(binding); } catch (cause: unknown) { wakeLogger.error`drift check for ${binding.roomAddress} (run ${live.run.id}) failed, leaving it as-is: ${ cause instanceof Error ? cause.message : String(cause) }`; - return; + return false; } - if (driftedFoldedBody === undefined) return; - wakeLogger.info`relaunching ${binding.roomAddress}: run ${live.run.id} is routable but its deployed definition has drifted from the current authored projection; minting a fresh run`; + if (driftedFoldedBody === undefined && !driftedSources) return false; + const reason = + driftedFoldedBody !== undefined + ? "its deployed definition has drifted from the current authored projection" + : "its deployed inference chain no longer matches the tenant catalog"; + wakeLogger.info`relaunching ${binding.roomAddress}: run ${live.run.id} is routable but ${reason}; minting a fresh run`; + sourcesCheckedAt.delete(binding.stableId); await relaunchTerminalRun({ - binding: { ...binding, foldedBody: driftedFoldedBody }, + binding: { + ...binding, + foldedBody: driftedFoldedBody ?? binding.foldedBody, + }, run: live.run, }); + return true; + } + + const RECONCILE_SOURCES_LIMIT = 200; + + async function reconcileInferenceSources( + tenantId: string, + ): Promise<{ scanned: number; relaunched: number }> { + const participants = await listLaunchesForTenant( + deps.db, + tenantId, + RECONCILE_SOURCES_LIMIT, + ); + let relaunched = 0; + for (const { binding, run } of participants) { + // A recent send may have stamped `sourcesCheckedAt` so + // `hasDriftedSources` would skip this pass. A provider connect is + // exactly the catalog change that interval is meant to wait for — + // drop the stamp so this sweep (and the next send) can see it. + sourcesCheckedAt.delete(binding.stableId); + try { + if (await reconcileDriftedRun(binding.roomAddress)) relaunched++; + } catch (cause: unknown) { + wakeLogger.error`inference-source reconcile for ${binding.roomAddress} (run ${run.id}) failed, leaving it as-is: ${ + cause instanceof Error ? cause.message : String(cause) + }`; + } + } + return { scanned: participants.length, relaunched }; } /** @@ -1154,5 +1270,6 @@ export function createHubChatPlatform( return Object.assign(platform, { recordActivity: (address: string) => lifecycle?.recordActivity(address), sweepTerminalRuns, + reconcileInferenceSources, }); } diff --git a/packages/chat/src/schema.ts b/packages/chat/src/schema.ts index d27578b5e..a8883afc0 100644 --- a/packages/chat/src/schema.ts +++ b/packages/chat/src/schema.ts @@ -116,6 +116,14 @@ export const workbenchLaunch = chatSchema.table("workbench_launch", { */ priorRunIds: jsonb("prior_run_ids").notNull().default([]), foldedBody: jsonb("folded_body").notNull(), + /** + * `@corbits/folded-runs`' `inferenceSourcesDigest` of the chain the + * current run last deployed with — secret included, hashed. A send + * compares it against today's resolution so a rotated API key + * reaches a live agent (CL-6687); `null` until the run's first + * deploy, and for rows that predate the column. + */ + sourcesDigest: text("sources_digest"), createdAt: timestamp("created_at", { withTimezone: true }) .notNull() .defaultNow(), diff --git a/packages/chat/test/migrations.test.ts b/packages/chat/test/migrations.test.ts index b39807782..b386b0fbf 100644 --- a/packages/chat/test/migrations.test.ts +++ b/packages/chat/test/migrations.test.ts @@ -44,6 +44,7 @@ const migrationNames = [ "0021_workbench_launch_prior_runs", "0022_agent_turns", "0023_drop_workbench_host_arm", + "0024_workbench_launch_sources_digest", ]; describeIfDb("applyChatMigrations", () => { diff --git a/packages/chat/test/platform-adapter.test.ts b/packages/chat/test/platform-adapter.test.ts index 1dddccd90..62b08bf2e 100644 --- a/packages/chat/test/platform-adapter.test.ts +++ b/packages/chat/test/platform-adapter.test.ts @@ -35,6 +35,7 @@ import { import { workbenchLaunch } from "../src/schema"; import { foldedRun, + inferenceSourcesDigest, DefinitionProjectionMissingError, } from "@corbits/folded-runs"; import { IDLE_HIBERNATE_UNDEPLOY_REASON } from "@corbits/agent-lifecycle"; @@ -215,6 +216,7 @@ function createFakeDb(opts: { currentRunId?: string; foldedBody: unknown; noopInference?: boolean; + sourcesDigest?: string | null; } | undefined; /** @@ -2766,6 +2768,16 @@ describe("createHubChatPlatform stale-definition reconciliation", () => { return { db, platform }; } + // A relaunch re-points `currentRunId`; the fixture's rows predate + // CL-6687's digest column, so the first check also records a baseline + // `sourcesDigest` — a write that is not a relaunch. + function repointOf(db: ReturnType) { + return db.updated + .filter((row) => row.table === workbenchLaunch) + .map((row) => row.values as { currentRunId?: string }) + .find((values) => values.currentRunId !== undefined); + } + test("a routable run whose deployed content differs from the current authored content is relaunched, not served as-is", async () => { const { db, platform } = createDriftFixture({ deployedSystemPrompt: STALE_SYSTEM_PROMPT, @@ -2775,8 +2787,7 @@ describe("createHubChatPlatform stale-definition reconciliation", () => { await platform.ensureAwake("run_stale@ten1.workbench.test"); - const repointed = db.updated.find((row) => row.table === workbenchLaunch) - ?.values as { currentRunId: string } | undefined; + const repointed = repointOf(db); expect(repointed?.currentRunId).toBeDefined(); expect(repointed?.currentRunId).not.toBe("run_stale"); }); @@ -2797,9 +2808,7 @@ describe("createHubChatPlatform stale-definition reconciliation", () => { await platform.ensureAwake("run_stale@ten1.workbench.test"); - expect( - db.updated.find((row) => row.table === workbenchLaunch), - ).toBeUndefined(); + expect(repointOf(db)).toBeUndefined(); }); test("sendMail redeploys an already-routable-but-drifted target before delivering — lifecycle.ensureAwake's routability check alone would have missed it", async () => { @@ -3101,3 +3110,221 @@ describe("createHubChatPlatform relaunch sweep", () => { expect(notices).toEqual([]); }); }); + +// CL-6687: inference sources — the decrypted API key included — are +// rendered into a run's deployed bytes at deploy time and never re-read. +// `foldedBody` comparison cannot see a rotated key, so a live agent kept +// sending the dead one after Settings said the new key was saved. The +// deploy now records a digest of the chain it pinned; a send (or a +// provider connect) compares it against today's resolution and relaunches +// on a mismatch. +describe("createHubChatPlatform inference-source rotation reconciliation", () => { + const FOLDED_BODY = { + systemPrompt: "be helpful", + toolPackagePins: [], + grantRequirements: [], + credentialBindings: [], + model: null, + }; + + function sourcesFor(apiKey: string) { + return { + sources: [ + { + id: "off_1", + provider: "anthropic", + baseURL: "https://inference.invalid", + apiKey, + model: "claude-sonnet-5", + }, + ], + defaultSource: "off_1", + }; + } + + function createRotationFixture(opts: { + deployedWithKey: string | null; + catalogKey: string; + }) { + resolveDefinitionSourcesResult = { + ok: true, + ...sourcesFor(opts.catalogKey), + }; + const db = createFakeDb({ + assetRow: { + tenantId: "ten_1", + creatorPrincipalId: "prin_creator", + name: "workbench-1", + displayName: null, + }, + definitionId: "wfd_unused", + workflowRunRow: { + id: "run_live", + address: "run_live@ten1.workbench.test", + principalId: "prin_room1", + definitionId: "wfd_run_clone", + status: "running", + }, + workflowDefinitionRow: { + id: "wfd_run_clone", + tenantId: "ten_1", + status: "deployed", + assetId: "asst_myra", + name: "myra", + origin: "run", + }, + workflowDefinitionRows: [ + { + id: "wfd_run_clone", + tenantId: "ten_1", + status: "deployed", + name: "myra", + assetId: "asst_myra", + origin: "run", + }, + { + id: "wfd_myra_authored", + tenantId: "ten_1", + status: "deployed", + name: "myra", + assetId: "asst_myra", + origin: "authored", + }, + ], + workbenchLaunchRow: { + tenantId: "ten_1", + instanceId: "run_live", + currentRunId: "run_live", + foldedBody: FOLDED_BODY, + sourcesDigest: + opts.deployedWithKey === null + ? null + : inferenceSourcesDigest(sourcesFor(opts.deployedWithKey)), + }, + wireProjectionsByDefinitionId: { + wfd_run_clone: inertProjection({ + id: "wfd_run_clone", + systemPrompt: FOLDED_BODY.systemPrompt, + model: null, + }), + wfd_myra_authored: inertProjection({ + id: "wfd_myra_authored", + systemPrompt: FOLDED_BODY.systemPrompt, + model: null, + }), + }, + }); + const platform = createHubChatPlatform({ + toolGrantsForPins: () => [], + db: db as never, + sessionService: createFakeSessionService(), + assetService: createFakeAssetService(), + sidecarRouter: createFakeSidecarRouter({ + routableAddresses: ["run_live@ten1.workbench.test"], + }), + eventCollectors: createFakeEventCollectors(), + }); + return { db, platform }; + } + + function repointedLaunch(db: ReturnType) { + return db.updated.find((row) => row.table === workbenchLaunch)?.values as + { currentRunId: string; sourcesDigest: string } | undefined; + } + + test("a routable run deployed with a key the catalog has since rotated is relaunched on the new key", async () => { + const { db, platform } = createRotationFixture({ + deployedWithKey: "sk-ant-expired", + catalogKey: "sk-ant-fresh", + }); + + await platform.ensureAwake("run_live@ten1.workbench.test"); + + const repointed = repointedLaunch(db); + expect(repointed?.currentRunId).toBeDefined(); + expect(repointed?.currentRunId).not.toBe("run_live"); + expect(repointed?.sourcesDigest).toBe( + inferenceSourcesDigest(sourcesFor("sk-ant-fresh")), + ); + }); + + test("re-saving the same key is not a rotation: the run is left alone", async () => { + const { db, platform } = createRotationFixture({ + deployedWithKey: "sk-ant-same", + catalogKey: "sk-ant-same", + }); + + await platform.ensureAwake("run_live@ten1.workbench.test"); + + expect(repointedLaunch(db)).toBeUndefined(); + }); + + test("a run that predates digest recording gets today's chain as its baseline and is left alone", async () => { + const { db, platform } = createRotationFixture({ + deployedWithKey: null, + catalogKey: "sk-ant-fresh", + }); + + await platform.ensureAwake("run_live@ten1.workbench.test"); + + const launchWrites = db.updated.filter( + (row) => row.table === workbenchLaunch, + ); + expect(launchWrites).toHaveLength(1); + expect(launchWrites[0]?.values).toEqual({ + sourcesDigest: inferenceSourcesDigest(sourcesFor("sk-ant-fresh")), + }); + }); + + test("the per-send chain check is throttled: a second send inside the interval resolves nothing", async () => { + const { platform } = createRotationFixture({ + deployedWithKey: "sk-ant-same", + catalogKey: "sk-ant-same", + }); + resolveDefinitionSourcesCalls.length = 0; + + await platform.ensureAwake("run_live@ten1.workbench.test"); + const afterFirst = resolveDefinitionSourcesCalls.length; + await platform.ensureAwake("run_live@ten1.workbench.test"); + + expect(afterFirst).toBeGreaterThan(0); + expect(resolveDefinitionSourcesCalls).toHaveLength(afterFirst); + }); + + test("reconcileInferenceSources sweeps a tenant's live participants the moment a credential lands", async () => { + const { db, platform } = createRotationFixture({ + deployedWithKey: "sk-ant-expired", + catalogKey: "sk-ant-fresh", + }); + + const swept = await platform.reconcileInferenceSources("ten_1"); + + expect(swept).toEqual({ scanned: 1, relaunched: 1 }); + expect(repointedLaunch(db)?.currentRunId).not.toBe("run_live"); + }); + + test("connecting a provider then sending within the check interval relaunches onto the new key", async () => { + const { db, platform } = createRotationFixture({ + deployedWithKey: "sk-ant-expired", + catalogKey: "sk-ant-expired", + }); + + await platform.ensureAwake("run_live@ten1.workbench.test"); + expect(repointedLaunch(db)).toBeUndefined(); + + resolveDefinitionSourcesResult = { + ok: true, + ...sourcesFor("sk-ant-fresh"), + }; + + await platform.reconcileInferenceSources("ten_1"); + await platform.ensureAwake("run_live@ten1.workbench.test"); + + const repointed = repointedLaunch(db); + expect(repointed?.currentRunId).toBeDefined(); + expect(repointed?.currentRunId).not.toBe("run_live"); + expect(repointed?.sourcesDigest).toBe( + inferenceSourcesDigest(sourcesFor("sk-ant-fresh")), + ); + }); +}); diff --git a/packages/folded-runs/src/index.ts b/packages/folded-runs/src/index.ts index 5f061f1a9..f249d0335 100644 --- a/packages/folded-runs/src/index.ts +++ b/packages/folded-runs/src/index.ts @@ -32,15 +32,21 @@ export { export { deployAtHead, foldedRunSourceRef, + inferenceSourcesDigest, launchFoldedRun, mintFoldedRun, parseSourcesOverride, + resolveLaunchSources, SourcesOverride, InferenceResolutionError, + type DeployedAtHead, type FoldedRunMode, type LaunchFoldedRunParams, + type LaunchedAndDeployedFoldedRun, type MintFoldedRunParams, type LaunchedFoldedRun, + type ResolvedLaunchSources, + type ResolveLaunchSourcesParams, } from "./launch"; export { wakeFoldedRun, type WakeFoldedRunParams } from "./wake"; export { diff --git a/packages/folded-runs/src/launch.test.ts b/packages/folded-runs/src/launch.test.ts new file mode 100644 index 000000000..26085c827 --- /dev/null +++ b/packages/folded-runs/src/launch.test.ts @@ -0,0 +1,58 @@ +import { describe, expect, test } from "bun:test"; + +import { inferenceSourcesDigest, type ResolvedLaunchSources } from "./launch"; + +const anthropic = { + id: "off_anthropic", + provider: "anthropic", + baseURL: "https://api.anthropic.com", + apiKey: "sk-ant-old", + model: "claude-sonnet-5", +}; + +const resolved: ResolvedLaunchSources = { + sources: [anthropic], + defaultSource: "off_anthropic", +}; + +describe("inferenceSourcesDigest", () => { + test("is stable across key insertion order", () => { + const reordered: ResolvedLaunchSources = { + defaultSource: "off_anthropic", + sources: [ + { + model: anthropic.model, + apiKey: anthropic.apiKey, + baseURL: anthropic.baseURL, + provider: anthropic.provider, + id: anthropic.id, + }, + ], + }; + expect(inferenceSourcesDigest(reordered)).toBe( + inferenceSourcesDigest(resolved), + ); + }); + + test("changes when the secret rotates, without carrying the secret", () => { + const rotated: ResolvedLaunchSources = { + ...resolved, + sources: [{ ...anthropic, apiKey: "sk-ant-new" }], + }; + const digest = inferenceSourcesDigest(rotated); + expect(digest).not.toBe(inferenceSourcesDigest(resolved)); + expect(digest).not.toContain("sk-ant"); + expect(digest).toMatch(/^[0-9a-f]{64}$/); + }); + + test("changes when the default source moves", () => { + const openai = { ...anthropic, id: "off_openai", provider: "openai" }; + const chain: ResolvedLaunchSources = { + sources: [anthropic, openai], + defaultSource: "off_anthropic", + }; + expect( + inferenceSourcesDigest({ ...chain, defaultSource: "off_openai" }), + ).not.toBe(inferenceSourcesDigest(chain)); + }); +}); diff --git a/packages/folded-runs/src/launch.ts b/packages/folded-runs/src/launch.ts index 1d6f0e551..2fdfe0d3d 100644 --- a/packages/folded-runs/src/launch.ts +++ b/packages/folded-runs/src/launch.ts @@ -11,6 +11,7 @@ // share that shape; what still distinguishes a folded run is its // `principalId`, whose `agent_session` join is what makes the run's // mailbox listable through the platform's sanctioned per-run surfaces. +import { createHash } from "node:crypto"; import { eq } from "drizzle-orm"; import { type } from "arktype"; import type { DBExecutor } from "@intx/db"; @@ -234,63 +235,56 @@ async function markRunDeployClone( * failure-path rollback of those rows — this function only throws. */ -export async function deployAtHead( - deps: Pick< - FoldedRunsDeps, - | "db" - | "sessionService" - | "sidecarRouter" - | "eventCollectors" - | "credentialCipher" - | "assetService" - | "toolGrantsForPins" - | "mcpCredentialBindingsFor" - >, - params: { - tenantId: string; - instanceId: string; - triggerAddress: string; - principalId: string; - sessionId: string; - foldedBody: FoldedBody; - /** Named in the "seed a tenant catalog source" error, e.g. "the workbench host", "the invited agent", or "the woken instance". */ - launchLabel: string; - /** - * When present, used verbatim in place of `resolveDefinitionSources` - * — the tenant catalog is never touched, so a launch pinned this - * way needs no catalog source to exist at all. Absent, this is - * the ordinary catalog-resolved path every launch used before this - * override existed. - */ - sources?: SourcesOverride; - /** - * The tenant's current live default model — the same (provider, - * model) Settings' "Default model & fallbacks" row heads today - * (see `@corbits/chat`'s `resolveFallbackModel`). Tried first when - * `foldedBody.model` is `null` (a definition that declares no model - * of its own); tried SECOND, as a retry, when `foldedBody.model` is - * set but does not resolve against the tenant's current catalog — - * e.g. it was pinned (by a person, or by the seed that created this - * definition) against a provider the tenant has since disconnected, - * or before it connected any provider at all. Without that retry an - * agent whose pin predates the tenant's current connection stays - * permanently dead even after a working provider is connected, - * because `foldedBody.model` alone would keep resolving to the same - * unlaunchable name forever. Absent, a `foldedBody.model` that fails - * to resolve — or a `null` one with nothing to fall back to — - * resolves to the loud `InferenceResolutionError` it always has. - */ - fallbackModel?: string; - /** - * The shape the run's deployed definition takes. Defaults to the - * folded conversational step every launcher uses today; CL-6329's - * per-turn swap passes `{ kind: "section", turnTimeoutMs }` and - * nothing else about this call changes, because the mode is config - * data rendered into the deployed bytes rather than a branch here. - */ - mode?: FoldedRunMode; - }, -): Promise { +/** The inference chain a launch pins into its deployed bytes. */ +export type ResolvedLaunchSources = { + sources: InferenceSource[]; + defaultSource: string; +}; + +export type ResolveLaunchSourcesParams = { + tenantId: string; + foldedBody: FoldedBody; + /** Named in the "seed a tenant catalog source" error, e.g. "the workbench host", "the invited agent", or "the woken instance". */ + launchLabel: string; + /** + * When present, used verbatim in place of `resolveDefinitionSources` + * — the tenant catalog is never touched, so a launch pinned this + * way needs no catalog source to exist at all. Absent, this is + * the ordinary catalog-resolved path every launch used before this + * override existed. + */ + sources?: SourcesOverride; + /** + * The tenant's current live default model — the same (provider, + * model) Settings' "Default model & fallbacks" row heads today + * (see `@corbits/chat`'s `resolveFallbackModel`). Tried first when + * `foldedBody.model` is `null` (a definition that declares no model + * of its own); tried SECOND, as a retry, when `foldedBody.model` is + * set but does not resolve against the tenant's current catalog — + * e.g. it was pinned (by a person, or by the seed that created this + * definition) against a provider the tenant has since disconnected, + * or before it connected any provider at all. Without that retry an + * agent whose pin predates the tenant's current connection stays + * permanently dead even after a working provider is connected, + * because `foldedBody.model` alone would keep resolving to the same + * unlaunchable name forever. Absent, a `foldedBody.model` that fails + * to resolve — or a `null` one with nothing to fall back to — + * resolves to the loud `InferenceResolutionError` it always has. + */ + fallbackModel?: string; +}; + +/** + * Resolves the inference chain a folded run deploys with, exactly as + * `deployAtHead` pins it — the same call a drift check makes later to + * ask "would this run deploy with different sources today?" (CL-6687: + * a rotated API key is only ever picked up at deploy time, so the + * deployed chain has to be comparable against the current one). + */ +export async function resolveLaunchSources( + deps: Pick, + params: ResolveLaunchSourcesParams, +): Promise { const sourcesOverride = parseSourcesOverride(params.sources); const resolveAgainst = (model: string | null) => resolveDefinitionSources({ @@ -338,18 +332,92 @@ export async function deployAtHead( // A caller-supplied override already states the adapter key it wants // (see `SourcesOverride`'s doc) — only a catalog-resolved chain needs // its `provider` field corrected for Ollama offerings. - if (sourcesOverride === undefined) { - const offerings = await listVisibleOfferings(deps.db, params.tenantId); - const ollamaOfferingIds = new Set( - offerings - .filter((resolved) => resolved.provider.name === "ollama") - .map((resolved) => resolved.offering.id), - ); - resolution.sources = withOllamaAdapterKey( - resolution.sources, - ollamaOfferingIds, - ); + if (sourcesOverride !== undefined) { + return { + sources: [...resolution.sources], + defaultSource: resolution.defaultSource, + }; + } + const offerings = await listVisibleOfferings(deps.db, params.tenantId); + const ollamaOfferingIds = new Set( + offerings + .filter((resolved) => resolved.provider.name === "ollama") + .map((resolved) => resolved.offering.id), + ); + return { + sources: withOllamaAdapterKey(resolution.sources, ollamaOfferingIds), + defaultSource: resolution.defaultSource, + }; +} + +function canonicalJSON(value: unknown): string { + if (Array.isArray(value)) { + return `[${value.map(canonicalJSON).join(",")}]`; + } + if (value !== null && typeof value === "object") { + const entries = Object.entries(value as Record) + .filter(([, v]) => v !== undefined) + .sort(([a], [b]) => (a < b ? -1 : a > b ? 1 : 0)); + return `{${entries + .map(([k, v]) => `${JSON.stringify(k)}:${canonicalJSON(v)}`) + .join(",")}}`; } + return JSON.stringify(value); +} + +/** + * A stable fingerprint of a resolved inference chain — every field the + * deployed bytes carry, the decrypted secret included, so a rotated key + * (or a moved base URL, or a re-pointed default) changes the digest + * while a re-save of the same key does not. Only the hash is ever + * stored; the secret itself never leaves the resolution. + */ +export function inferenceSourcesDigest( + resolved: ResolvedLaunchSources, +): string { + return createHash("sha256") + .update( + canonicalJSON({ + sources: resolved.sources, + defaultSource: resolved.defaultSource, + }), + ) + .digest("hex"); +} + +export type DeployedAtHead = { + /** `inferenceSourcesDigest` of the chain this deploy pinned. */ + readonly sourcesDigest: string; +}; + +export async function deployAtHead( + deps: Pick< + FoldedRunsDeps, + | "db" + | "sessionService" + | "sidecarRouter" + | "eventCollectors" + | "credentialCipher" + | "assetService" + | "toolGrantsForPins" + | "mcpCredentialBindingsFor" + >, + params: ResolveLaunchSourcesParams & { + instanceId: string; + triggerAddress: string; + principalId: string; + sessionId: string; + /** + * The shape the run's deployed definition takes. Defaults to the + * folded conversational step every launcher uses today; CL-6329's + * per-turn swap passes `{ kind: "section", turnTimeoutMs }` and + * nothing else about this call changes, because the mode is config + * data rendered into the deployed bytes rather than a branch here. + */ + mode?: FoldedRunMode; + }, +): Promise { + const resolution = await resolveLaunchSources(deps, params); // `create` replaces any collector already registered for this // address (see `EventCollectorRegistry.create`), so this is @@ -580,6 +648,7 @@ export async function deployAtHead( `${params.launchLabel}: deployment ${params.triggerAddress} is not routable for run ${params.instanceId}; cannot deliver its run grants, so the run would start under-authorized`, ); } + return { sourcesDigest: inferenceSourcesDigest(resolution) }; } export type LaunchFoldedRunParams = { @@ -615,6 +684,8 @@ export type LaunchedFoldedRun = { readonly sessionId: string; }; +export type LaunchedAndDeployedFoldedRun = LaunchedFoldedRun & DeployedAtHead; + export type MintFoldedRunParams = Pick< LaunchFoldedRunParams, "tenantId" | "instanceId" | "triggerAddress" | "definitionId" | "persistExtra" @@ -735,7 +806,7 @@ export async function mintFoldedRun( export async function launchFoldedRun( deps: FoldedRunsDeps, params: LaunchFoldedRunParams, -): Promise { +): Promise { const { instancePrincipalId, sessionId } = await mintFoldedRun(deps, { tenantId: params.tenantId, instanceId: params.instanceId, @@ -746,6 +817,7 @@ export async function launchFoldedRun( : {}), }); + let deployed: DeployedAtHead; try { // Opened/deployed via the shared `deployAtHead` step, mirroring // the reference route: the run's runtime status/readiness @@ -761,7 +833,7 @@ export async function launchFoldedRun( foldedBody: params.foldedBody, launchLabel: params.launchLabel, }; - await deployAtHead(deps, { + deployed = await deployAtHead(deps, { ...deployAtHeadParams, ...(params.sources !== undefined ? { sources: params.sources } : {}), ...(params.fallbackModel !== undefined @@ -810,5 +882,5 @@ export async function launchFoldedRun( throw err; } - return { instancePrincipalId, sessionId }; + return { instancePrincipalId, sessionId, ...deployed }; } diff --git a/packages/folded-runs/src/wake.ts b/packages/folded-runs/src/wake.ts index 12813ddce..8fe0404de 100644 --- a/packages/folded-runs/src/wake.ts +++ b/packages/folded-runs/src/wake.ts @@ -10,6 +10,7 @@ import { eq } from "drizzle-orm"; import { sessionAsset } from "@intx/db/schema"; import { deployAtHead, + type DeployedAtHead, type FoldedRunMode, type SourcesOverride, } from "./launch"; @@ -53,7 +54,7 @@ export type WakeFoldedRunParams = { export async function wakeFoldedRun( deps: FoldedRunsDeps, params: WakeFoldedRunParams, -): Promise { +): Promise { if (params.principalId === null) { throw new Error(`Run "${params.instanceId}" has no principal`); } @@ -84,7 +85,7 @@ export async function wakeFoldedRun( launchLabel: "the woken instance", }; try { - await deployAtHead(deps, { + return await deployAtHead(deps, { ...deployAtHeadParams, ...(params.sources !== undefined ? { sources: params.sources } : {}), ...(params.mode !== undefined ? { mode: params.mode } : {}), diff --git a/packages/webhook-triggers/src/launch.ts b/packages/webhook-triggers/src/launch.ts index ae847b70a..2edb06c62 100644 --- a/packages/webhook-triggers/src/launch.ts +++ b/packages/webhook-triggers/src/launch.ts @@ -63,6 +63,17 @@ export type LaunchWebhookTriggerDeps = FoldedRunsDeps & { readonly instanceId: string; readonly foldedBody: LaunchFoldedRunParams["foldedBody"]; }) => NonNullable; + /** + * Records the inference chain the launch just deployed with on that + * same mapping row — the host wires this to `@corbits/chat`'s + * `recordSourcesDigest`. `persistLaunch` runs before the deploy + * resolves the chain, so the digest lands in a second write; without + * it a rotated provider key never reaches this run (CL-6687). + */ + recordLaunchSources: (input: { + readonly instanceId: string; + readonly sourcesDigest: string; + }) => Promise; }; export type LaunchedWebhookTrigger = { @@ -147,6 +158,10 @@ export async function launchWebhookTrigger( foldedBody, }), }); + await deps.recordLaunchSources({ + instanceId, + sourcesDigest: launched.sourcesDigest, + }); const content = renderInputTemplate(trigger.inputTemplate, payload); const cryptoProvider = await deps.cryptoProviderCache.get(instanceId); diff --git a/packages/webhook-triggers/test/launch.test.ts b/packages/webhook-triggers/test/launch.test.ts index 25c9393f6..d92b0410c 100644 --- a/packages/webhook-triggers/test/launch.test.ts +++ b/packages/webhook-triggers/test/launch.test.ts @@ -29,7 +29,11 @@ mock.module("@corbits/folded-runs", () => ({ readFoldedBody: () => FOLDED_BODY, launchFoldedRun: async (...args: unknown[]) => { launchFoldedRunCalls.push(args); - return { instancePrincipalId: "prn_run1", sessionId: "ses_run1" }; + return { + instancePrincipalId: "prn_run1", + sessionId: "ses_run1", + sourcesDigest: "digest_run1", + }; }, sendFoldedMailWithRetry: async (...args: unknown[]) => { sendFoldedMailWithRetryCalls.push(args); @@ -88,6 +92,7 @@ const TRIGGER = { }; let persistLaunchCalls: unknown[] = []; +let recordLaunchSourcesCalls: unknown[] = []; const persistedLaunchExtra = async () => {}; function baseDeps() { @@ -104,6 +109,9 @@ function baseDeps() { persistLaunchCalls.push(input); return persistedLaunchExtra; }, + recordLaunchSources: async (input: unknown) => { + recordLaunchSourcesCalls.push(input); + }, }; } @@ -202,6 +210,7 @@ describe("launchWebhookTrigger", () => { launchFoldedRunCalls = []; sendFoldedMailWithRetryCalls = []; persistLaunchCalls = []; + recordLaunchSourcesCalls = []; sendFoldedMailWithRetryResult = { ok: true, mail: { id: "m_1" } }; const result = await launchWebhookTrigger(baseDeps(), TRIGGER, { @@ -221,5 +230,11 @@ describe("launchWebhookTrigger", () => { foldedBody: FOLDED_BODY, }, ]); + // CL-6687: the mapping row is written before the deploy resolves the + // inference chain, so the digest a rotation check compares against + // has to land in a second write once the launch returns. + expect(recordLaunchSourcesCalls).toEqual([ + { instanceId: result.instanceId, sourcesDigest: "digest_run1" }, + ]); }); }); diff --git a/scripts/e2e/harness.test.ts b/scripts/e2e/harness.test.ts index cb7ccbc9e..5b4d14800 100644 --- a/scripts/e2e/harness.test.ts +++ b/scripts/e2e/harness.test.ts @@ -9,6 +9,7 @@ import { existsSync } from "node:fs"; import { createCleanupHarness, + parseEnvFileDatabaseUrl, runCleanups, type SpawnedApp, } from "./harness.ts"; @@ -65,3 +66,40 @@ describe("createCleanupHarness", () => { expect(cleanups).toHaveLength(0); }); }); + +describe("parseEnvFileDatabaseUrl", () => { + test("reads DATABASE_URL and ignores comments and blanks", () => { + expect( + parseEnvFileDatabaseUrl( + "# comment\n\nFOO=bar\nDATABASE_URL=postgres://localhost:5432/workbench\n", + ), + ).toBe("postgres://localhost:5432/workbench"); + }); + + test("strips surrounding quotes and trims", () => { + expect( + parseEnvFileDatabaseUrl( + ` DATABASE_URL="postgres://localhost:5432/workbench" \n`, + ), + ).toBe("postgres://localhost:5432/workbench"); + expect( + parseEnvFileDatabaseUrl( + `DATABASE_URL='postgres://localhost:5432/workbench'\n`, + ), + ).toBe("postgres://localhost:5432/workbench"); + }); + + test("ignores a commented DATABASE_URL and returns undefined when none is set", () => { + expect( + parseEnvFileDatabaseUrl("# DATABASE_URL=postgres://commented\nFOO=bar\n"), + ).toBeUndefined(); + }); + + test("last DATABASE_URL wins", () => { + expect( + parseEnvFileDatabaseUrl( + "DATABASE_URL=postgres://first/db\nDATABASE_URL=postgres://second/db\n", + ), + ).toBe("postgres://second/db"); + }); +}); diff --git a/scripts/e2e/harness.ts b/scripts/e2e/harness.ts index 3704c9ae7..a4d7f0caf 100644 --- a/scripts/e2e/harness.ts +++ b/scripts/e2e/harness.ts @@ -5,6 +5,7 @@ // real the suite fails and says which hop, it never fakes the result. import { afterAll } from "bun:test"; +import { readFileSync } from "node:fs"; import { mkdtemp, rm } from "node:fs/promises"; import { tmpdir } from "node:os"; import path from "node:path"; @@ -23,9 +24,17 @@ const SIDECAR_DIR = path.join(REPO_ROOT, "apps", "sidecar"); * database still runs the unit gates); in CI, E2E_REQUIRED=1 turns * that skip into a loud failure so the suite can never silently * vanish from the pipeline. + * + * `bun test` workers do not always inherit process.env.DATABASE_URL + * even when the shell sourced `.env`. Fall back to the same repo-root + * `.env` file `bun run` would load. */ export function e2eDatabaseUrl(): string | undefined { - const url = process.env["DATABASE_URL"]; + const fromProcess = process.env["DATABASE_URL"]; + const url = + fromProcess !== undefined && fromProcess !== "" + ? fromProcess + : databaseUrlFromRepoEnvFile(); if (url !== undefined && url !== "") return baseUrlToE2eUrl(url); if (process.env["E2E_REQUIRED"] === "1") { throw new Error( @@ -36,6 +45,51 @@ export function e2eDatabaseUrl(): string | undefined { return undefined; } +/** + * Pull DATABASE_URL from a dotenv-style file body. Ignores blanks and + * comments; trims; strips one layer of surrounding quotes. Last match + * wins. Does not expand interpolations. + */ +export function parseEnvFileDatabaseUrl(text: string): string | undefined { + let found: string | undefined; + for (const rawLine of text.split(/\r?\n/)) { + const line = rawLine.trim(); + if (line === "" || line.startsWith("#")) continue; + if (!line.startsWith("DATABASE_URL=")) continue; + let value = line.slice("DATABASE_URL=".length).trim(); + if (value.length >= 2) { + const start = value[0]; + const end = value[value.length - 1]; + if ((start === '"' && end === '"') || (start === "'" && end === "'")) { + value = value.slice(1, -1); + } + } + found = value; + } + if (found === undefined || found === "") return undefined; + return found; +} + +function databaseUrlFromRepoEnvFile(): string | undefined { + let text: string; + try { + text = readFileSync(path.join(REPO_ROOT, ".env"), "utf8"); + } catch (error) { + if (isEnoent(error)) return undefined; + throw error; + } + return parseEnvFileDatabaseUrl(text); +} + +function isEnoent(error: unknown): boolean { + return ( + typeof error === "object" && + error !== null && + "code" in error && + error.code === "ENOENT" + ); +} + /** * The suite tears its schema down and rebuilds it on every run, so it * must never run inside the developer's own database. It derives a