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
47 changes: 35 additions & 12 deletions apps/server/src/provider/Layers/ProviderRegistry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import {
ClaudeSettings,
CodexSettings,
DEFAULT_SERVER_SETTINGS,
EnvironmentId,
ProviderDriverKind,
ProviderInstanceId,
ServerSettings,
Expand Down Expand Up @@ -48,6 +49,7 @@ import {
ProviderRegistryLive,
} from "./ProviderRegistry.ts";
import * as ServerConfig from "../../config.ts";
import * as ServerEnvironment from "../../environment/ServerEnvironment.ts";
import * as ServerSettingsModule from "../../serverSettings.ts";
import {
readProviderStatusCache,
Expand All @@ -63,6 +65,16 @@ const decodeServerSettings = Schema.decodeSync(ServerSettings);
const encodeServerSettings = Schema.encodeSync(ServerSettings);
const encodedDefaultServerSettings = encodeServerSettings(DEFAULT_SERVER_SETTINGS);

const withProviderHost = <Provider extends { readonly instanceId: string }>(
provider: Provider,
) => ({
...provider,
providerHostInstance: {
environmentId: EnvironmentId.make("provider-registry-test"),
providerInstanceId: ProviderInstanceId.make(provider.instanceId),
},
});

const defaultClaudeSettings: ClaudeSettings = Schema.decodeSync(ClaudeSettings)({});
const defaultCodexSettings: CodexSettings = Schema.decodeSync(CodexSettings)({});
const decodeCodexSettings = Schema.decodeSync(CodexSettings);
Expand All @@ -77,11 +89,16 @@ process.env.T3CODE_CURSOR_ENABLED = "1";
const encoder = new TextEncoder();
const TEST_EPOCH = DateTime.makeUnsafe("1970-01-01T00:00:00.000Z");

const TestHttpClientLive = Layer.succeed(
HttpClient.HttpClient,
HttpClient.make((request) =>
Effect.succeed(HttpClientResponse.fromWeb(request, Response.json({ version: "0.0.0" }))),
const TestHttpClientLive = Layer.mergeAll(
Layer.succeed(
HttpClient.HttpClient,
HttpClient.make((request) =>
Effect.succeed(HttpClientResponse.fromWeb(request, Response.json({ version: "0.0.0" }))),
),
),
Layer.succeed(ServerEnvironment.ServerEnvironmentIdentity, {
getEnvironmentId: Effect.succeed(EnvironmentId.make("provider-registry-test")),
}),
);

const BackgroundPolicyAlwaysRunLayer = Layer.mock(BackgroundPolicy.BackgroundPolicy)({
Expand Down Expand Up @@ -1422,7 +1439,9 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
).pipe(Scope.provide(scope));
yield* Effect.gen(function* () {
const registry = yield* ProviderRegistry.ProviderRegistry;
assert.deepStrictEqual(yield* registry.getProviders, [initialProvider]);
assert.deepStrictEqual(yield* registry.getProviders, [
withProviderHost(initialProvider),
]);
assert.strictEqual(yield* Ref.get(refreshCalls), 0);
}).pipe(Effect.provide(runtimeServices));
}),
Expand Down Expand Up @@ -1753,7 +1772,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
);
assert.deepStrictEqual(
recoveredProviders.find((provider) => provider.instanceId === codexInstanceId),
codexProvider,
withProviderHost(codexProvider),
);

yield* Ref.set(catalogSnapshot, changedCatalogProvider);
Expand All @@ -1765,7 +1784,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
);
assert.deepStrictEqual(
changedProviders.find((provider) => provider.instanceId === codexInstanceId),
codexProvider,
withProviderHost(codexProvider),
);
}).pipe(Effect.provide(runtimeServices));

Expand Down Expand Up @@ -1880,7 +1899,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
const cachedProvider = yield* readProviderStatusCache(filePath);

assert.deepStrictEqual(cachedProvider, {
...refreshedProvider,
...withProviderHost(refreshedProvider),
models: [...initialProvider.models],
});
}).pipe(Effect.provide(runtimeServices));
Expand Down Expand Up @@ -2094,10 +2113,14 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te
yield* Effect.gen(function* () {
const registry = yield* ProviderRegistry.ProviderRegistry;

assert.deepStrictEqual(yield* registry.getProviders, [cachedProvider]);
assert.deepStrictEqual(yield* registry.refresh(codexDriver), [cachedProvider]);
assert.deepStrictEqual(yield* registry.getProviders, [
withProviderHost(cachedProvider),
]);
assert.deepStrictEqual(yield* registry.refresh(codexDriver), [
withProviderHost(cachedProvider),
]);
assert.deepStrictEqual(yield* registry.refreshInstance(codexInstanceId), [
cachedProvider,
withProviderHost(cachedProvider),
]);
}).pipe(Effect.provide(runtimeServices));
}),
Expand Down Expand Up @@ -2205,7 +2228,7 @@ it.layer(Layer.mergeAll(NodeServices.layer, ServerSettingsModule.layerTest(), Te

yield* Effect.gen(function* () {
const registry = yield* ProviderRegistry.ProviderRegistry;
assert.deepStrictEqual(yield* registry.getProviders, [codexProvider]);
assert.deepStrictEqual(yield* registry.getProviders, [withProviderHost(codexProvider)]);

yield* Ref.set(failNextList, true);
yield* PubSub.publish(changes, undefined);
Expand Down
52 changes: 42 additions & 10 deletions apps/server/src/provider/Layers/ProviderRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@
*/
import {
defaultInstanceIdForDriver,
type EnvironmentId,
makeProviderHostInstanceIdentity,
ProviderDriverKind,
type ProviderInstanceId,
type ServerProvider,
Expand All @@ -41,6 +43,7 @@ import * as Stream from "effect/Stream";
import * as Semaphore from "effect/Semaphore";

import { ServerConfig } from "../../config.ts";
import * as ServerEnvironment from "../../environment/ServerEnvironment.ts";
import { ProviderInstanceRegistry } from "../Services/ProviderInstanceRegistry.ts";
import { ProviderRegistry, type ProviderRegistryShape } from "../Services/ProviderRegistry.ts";
import {
Expand All @@ -57,12 +60,15 @@ import type { ProviderSnapshotSource } from "../builtInProviderCatalog.ts";

const loadProviders = (
providerSources: ReadonlyArray<ProviderSnapshotSource>,
environmentId: EnvironmentId,
): Effect.Effect<ReadonlyArray<ServerProvider>> =>
Effect.forEach(
providerSources,
(providerSource) =>
providerSource.getSnapshot.pipe(
Effect.flatMap((snapshot) => correlateSnapshotWithSource(providerSource, snapshot)),
Effect.flatMap((snapshot) =>
correlateSnapshotWithSource(providerSource, snapshot, environmentId),
),
),
{
concurrency: "unbounded",
Expand Down Expand Up @@ -231,9 +237,22 @@ const haveProvidersChanged = (
nextProviders: ReadonlyArray<ServerProvider>,
): boolean => !Equal.equals(previousProviders, nextProviders);

/** Stamp the persisted environment and configured instance at the registry boundary. */
export const stampProviderHostInstance = (
provider: ServerProvider,
environmentId: EnvironmentId,
): ServerProvider => ({
...provider,
providerHostInstance: makeProviderHostInstanceIdentity({
environmentId,
providerInstanceId: provider.instanceId,
}),
});

const correlateSnapshotWithSource = (
source: ProviderSnapshotSource,
snapshot: ServerProvider,
environmentId: EnvironmentId,
): Effect.Effect<ServerProvider> => {
if (snapshot.instanceId !== source.instanceId) {
return Effect.die(
Expand All @@ -249,7 +268,10 @@ const correlateSnapshotWithSource = (
),
);
}
return Effect.succeed(snapshot);
// Provider drivers only know their configured instance. The registry owns
// the persisted environment identity and stamps the pair centrally so
// unavailable, cached, and live snapshots use the same host contract.
return Effect.succeed(stampProviderHostInstance(snapshot, environmentId));
};

/**
Expand Down Expand Up @@ -278,6 +300,9 @@ export const ProviderRegistryLive = Layer.effect(
ProviderRegistry,
Effect.gen(function* () {
const instanceRegistry = yield* ProviderInstanceRegistry;
const environmentId = yield* ServerEnvironment.ServerEnvironmentIdentity.pipe(
Effect.flatMap((identity) => identity.getEnvironmentId),
);
const config = yield* ServerConfig;
const fileSystem = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
Expand All @@ -296,7 +321,7 @@ export const ProviderRegistryLive = Layer.effect(
// below.
const bootInstances = yield* instanceRegistry.listInstances;
const bootSources = bootInstances.map(buildSnapshotSource);
const fallbackProviders = yield* loadProviders(bootSources);
const fallbackProviders = yield* loadProviders(bootSources, environmentId);
const fallbackByInstance = new Map<ProviderInstanceId, ServerProvider>();
for (let index = 0; index < fallbackProviders.length; index++) {
const provider = fallbackProviders[index];
Expand Down Expand Up @@ -526,7 +551,7 @@ export const ProviderRegistryLive = Layer.effect(
) {
return yield* providerSource.refresh.pipe(
Effect.flatMap((nextProvider) =>
correlateSnapshotWithSource(providerSource, nextProvider).pipe(
correlateSnapshotWithSource(providerSource, nextProvider, environmentId).pipe(
Effect.flatMap(syncProvider),
),
),
Expand Down Expand Up @@ -669,7 +694,9 @@ export const ProviderRegistryLive = Layer.effect(
for (const [, instance] of newlyAdded) {
const source = buildSnapshotSource(instance);
yield* Stream.runForEach(source.streamChanges, (provider) =>
correlateSnapshotWithSource(source, provider).pipe(Effect.flatMap(syncProvider)),
correlateSnapshotWithSource(source, provider, environmentId).pipe(
Effect.flatMap(syncProvider),
),
).pipe(Effect.forkScoped);
}
yield* Effect.yieldNow;
Expand All @@ -684,16 +711,21 @@ export const ProviderRegistryLive = Layer.effect(
Effect.gen(function* () {
const source = buildSnapshotSource(instance);
const provider = yield* source.getSnapshot;
yield* correlateSnapshotWithSource(source, provider).pipe(
yield* correlateSnapshotWithSource(source, provider, environmentId).pipe(
Effect.flatMap(syncProvider),
);
}).pipe(Effect.ignoreCause({ log: true })),
{ concurrency: "unbounded", discard: true },
);
yield* upsertProviders(unavailableProviders, {
persist: false,
replace: true,
});
yield* upsertProviders(
unavailableProviders.map((provider) =>
stampProviderHostInstance(provider, environmentId),
),
{
persist: false,
replace: true,
},
);

const nextSubs = new Map(carriedOver);
for (const [instanceId, instance] of newlyAdded) {
Expand Down
20 changes: 20 additions & 0 deletions apps/server/src/provider/Layers/ProviderRegistryIdentity.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
import { assert, describe, it } from "@effect/vitest";
import { EnvironmentId, ProviderInstanceId, type ServerProvider } from "@t3tools/contracts";

import { stampProviderHostInstance } from "./ProviderRegistry.ts";

describe("provider host identity stamping", () => {
it("uses the persisted environment plus instance and never continuation metadata", () => {
const provider = {
instanceId: ProviderInstanceId.make("codex_work"),
continuation: { groupKey: "codex:home:/shared" },
} as ServerProvider;
assert.deepStrictEqual(stampProviderHostInstance(provider, EnvironmentId.make("env-a")), {
...provider,
providerHostInstance: {
environmentId: EnvironmentId.make("env-a"),
providerInstanceId: ProviderInstanceId.make("codex_work"),
},
});
});
});
1 change: 1 addition & 0 deletions packages/contracts/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ export * from "./ipc.ts";
export * from "./terminal.ts";
export * from "./provider.ts";
export * from "./providerInstance.ts";
export * from "./providerIdentity.ts";
export * from "./providerSetup.ts";
export * from "./providerRuntime.ts";
export * from "./providerUsageLimits.ts";
Expand Down
105 changes: 105 additions & 0 deletions packages/contracts/src/providerIdentity.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,105 @@
import * as Schema from "effect/Schema";
import { describe, expect, it } from "vite-plus/test";
import {
compareManagedProviderLaunchIdentity,
ManagedProviderLaunchIdentity,
qualifyManagedProviderLaunch,
} from "./providerIdentity.ts";
import { ServerProvider } from "./server.ts";
import { EnvironmentId } from "./baseSchemas.ts";
import { ProviderInstanceId } from "./providerInstance.ts";

const decodeServerProvider = Schema.decodeUnknownSync(ServerProvider);
const decodeManagedProviderLaunchIdentity = Schema.decodeUnknownSync(ManagedProviderLaunchIdentity);
const host = {
environmentId: EnvironmentId.make("env-a"),
providerInstanceId: ProviderInstanceId.make("codex"),
};
const keyIds = [{ namespace: "provider.example", id: "key-a" }] as const;
const identity = {
subject: "subject-a",
keyIds: [...keyIds],
providerHostInstance: host,
} satisfies typeof ManagedProviderLaunchIdentity.Type;

describe("managed provider launch identity", () => {
it("keeps old provider snapshots compatible when optional identity fields are absent", () => {
const provider = decodeServerProvider({
instanceId: "codex",
driver: "codex",
enabled: true,
installed: true,
version: "1.0.0",
status: "ready",
auth: { status: "authenticated" },
checkedAt: "2026-01-01T00:00:00.000Z",
models: [],
});
expect(provider.auth.keyIds).toBeUndefined();
expect(provider.providerHostInstance).toBeUndefined();
});

it("fails closed when any managed-launch component is unavailable", () => {
expect(qualifyManagedProviderLaunch({ keyIds, providerHostInstance: host })).toEqual({
qualified: false,
reason: "missing-subject",
});
expect(
qualifyManagedProviderLaunch({ subject: "subject-a", providerHostInstance: host }),
).toEqual({
qualified: false,
reason: "missing-key-id",
});
expect(qualifyManagedProviderLaunch({ subject: "subject-a", keyIds })).toEqual({
qualified: false,
reason: "missing-host-identity",
});
expect(() =>
decodeManagedProviderLaunchIdentity({
subject: "subject-a",
keyIds: [],
providerHostInstance: host,
}),
).toThrow();
});

it("compares namespace-qualified key sets and reports exact drift", () => {
expect(compareManagedProviderLaunchIdentity(identity, identity)).toBe("match");
expect(
compareManagedProviderLaunchIdentity(identity, { ...identity, subject: "subject-b" }),
).toBe("subject-drift");
expect(
compareManagedProviderLaunchIdentity(identity, {
...identity,
keyIds: [{ namespace: "provider.example", id: "key-b" }],
}),
).toBe("key-id-drift");
expect(
compareManagedProviderLaunchIdentity(identity, {
...identity,
providerHostInstance: { ...host, environmentId: EnvironmentId.make("env-b") },
}),
).toBe("host-drift");
expect(compareManagedProviderLaunchIdentity(identity, undefined)).toBe("unknown-identity");
});

it("does not collide on embedded NULs and normalizes duplicate IDs", () => {
const first = {
...identity,
keyIds: [
{ namespace: "a\u0000b", id: "c" },
{ namespace: "a\u0000b", id: "c" },
],
};
const second = {
...identity,
keyIds: [{ namespace: "a", id: "b\u0000c" }],
};
expect(compareManagedProviderLaunchIdentity(first, second)).toBe("key-id-drift");
const qualification = qualifyManagedProviderLaunch(first);
expect(qualification.qualified).toBe(true);
if (qualification.qualified) {
expect(qualification.identity.keyIds).toHaveLength(1);
}
});
});
Loading
Loading