From 4edfca82404fc88d8f22fca4da17cd5e006bf17f Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 28 Aug 2026 05:17:52 -0700 Subject: [PATCH 01/12] Add tests for inference-source rotation reaching live agents --- apps/hub/src/routine-launcher.test.ts | 16 +- packages/chat/test/platform-adapter.test.ts | 212 +++++++++++++++++- packages/folded-runs/src/launch.test.ts | 58 +++++ packages/webhook-triggers/test/launch.test.ts | 17 +- 4 files changed, 296 insertions(+), 7 deletions(-) create mode 100644 packages/folded-runs/src/launch.test.ts 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/packages/chat/test/platform-adapter.test.ts b/packages/chat/test/platform-adapter.test.ts index 1dddccd90..1745be0cc 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,196 @@ 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"); + }); +}); 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/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" }, + ]); }); }); From 8b987aaf8d3b92bb88dfd3211de0a545bb108461 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 28 Aug 2026 05:17:58 -0700 Subject: [PATCH 02/12] Relaunch live agents when their inference credential rotates MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Inference sources are resolved and rendered into a run's deployed bytes once; a rotated provider key only ever reached the next deploy, so an open workbench kept sending the dead key after Settings said the new one was saved. Every deploy now records a digest of the resolved chain on the run's launch row — room invites, routine fires and webhook launches alike; the send-time drift check and a new provider-connect hook compare it against the current catalog and relaunch on a mismatch. Fixes CL-6687. --- apps/hub/src/index.ts | 29 +++- apps/hub/src/routine-launcher.ts | 2 + packages/chat/src/agent-binding.ts | 49 +++++- packages/chat/src/index.ts | 1 + packages/chat/src/migrations.ts | 7 + packages/chat/src/platform-adapter.ts | 136 +++++++++++++-- packages/chat/src/schema.ts | 8 + packages/folded-runs/src/index.ts | 6 + packages/folded-runs/src/launch.ts | 214 ++++++++++++++++-------- packages/folded-runs/src/wake.ts | 5 +- packages/webhook-triggers/src/launch.ts | 15 ++ 11 files changed, 384 insertions(+), 88 deletions(-) diff --git a/apps/hub/src/index.ts b/apps/hub/src/index.ts index 483887d8e..1a69adffb 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.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/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..2e7d01d59 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,58 @@ 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) { + 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 +1265,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/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.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); From d733e0d3bfca2849f26d0487c2c7cf8ab2275d49 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 28 Aug 2026 05:18:03 -0700 Subject: [PATCH 03/12] Update docs: credential rotation reaches live agents --- docs/credential-wiring.md | 28 ++++++++++++++++++++++++++++ 1 file changed, 28 insertions(+) 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. From 73925dd527242b3322de47b1e5d9fc001c723adc Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 08:33:56 -0700 Subject: [PATCH 04/12] Restore the lockfile to match main The credential-rotation rebase carried a one-line bun.lock drift that would fail bun install --frozen-lockfile. This PR does not add a workspace package. --- bun.lock | 1 - 1 file changed, 1 deletion(-) diff --git a/bun.lock b/bun.lock index 74c1a084f..96da4beb8 100644 --- a/bun.lock +++ b/bun.lock @@ -1251,7 +1251,6 @@ "@corbits/api-query": "workspace:*", "@corbits/bench-ui": "workspace:*", "@corbits/chat-ui": "workspace:*", - "@corbits/error-sink": "workspace:*", "@corbits/icons": "workspace:*", "@corbits/inference-settings": "workspace:*", "@corbits/react-ui": "github:corbitsdev/react-ui#3b122812a307ccb35be31386f7696020c5a84635", From e7ace847637f3477f993de6a968961bd7e983f1a Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 10:17:42 -0700 Subject: [PATCH 05/12] Add 0024 to the chat migrations test list --- packages/chat/test/migrations.test.ts | 1 + 1 file changed, 1 insertion(+) 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", () => { From 969c8eb121bb8ee515d2903aaed99ce3f29a7390 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 10:31:03 -0700 Subject: [PATCH 06/12] Pass DATABASE_URL into walking-skeleton package suites --- .github/workflows/ci.yml | 3 +++ 1 file changed, 3 insertions(+) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index b65995898..fd97c7d47 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 @@ -191,6 +192,7 @@ jobs: run: 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 @@ -203,3 +205,4 @@ jobs: bun 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 From 634e76f4ace29cd8695ce0222bef17707e49cf6f Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 10:43:56 -0700 Subject: [PATCH 07/12] Restore the lockfile to match main --- bun.lock | 1 + 1 file changed, 1 insertion(+) diff --git a/bun.lock b/bun.lock index 96da4beb8..74c1a084f 100644 --- a/bun.lock +++ b/bun.lock @@ -1251,6 +1251,7 @@ "@corbits/api-query": "workspace:*", "@corbits/bench-ui": "workspace:*", "@corbits/chat-ui": "workspace:*", + "@corbits/error-sink": "workspace:*", "@corbits/icons": "workspace:*", "@corbits/inference-settings": "workspace:*", "@corbits/react-ui": "github:corbitsdev/react-ui#3b122812a307ccb35be31386f7696020c5a84635", From 8aa9982a9c2ed30ea7b929103bd03966692b6a14 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 10:47:36 -0700 Subject: [PATCH 08/12] Export .env before walking-skeleton package suites --- .github/workflows/ci.yml | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index fd97c7d47..4d0835dca 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -189,7 +189,11 @@ 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 @@ -202,6 +206,9 @@ jobs: # would-be skip into a hard failure. - name: Run the database-backed package suites run: | + set -a + . ./.env + set +a bun test $(grep -rl DATABASE_URL --include='*.test.ts' packages apps | grep -v apps/hub/test) env: E2E_REQUIRED: "1" From b12d190a93fd13f1fa937c84a7037f1f268cf1f6 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 10:57:46 -0700 Subject: [PATCH 09/12] Load root .env for walking-skeleton package suites --- .github/workflows/ci.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 4d0835dca..07a378850 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -209,7 +209,7 @@ jobs: set -a . ./.env set +a - bun test $(grep -rl DATABASE_URL --include='*.test.ts' packages apps | grep -v apps/hub/test) + 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 From 0ff30169c496f47d97f19866321209fa77ecc867 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 11:10:50 -0700 Subject: [PATCH 10/12] Read DATABASE_URL from .env in e2eDatabaseUrl --- scripts/e2e/harness.test.ts | 38 +++++++++++++++++++++++++ scripts/e2e/harness.ts | 56 ++++++++++++++++++++++++++++++++++++- 2 files changed, 93 insertions(+), 1 deletion(-) 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 From f9b021d8fff478cda9bb95e322899d386be86c03 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 14:36:02 -0700 Subject: [PATCH 11/12] Add tests for connect-path inference rotation --- packages/chat/test/platform-adapter.test.ts | 25 +++++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/packages/chat/test/platform-adapter.test.ts b/packages/chat/test/platform-adapter.test.ts index 1745be0cc..62b08bf2e 100644 --- a/packages/chat/test/platform-adapter.test.ts +++ b/packages/chat/test/platform-adapter.test.ts @@ -3302,4 +3302,29 @@ describe("createHubChatPlatform inference-source rotation reconciliation", () => 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")), + ); + }); }); From 3c0179132534604ba0e8d4b2b6017dfa2f316d65 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Sat, 29 Aug 2026 14:36:05 -0700 Subject: [PATCH 12/12] Invalidate the sources-check cache on provider connect A recent send stamped the 30s interval, so a connect inside that window never reached the drift check. Drop the stamp when the tenant's inference sources are reconciled so the sweep and the next send can see the new key. --- packages/chat/src/platform-adapter.ts | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/packages/chat/src/platform-adapter.ts b/packages/chat/src/platform-adapter.ts index 2e7d01d59..7338e6198 100644 --- a/packages/chat/src/platform-adapter.ts +++ b/packages/chat/src/platform-adapter.ts @@ -737,6 +737,11 @@ export function createHubChatPlatform( ); 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) {