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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 12 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
29 changes: 27 additions & 2 deletions apps/hub/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -286,6 +288,7 @@ import {
} from "@workbench/onboarding";
import {
createConnectionRoutes,
isInferenceProvider,
createMcpOAuthRoutes,
createMcpServerRoutes,
createOAuthConnectRoutes,
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand All @@ -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-
Expand Down
16 changes: 15 additions & 1 deletion apps/hub/src/routine-launcher.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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);
Expand Down Expand Up @@ -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);
},
}),
}),
};
}

Expand Down
2 changes: 2 additions & 0 deletions apps/hub/src/routine-launcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -199,6 +200,7 @@ export function createHubRoutineLauncher(
foldedBody,
}),
});
await recordSourcesDigest(deps.db, instanceId, launched.sourcesDigest);

if (
input.deliveryWorkbenchId !== undefined &&
Expand Down
28 changes: 28 additions & 0 deletions docs/credential-wiring.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
49 changes: 48 additions & 1 deletion packages/chat/src/agent-binding.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

/**
Expand Down Expand Up @@ -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,
};
}

Expand Down Expand Up @@ -288,16 +291,60 @@ export async function repointBinding(
db: DB["db"],
binding: AgentBinding,
newRunId: string,
sourcesDigest: string,
): Promise<void> {
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<void> {
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<LiveAgent[]> {
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
Expand Down
1 change: 1 addition & 0 deletions packages/chat/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
7 changes: 7 additions & 0 deletions packages/chat/src/migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
`,
},
];

/**
Expand Down
Loading
Loading