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
Original file line number Diff line number Diff line change
Expand Up @@ -144,6 +144,7 @@ const fakeInstanceRegistryLayer = Layer.succeed(ProviderInstanceRegistry.Provide
streamChanges: Stream.empty,
// Tests never drive changes through this fake; acquire a throwaway
// subscription on an unused PubSub so the shape is satisfied.
listEnvironments: Effect.succeed(new Map()),
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), (pubsub) => PubSub.subscribe(pubsub)),
});

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,9 @@
* ----------
* On layer build we:
* 1. Read the current `ServerSettings` once and use it to seed the
* registry's initial state via `ProviderInstanceRegistryMutableLayer`.
* registry's initial state via `ProviderInstanceRegistryMutableLayer`,
* priming that snapshot's secret references first so the whole boot
* fleet costs one unlock rather than one per instance.
* 2. Fork a daemon fiber (lifetime tied to the layer's scope) that
* subscribes to `ServerSettingsService.streamChanges` and calls
* `ProviderInstanceRegistryMutator.reconcile` on every emission.
Expand All @@ -53,8 +55,13 @@ import * as Stream from "effect/Stream";

import { ServerSettingsService } from "../../serverSettings.ts";
import { BUILT_IN_DRIVERS, type BuiltInDriversEnv } from "../builtInDrivers.ts";
import { collectProviderSecretReferences } from "../ProviderSecretReference.ts";
import { ProviderInstanceRegistry } from "../Services/ProviderInstanceRegistry.ts";
import { ProviderInstanceRegistryMutator } from "../Services/ProviderInstanceRegistryMutator.ts";
import {
ProviderSecretResolver,
type ProviderSecretResolverShape,
} from "../Services/ProviderSecretResolver.ts";
import { ProviderInstanceRegistryMutableLayer } from "./ProviderInstanceRegistryLive.ts";

/**
Expand Down Expand Up @@ -103,6 +110,20 @@ export const deriveProviderInstanceConfigMap = (
return merged as ProviderInstanceConfigMap;
};

/**
* Read every secret the config map's instances are about to need, before any
* of them is built. Each instance resolves its own environment, and the secret
* store charges an unlock per read rather than per secret, so without this a
* fleet of five reference-backed providers is five authorizations.
*/
const primeConfigMapSecrets = (
secretResolver: ProviderSecretResolverShape,
configMap: ProviderInstanceConfigMap,
) =>
secretResolver.prime(
collectProviderSecretReferences(Object.values(configMap).map((entry) => entry.environment)),
);

/**
* Layer that consumes `ProviderInstanceRegistryMutator` and forks a
* settings-watcher fiber. The fiber's lifetime is tied to the enclosing
Expand All @@ -118,16 +139,18 @@ const SettingsWatcherLive = Layer.effectDiscard(
Effect.gen(function* () {
const mutator = yield* ProviderInstanceRegistryMutator;
const serverSettings = yield* ServerSettingsService;
const secretResolver = yield* ProviderSecretResolver;
yield* serverSettings.streamChanges.pipe(
Stream.runForEach((next) =>
mutator
.reconcile(deriveProviderInstanceConfigMap(next))
Stream.runForEach((next) => {
const configMap = deriveProviderInstanceConfigMap(next);
return primeConfigMapSecrets(secretResolver, configMap)
.pipe(Effect.andThen(mutator.reconcile(configMap)))
.pipe(
Effect.catchCause((cause) =>
Effect.logError("ProviderInstanceRegistry reconcile failed", cause),
),
),
),
);
}),
Comment thread
yordis marked this conversation as resolved.
Effect.forkScoped,
);
}),
Expand Down Expand Up @@ -156,6 +179,7 @@ export const ProviderInstanceRegistryHydrationLive: Layer.Layer<
> = Layer.unwrap(
Effect.gen(function* () {
const serverSettings = yield* ServerSettingsService;
const secretResolver = yield* ProviderSecretResolver;
const initialSettings: ServerSettings | undefined = yield* serverSettings.getSettings.pipe(
Effect.orElseSucceed(() => undefined),
);
Expand All @@ -164,6 +188,10 @@ export const ProviderInstanceRegistryHydrationLive: Layer.Layer<
? ({} as ProviderInstanceConfigMap)
: deriveProviderInstanceConfigMap(initialSettings);

// The watcher only sees later writes, so this snapshot is the whole boot
// fleet and the one place where nothing is cached yet.
yield* primeConfigMapSecrets(secretResolver, initialConfigMap);

const mutableLayer = ProviderInstanceRegistryMutableLayer({
drivers: BUILT_IN_DRIVERS,
configMap: initialConfigMap,
Expand Down
13 changes: 13 additions & 0 deletions apps/server/src/provider/Layers/ProviderInstanceRegistryLive.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ import {
ProviderInstanceId,
type ProviderInstanceConfig,
type ProviderInstanceConfigMap,
type ProviderInstanceEnvironment,
type ProviderDriverKind,
type ServerProvider,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -520,6 +521,18 @@ export const makeProviderInstanceRegistry = <R>(input: {
listUnavailable: Ref.get(unavailable).pipe(
Effect.map((map) => Array.from(map.values()) as ReadonlyArray<ServerProvider>),
),
listEnvironments: Effect.gen(function* () {
const environments = new Map<ProviderInstanceId, ProviderInstanceEnvironment | undefined>();
for (const [instanceId, entry] of yield* Ref.get(rebuildable)) {
environments.set(instanceId, entry.environment);
}
// Live entries are written second so an instance that is both live and
// pending a retry reports the configuration it is actually running.
for (const [instanceId, live] of yield* Ref.get(entries)) {
environments.set(instanceId, live.entry.environment);
}
return environments;
}),
rebuildInstanceWhen,
// Getters: each read constructs a fresh Stream / Effect descriptor
// so multiple consumers don't share a single already-started
Expand Down
193 changes: 193 additions & 0 deletions apps/server/src/provider/Layers/ProviderRegistry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -881,6 +881,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
listUnavailable: Effect.succeed([]),
rebuildInstanceWhen: () => Effect.succeed(false),
streamChanges: Stream.empty,
listEnvironments: Effect.succeed(new Map()),
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), PubSub.subscribe),
},
);
Expand Down Expand Up @@ -955,6 +956,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
const rebuiltIds = yield* Ref.make<ReadonlyArray<ProviderInstanceId>>([]);
const secretResolverLayer = Layer.succeed(ProviderSecretResolver, {
resolve: (environment) => Effect.succeed({ variables: environment, unresolved: [] }),
prime: () => Effect.void,
invalidate: Ref.update(invalidations, (count) => count + 1),
});
const instanceRegistryLayer = Layer.succeed(
Expand All @@ -976,6 +978,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
)
: Effect.succeed(false),
streamChanges: Stream.empty,
listEnvironments: Effect.succeed(new Map()),
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), PubSub.subscribe),
},
);
Expand Down Expand Up @@ -1031,6 +1034,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
const rebuiltIds = yield* Ref.make<ReadonlyArray<ProviderInstanceId>>([]);
const secretResolverLayer = Layer.succeed(ProviderSecretResolver, {
resolve: (environment) => Effect.succeed({ variables: environment, unresolved: [] }),
prime: () => Effect.void,
invalidate: Effect.void,
});
const instanceRegistryLayer = Layer.succeed(
Expand All @@ -1051,6 +1055,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
)
: Effect.succeed(false),
streamChanges: Stream.empty,
listEnvironments: Effect.succeed(new Map()),
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), PubSub.subscribe),
},
);
Expand Down Expand Up @@ -1079,6 +1084,101 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
}),
);

it.effect("reads every instance's secret in one go before rebuilding any of them", () =>
Effect.gen(function* () {
const codexDriver = ProviderDriverKind.make("codex");
const claudeInstanceId = ProviderInstanceId.make("claude");
const codexInstanceId = ProviderInstanceId.make("codex");
const claudeReference = "op://Vault/claude/token";
const codexReference = "op://Vault/codex/token";
const unavailableProvider = (instanceId: ProviderInstanceId) =>
({
instanceId,
driver: codexDriver,
status: "error",
enabled: true,
installed: false,
auth: { status: "unknown" },
checkedAt: "2026-06-10T00:00:00.000Z",
version: null,
models: [],
slashCommands: [],
skills: [],
message: "Driver 'codex' failed to create instance: secret store is locked",
}) as const satisfies ServerProvider;

const primed = yield* Ref.make<ReadonlyArray<ReadonlyArray<string>>>([]);
const rebuiltIds = yield* Ref.make<ReadonlyArray<ProviderInstanceId>>([]);
const secretResolverLayer = Layer.succeed(ProviderSecretResolver, {
resolve: (environment) => Effect.succeed({ variables: environment, unresolved: [] }),
prime: (references) =>
Ref.update(primed, (previous) => [...previous, references]).pipe(Effect.asVoid),
invalidate: Effect.void,
});
const environmentFor = (reference: string) => [
{ name: "TOKEN", value: reference, sensitive: true },
];
const instanceRegistryLayer = Layer.succeed(
ProviderInstanceRegistry.ProviderInstanceRegistry,
{
getInstance: () => Effect.succeed(undefined),
listInstances: Effect.succeed([]),
listUnavailable: Effect.succeed([
unavailableProvider(claudeInstanceId),
unavailableProvider(codexInstanceId),
]),
listEnvironments: Effect.succeed(
new Map([
[claudeInstanceId, environmentFor(claudeReference)],
[codexInstanceId, environmentFor(codexReference)],
]),
),
rebuildInstanceWhen: (instanceId, shouldRebuild) =>
shouldRebuild({
driver: codexDriver,
environment: environmentFor(
instanceId === claudeInstanceId ? claudeReference : codexReference,
),
})
? Ref.update(rebuiltIds, (previous) => [...previous, instanceId]).pipe(
Effect.as(true),
)
: Effect.succeed(false),
streamChanges: Stream.empty,
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), PubSub.subscribe),
},
);

const scope = yield* Scope.make();
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void));
const runtimeServices = yield* Layer.build(
ProviderRegistryLive.pipe(
Layer.provideMerge(instanceRegistryLayer),
Layer.provideMerge(
ServerConfig.layerTest(process.cwd(), {
prefix: "t3-provider-registry-secret-prime-",
}),
),
Layer.provideMerge(NodeServices.layer),
Layer.provideMerge(secretResolverLayer),
),
).pipe(Scope.provide(scope));

yield* Effect.gen(function* () {
const registry = yield* ProviderRegistry.ProviderRegistry;
yield* registry.refresh();

// One call carrying both references. Priming per instance would be
// one unlock prompt per provider, which is the thing this exists to
// avoid, so the count matters as much as the contents.
const calls = yield* Ref.get(primed);
assert.strictEqual(calls.length, 1);
assert.deepStrictEqual(Array.from(calls[0] ?? []), [claudeReference, codexReference]);
assert.deepStrictEqual(yield* Ref.get(rebuiltIds), [claudeInstanceId, codexInstanceId]);
}).pipe(Effect.provide(runtimeServices));
}),
);

it.effect("persists the merged snapshot when a live update has empty models", () =>
Effect.gen(function* () {
const cursorDriver = ProviderDriverKind.make("cursor");
Expand Down Expand Up @@ -1145,6 +1245,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
listUnavailable: Effect.succeed([]),
rebuildInstanceWhen: () => Effect.succeed(false),
streamChanges: Stream.empty,
listEnvironments: Effect.succeed(new Map()),
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), (pubsub) =>
PubSub.subscribe(pubsub),
),
Expand Down Expand Up @@ -1276,6 +1377,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
listUnavailable: Effect.succeed([]),
rebuildInstanceWhen: () => Effect.succeed(false),
streamChanges: Stream.empty,
listEnvironments: Effect.succeed(new Map()),
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), (pubsub) =>
PubSub.subscribe(pubsub),
),
Expand Down Expand Up @@ -1385,6 +1487,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
listUnavailable: Effect.succeed([]),
rebuildInstanceWhen: () => Effect.succeed(false),
streamChanges: Stream.empty,
listEnvironments: Effect.succeed(new Map()),
subscribeChanges: Effect.flatMap(PubSub.unbounded<void>(), (pubsub) =>
PubSub.subscribe(pubsub),
),
Expand Down Expand Up @@ -1497,6 +1600,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
listUnavailable: Effect.succeed([]),
rebuildInstanceWhen: () => Effect.succeed(false),
streamChanges: Stream.fromPubSub(changes),
listEnvironments: Effect.succeed(new Map()),
subscribeChanges: PubSub.subscribe(changes),
},
);
Expand Down Expand Up @@ -1665,6 +1769,95 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
}),
);

it.effect("reads every boot instance's secret before building any of them", () =>
Effect.gen(function* () {
const claudeReference = "op://Vault/claude/token";
const codexReference = "op://Vault/codex/token";
// `streamChanges` carries later writes only, so the watcher never
// sees this snapshot. Boot is the run where nothing is cached yet
// and therefore the one that pays the most prompts without priming.
const serverSettings = yield* makeMutableServerSettingsService(
decodeServerSettings(
deepMerge(encodedDefaultServerSettings, {
providers: {
codex: { enabled: false },
claudeAgent: { enabled: false },
cursor: { enabled: false },
grok: { enabled: false },
opencode: { enabled: false },
},
providerInstances: {
claude_secret: {
driver: "claudeAgent",
displayName: "Claude Secret",
enabled: false,
environment: [
{ name: "CLAUDE_CODE_OAUTH_TOKEN", value: claudeReference, sensitive: true },
],
},
codex_secret: {
driver: "codex",
displayName: "Codex Secret",
enabled: false,
environment: [{ name: "TOKEN", value: codexReference, sensitive: true }],
},
} as unknown as ContractServerSettings["providerInstances"],
}),
),
);

const calls = yield* Ref.make<ReadonlyArray<string>>([]);
const recordingSecretResolverLayer = Layer.succeed(ProviderSecretResolver, {
resolve: (environment) =>
Ref.update(calls, (previous) => [...previous, "resolve"]).pipe(
Effect.as({ variables: environment, unresolved: [] }),
),
prime: (references) =>
Ref.update(calls, (previous) => [
...previous,
`prime:${Array.from(references).join(",")}`,
]).pipe(Effect.asVoid),
invalidate: Effect.void,
});

const scope = yield* Scope.make();
yield* Effect.addFinalizer(() => Scope.close(scope, Exit.void));
yield* Layer.build(
ProviderInstanceRegistryHydrationLive.pipe(
Layer.provideMerge(
Layer.succeed(ServerSettingsModule.ServerSettingsService, serverSettings),
),
Layer.provideMerge(
ServerConfig.layerTest(process.cwd(), {
prefix: "t3-provider-registry-boot-prime-",
}),
),
Layer.provideMerge(TestHttpClientLive),
Layer.provideMerge(
Layer.succeed(
ProviderEventLoggers.ProviderEventLoggers,
ProviderEventLoggers.NoOpProviderEventLoggers,
),
),
Layer.provideMerge(OpenCodeRuntime.OpenCodeRuntimeLive),
Layer.provideMerge(NodeServices.layer),
Layer.provideMerge(BackgroundPolicyAlwaysRunLayer),
Layer.provideMerge(recordingSecretResolverLayer),
),
).pipe(Scope.provide(scope));

// Both references in one call, and that call ahead of the first
// instance that would have read one on its own.
const recorded = yield* Ref.get(calls);
assert.strictEqual(
recorded[0],
`prime:${claudeReference},${codexReference}`,
`Expected boot to prime both references first; instead saw: ${recorded.join(" | ")}`,
);
assert.strictEqual(recorded.filter((entry) => entry.startsWith("prime")).length, 1);
}),
);

// Guards the second half of the reported bug: changing
// `providers.codex.binaryPath` in settings must tear down the live
// instance and rebuild it so a fresh probe runs with the new binary.
Expand Down
Loading
Loading