Skip to content
Open
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
25 changes: 25 additions & 0 deletions apps/server/src/provider/Drivers/CodexDriver.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import { expect, it } from "@effect/vitest";
import { EnvironmentId, ProviderInstanceId, ProviderSessionId, ThreadId } from "@t3tools/contracts";
import { HostProcessPlatform } from "@t3tools/shared/hostProcess";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Deferred from "effect/Deferred";
import * as Fiber from "effect/Fiber";
import * as FileSystem from "effect/FileSystem";
Expand Down Expand Up @@ -139,9 +140,17 @@ it.layer(layerTest)("CodexDriver", (it) => {
// Interrupt the two startup checks that way (one is the sign-in listener's);
// the disconnect below must still refresh.
let interruptedChecks = 0;
let failAcquire = false;
const startupChecksInterrupted = yield* Deferred.make<void>();
const acquire = () =>
Effect.suspend(() => {
if (failAcquire)
return Effect.fail(
new CodexInstallation.CodexInstallationError({
operation: "acquire",
detail: "The fixture runtime is unavailable",
}),
);
if (interruptedChecks === 2) return Effect.succeed(executable);
interruptedChecks += 1;
return (
Expand Down Expand Up @@ -206,6 +215,18 @@ it.layer(layerTest)("CodexDriver", (it) => {
expect(restored.runtimePaths?.shadowHomePath).toContain(
`providers/codex/${instanceId}/shadow`,
);
// The fixture app-server cannot start, so the workspace scan fails
// instead of passing the machine snapshot off as the workspace's skills.
const scan = yield* instance.snapshotForCwd!(serverConfig.stateDir).pipe(Effect.exit);
expect(Exit.isFailure(scan)).toBe(true);
// The account is still signed in, so a runtime that cannot start is a
// failed scan too, not the machine snapshot's inventory.
failAcquire = true;
const runtimeScan = yield* instance.snapshotForCwd!(serverConfig.stateDir).pipe(
Effect.exit,
);
failAcquire = false;
expect(Exit.isFailure(runtimeScan)).toBe(true);
yield* Deferred.await(observedAccount);
// Sessions launch the T3-installed Codex with the account's token, not ambient credentials.
const threadId = ThreadId.make("managed-account-thread");
Expand Down Expand Up @@ -237,6 +258,10 @@ it.layer(layerTest)("CodexDriver", (it) => {
expect(after.installed).toBe(true);
expect(after.models).toEqual([]);
expect(Option.isNone(yield* store.get)).toBe(true);
// Signed out, there is no runtime to scan with, so the machine
// snapshot's empty inventory stands instead of a failed scan.
const signedOutScan = yield* instance.snapshotForCwd!(serverConfig.stateDir);
expect(signedOutScan.skills).toEqual([]);
}).pipe(
Effect.provideService(ServerSecretStore.ServerSecretStore, secrets),
Effect.provideService(
Expand Down
49 changes: 36 additions & 13 deletions apps/server/src/provider/Drivers/CodexManagedProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -275,21 +275,44 @@ export const makeManagedCodexProvider = Effect.fn("makeManagedCodexProvider")(fu
snapshotForCwd: (cwd: string) =>
enabled
? resolveRuntime.pipe(
Effect.flatMap((effective) =>
probeCodexSkillsForCwd({
binaryPath: effective.config.binaryPath,
homePath: effective.config.homePath,
launchArgs: effective.config.launchArgs,
cwd,
environment: effective.environment,
}),
),
Effect.flatMap((skills) =>
snapshot.getSnapshot.pipe(Effect.map((draft) => ({ ...draft, skills }))),
),
Effect.matchEffect({
// Signed out or not set up, there is nothing to scan with and the
// machine snapshot's empty inventory says as much. Otherwise the
// runtime failed for a signed-in account (a token refresh or
// reconnect) or before the first check, and the scan failed.
onFailure: (cause) =>
snapshot.getSnapshot.pipe(
Effect.filterOrFail(
(machine) => machine.auth.status === "unauthenticated",
() => cause,
),
),
onSuccess: (effective) =>
probeCodexSkillsForCwd({
binaryPath: effective.config.binaryPath,
homePath: effective.config.homePath,
launchArgs: effective.config.launchArgs,
cwd,
environment: effective.environment,
}).pipe(
Effect.flatMap((skills) =>
snapshot.getSnapshot.pipe(Effect.map((draft) => ({ ...draft, skills }))),
),
),
}),
Effect.scoped,
Effect.provideService(ChildProcessSpawner.ChildProcessSpawner, spawner),
Effect.catch(() => snapshot.getSnapshot),
// Fail rather than return the machine snapshot: the registry
// publishes any returned snapshot as this workspace's inventory.
Effect.mapError(
(cause) =>
new ProviderDriverError({
driver: DRIVER,
instanceId,
detail: `Failed to probe Codex skills for '${cwd}'`,
cause,
}),
),
)
: snapshot.getSnapshot,
} satisfies ProviderInstance;
Expand Down
151 changes: 142 additions & 9 deletions apps/server/src/provider/ProviderRegistry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ import type {
} from "@t3tools/provider-core/server/driver";
import * as ProviderInstanceRegistry from "./ProviderInstanceRegistry.ts";
import * as ProviderRegistry from "./ProviderRegistry.ts";
import { ProviderDriverError } from "./Errors.ts";
import { makeManualOnlyProviderMaintenanceCapabilities } from "@t3tools/provider-core/server/maintenanceResolver";
const decodeServerSettings = Schema.decodeSync(ServerSettings);
const encodeServerSettings = Schema.encodeSync(ServerSettings);
Expand Down Expand Up @@ -1661,12 +1662,6 @@ it.layer(
slashCommands: [{ name: "project" }],
skills: [{ name: "project", path: "/workspace/SKILL.md", enabled: true }],
} as const satisfies ServerProvider;
const pendingScopedProvider = {
...scopedProvider,
status: "error",
installed: false,
slashCommands: [],
} as const satisfies ServerProvider;
const snapshotCalls = yield* Ref.make(0);
const scopedResult = yield* Ref.make<ProviderWorkspaceSnapshot>({
...scopedProvider,
Expand All @@ -1679,7 +1674,7 @@ it.layer(
readonly started: Deferred.Deferred<void>;
readonly release: Deferred.Deferred<void>;
} | null>(null);
const returnPendingSnapshot = yield* Ref.make(true);
const failScan = yield* Ref.make(true);
const probeStarted = yield* Deferred.make<void>();
const releaseProbe = yield* Deferred.make<void>();
const makeInstance = (
Expand Down Expand Up @@ -1715,7 +1710,13 @@ it.layer(
const firstInstance = makeInstance(machineProvider, () =>
Effect.gen(function* () {
yield* Ref.update(snapshotCalls, (count) => count + 1);
if (yield* Ref.get(returnPendingSnapshot)) return pendingScopedProvider;
if (yield* Ref.get(failScan)) {
return yield* new ProviderDriverError({
driver,
instanceId,
detail: "The workspace scan failed.",
});
}
yield* Deferred.succeed(probeStarted, undefined);
yield* Deferred.await(releaseProbe);
const result = yield* Ref.get(scopedResult);
Expand Down Expand Up @@ -1772,7 +1773,7 @@ it.layer(
const registry = yield* ProviderRegistry.ProviderRegistry;
yield* registry.refreshWorkspaceSnapshot({ instanceId, cwd: "/workspace" });
assert.strictEqual((yield* registry.getProviders)[0]?.workspaceSnapshots, undefined);
yield* Ref.set(returnPendingSnapshot, false);
yield* Ref.set(failScan, false);
const workspaceUpdate = yield* registry.streamChanges.pipe(
Stream.runHead,
Effect.forkChild,
Expand Down Expand Up @@ -1887,6 +1888,138 @@ it.layer(
}),
);

it.effect("publishes workspace skills while machine health is in error", () =>
Effect.gen(function* () {
const driver = ProviderDriverKind.make("codex");
const instanceId = ProviderInstanceId.make("codex");
const healthyProvider = {
instanceId,
driver,
status: "ready",
enabled: true,
installed: true,
auth: { status: "authenticated" },
checkedAt: "2026-06-10T00:00:00.000Z",
version: "1.0.0",
models: [],
slashCommands: [COMPACT_SLASH_COMMAND],
skills: [],
} as const satisfies ServerProvider;
// Codex's account check failed; its skill scan did not.
const unhealthyProvider = {
...healthyProvider,
status: "error",
auth: { status: "unknown" },
message: "Codex app-server provider probe failed: unauthorized (401).",
slashCommands: [],
} as const satisfies ServerProvider;
const firstSkills = [{ name: "first", path: "/workspace/first/SKILL.md", enabled: true }];
const laterSkills = [{ name: "later", path: "/workspace/later/SKILL.md", enabled: true }];
const scanResult = yield* Ref.make<ProviderWorkspaceSnapshot>({
...healthyProvider,
skills: firstSkills,
});
const instance: ProviderInstance = {
instanceId,
driverKind: driver,
continuationIdentity: { driverKind: driver, continuationKey: "codex:instance:codex" },
displayName: undefined,
enabled: true,
snapshot: {
resolveMaintenance: () =>
Effect.succeed(
makeManualOnlyProviderMaintenanceCapabilities({
provider: driver,
packageName: null,
}),
),
getSnapshot: Effect.succeed(unhealthyProvider),
refresh: Effect.succeed(unhealthyProvider),
streamChanges: Stream.empty,
applyUsageLimits: () => Effect.void,
},
snapshotForCwd: () => Ref.get(scanResult),
orchestrationAdapter: {} as ProviderInstance["orchestrationAdapter"],
textGeneration: {} as ProviderInstance["textGeneration"],
};
const layerInstanceRegistry = Layer.succeed(
ProviderInstanceRegistry.ProviderInstanceRegistry,
{
getInstance: (requestedId) =>
Effect.succeed(requestedId === instanceId ? instance : undefined),
listInstances: Effect.succeed([instance]),
listUnavailable: Effect.succeed([]),
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(
ProviderRegistry.layer.pipe(
Layer.provideMerge(layerInstanceRegistry),
Layer.provideMerge(
ServerConfig.layerTest(process.cwd(), {
prefix: "t3-provider-registry-unhealthy-workspace-",
}),
),
Layer.provideMerge(NodeServices.layer),
),
).pipe(Scope.provide(scope));

yield* Effect.gen(function* () {
const registry = yield* ProviderRegistry.ProviderRegistry;
const workspaceOf = Effect.map(
registry.getProviders,
(providers) => providers[0]?.workspaceSnapshots?.[0],
);
yield* registry.refreshWorkspaceSnapshot({ instanceId, cwd: "/workspace" });
assert.deepStrictEqual((yield* workspaceOf)?.slashCommands, [COMPACT_SLASH_COMMAND]);

// An expired scan runs while the account check fails: the skills are
// current, but commands that ride on the health check are unknown.
yield* TestClock.adjust(PROVIDER_WORKSPACE_SNAPSHOT_TTL_MS);
yield* Ref.set(scanResult, { ...unhealthyProvider, skills: laterSkills });
yield* registry.refreshWorkspaceSnapshot({ instanceId, cwd: "/workspace" });
const unhealthy = yield* workspaceOf;
assert.deepStrictEqual(unhealthy?.skills, laterSkills);
assert.deepStrictEqual(unhealthy?.slashCommands, [COMPACT_SLASH_COMMAND]);
assert.strictEqual(unhealthy?.slashCommandsPending, true);
const provider = (yield* registry.getProviders)[0];
assert.strictEqual(provider?.status, "error");
assert.strictEqual(provider?.message, unhealthyProvider.message);

// A pending entry is rescanned on the next request; an empty scan
// removes the skills.
yield* Ref.set(scanResult, { ...unhealthyProvider, skills: [] });
yield* registry.refreshWorkspaceSnapshot({ instanceId, cwd: "/workspace" });
assert.deepStrictEqual((yield* workspaceOf)?.skills, []);

// A driver that discovers its own commands reports them even while
// health is in error.
const discoveredCommand = { name: "review", description: "Review changes" };
yield* Ref.set(scanResult, {
...unhealthyProvider,
slashCommands: [discoveredCommand],
slashCommandsPending: false,
});
yield* registry.refreshWorkspaceSnapshot({ instanceId, cwd: "/workspace" });
const discovered = yield* workspaceOf;
assert.deepStrictEqual(discovered?.slashCommands, [discoveredCommand]);
assert.strictEqual(discovered?.slashCommandsPending, undefined);

// Once health recovers, the next scan restores the full entry.
yield* TestClock.adjust(PROVIDER_WORKSPACE_SNAPSHOT_TTL_MS);
yield* Ref.set(scanResult, { ...healthyProvider, skills: laterSkills });
yield* registry.refreshWorkspaceSnapshot({ instanceId, cwd: "/workspace" });
const recovered = yield* workspaceOf;
assert.deepStrictEqual(recovered?.skills, laterSkills);
assert.deepStrictEqual(recovered?.slashCommands, [COMPACT_SLASH_COMMAND]);
assert.strictEqual(recovered?.slashCommandsPending, undefined);
}).pipe(Effect.provide(runtimeServices));
}),
);

it.effect("refreshes OpenCode catalogs and preserves other providers", () =>
Effect.gen(function* () {
const codexDriver = ProviderDriverKind.make("codex");
Expand Down
45 changes: 25 additions & 20 deletions apps/server/src/provider/ProviderRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1092,27 +1092,32 @@ export const layer = Layer.effect(
: Effect.void;
return yield* refreshMachineSnapshot.pipe(
Effect.andThen(instance.snapshotForCwd(input.cwd)),
// A failed scan fails this effect, so a returned snapshot is a real
// inventory even while the provider's health check reports an error.
Effect.flatMap((scopedSnapshot) =>
scopedSnapshot.status === "error" && scopedSnapshot.slashCommandsPending === undefined
? Ref.get(providersRef)
: instanceRegistry.getInstance(input.instanceId).pipe(
Effect.flatMap((currentInstance) => {
if (currentInstance !== instance) return Ref.get(providersRef);
// Write only if the cwd's snapshot did not change during the
// scan. A session event or another scan that landed first is newer.
return updateProviders((currentProviders) =>
currentProviders.map((candidate) =>
candidate.instanceId === input.instanceId &&
Equal.equals(workspaceSnapshotOf(candidate), scannedFrom)
? upsertProviderWorkspaceSnapshot(candidate, input.cwd, {
...scopedSnapshot,
checkedAt: scannedAt,
})
: candidate,
),
);
}),
),
instanceRegistry.getInstance(input.instanceId).pipe(
Effect.flatMap((currentInstance) => {
if (currentInstance !== instance) return Ref.get(providersRef);
// Write only if the cwd's snapshot did not change during the
// scan. A session event or another scan that landed first is newer.
return updateProviders((currentProviders) =>
currentProviders.map((candidate) =>
candidate.instanceId === input.instanceId &&
Equal.equals(workspaceSnapshotOf(candidate), scannedFrom)
? upsertProviderWorkspaceSnapshot(candidate, input.cwd, {
...scopedSnapshot,
checkedAt: scannedAt,
// Unless the driver reports its own command discovery,
// commands come from the health check. A failed check
// leaves them unknown, so keep the last ones and retry.
slashCommandsPending:
scopedSnapshot.slashCommandsPending ?? scopedSnapshot.status === "error",
})
Comment thread
coderabbitai[bot] marked this conversation as resolved.
: candidate,
),
);
}),
),
),
Effect.ensuring(
claimed
Expand Down
Loading
Loading