diff --git a/apps/server/src/provider/Drivers/CodexDriver.test.ts b/apps/server/src/provider/Drivers/CodexDriver.test.ts index 72340b7471cf..7a816ba3d3f5 100644 --- a/apps/server/src/provider/Drivers/CodexDriver.test.ts +++ b/apps/server/src/provider/Drivers/CodexDriver.test.ts @@ -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"; @@ -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(); 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 ( @@ -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"); @@ -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( diff --git a/apps/server/src/provider/Drivers/CodexManagedProvider.ts b/apps/server/src/provider/Drivers/CodexManagedProvider.ts index 32d88e644879..9d4ed6fe22a3 100644 --- a/apps/server/src/provider/Drivers/CodexManagedProvider.ts +++ b/apps/server/src/provider/Drivers/CodexManagedProvider.ts @@ -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; diff --git a/apps/server/src/provider/ProviderRegistry.test.ts b/apps/server/src/provider/ProviderRegistry.test.ts index 8653bd32ebd2..d2dc8d255b55 100644 --- a/apps/server/src/provider/ProviderRegistry.test.ts +++ b/apps/server/src/provider/ProviderRegistry.test.ts @@ -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); @@ -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({ ...scopedProvider, @@ -1679,7 +1674,7 @@ it.layer( readonly started: Deferred.Deferred; readonly release: Deferred.Deferred; } | null>(null); - const returnPendingSnapshot = yield* Ref.make(true); + const failScan = yield* Ref.make(true); const probeStarted = yield* Deferred.make(); const releaseProbe = yield* Deferred.make(); const makeInstance = ( @@ -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); @@ -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, @@ -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({ + ...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(), 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"); diff --git a/apps/server/src/provider/ProviderRegistry.ts b/apps/server/src/provider/ProviderRegistry.ts index 5f5120ee28bc..54da7d3d651a 100644 --- a/apps/server/src/provider/ProviderRegistry.ts +++ b/apps/server/src/provider/ProviderRegistry.ts @@ -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", + }) + : candidate, + ), + ); + }), + ), ), Effect.ensuring( claimed diff --git a/packages/provider-opencode/src/server/driver.test.ts b/packages/provider-opencode/src/server/driver.test.ts index 10410d3abab9..7389f8f1b09a 100644 --- a/packages/provider-opencode/src/server/driver.test.ts +++ b/packages/provider-opencode/src/server/driver.test.ts @@ -20,6 +20,7 @@ import * as ProviderEventLoggers from "@t3tools/provider-core/server/ProviderEve import * as ProviderMaintenance from "@t3tools/provider-core/server/maintenanceResolver"; import * as OpenCodeRuntime from "./OpenCodeRuntime.ts"; import { + OPENCODE_1_RESPONSES, OPENCODE_2_RESPONSES, OPENCODE_2_WORKSPACE_RESPONSES, replayOpenCodeServer, @@ -112,6 +113,36 @@ it.layer(layer)("OpenCodeDriver runtime selection", (it) => { }).pipe(Effect.scoped), ); + it.effect("keeps an OpenCode 1.x workspace's commands pending when listing them fails", () => + Effect.gen(function* () { + const openCode1Runtime = { + connectToOpenCodeServer: () => + Effect.succeed({ + url: "http://127.0.0.1:4096", + version: "1.18.32", + exitCode: null, + external: true, + }), + createOpenCodeSdkClient: () => ({ + command: { list: () => Promise.reject(new Error("command.list failed")) }, + }), + loadOpenCodeSkills: () => + Effect.succeed([{ name: "plum", location: "/work/.opencode/skills/plum/SKILL.md" }]), + } as unknown as OpenCodeRuntime.OpenCodeRuntimeShape; + const instance = yield* create( + { serverUrl: "http://127.0.0.1:4096", serverPassword: "secret" }, + replayOpenCodeServer(OPENCODE_1_RESPONSES, "secret"), + ).pipe(Effect.provideService(OpenCodeRuntime.OpenCodeRuntime, openCode1Runtime)); + + const workspace = yield* instance.snapshotForCwd!("/work"); + assert.deepStrictEqual( + workspace.skills.map((skill) => skill.name), + ["plum"], + ); + assert.strictEqual(workspace.slashCommandsPending, true); + }).pipe(Effect.scoped), + ); + it.effect("answers capability reads without waiting on an unreachable server", () => Effect.gen(function* () { const hang = HttpClient.make(() => Effect.never); diff --git a/packages/provider-opencode/src/server/driver.ts b/packages/provider-opencode/src/server/driver.ts index 9c37ab2abd48..2f3f1017d932 100644 --- a/packages/provider-opencode/src/server/driver.ts +++ b/packages/provider-opencode/src/server/driver.ts @@ -405,10 +405,18 @@ export const OpenCodeDriver: ProviderDriver skills: openCodeRuntime.loadOpenCodeSkills(client), commands: OpenCodeRuntime.loadOpenCodeCommands(client).pipe( Effect.timeout("10 seconds"), - Effect.orElseSucceed(() => []), + Effect.option, ), }, { concurrency: "unbounded" }, + ).pipe( + // A failed command lookup leaves the commands pending, so the + // registry keeps the last known ones and rescans. + Effect.map(({ skills, commands }) => ({ + skills, + commands: Option.getOrElse(commands, () => []), + slashCommandsPending: Option.isNone(commands), + })), ); const loadWorkspaceForCwd = (cwd: string) => effectiveConfig.serverUrl.trim().length > 0 @@ -511,16 +519,18 @@ export const OpenCodeDriver: ProviderDriver ...machineSnapshot, skills: openCode2SkillsToServerProviderSkills(skills), slashCommands: openCode2CommandsToServerProviderSlashCommands(commands), + slashCommandsPending: false, })), ), v1: Effect.all([ snapshot.getSnapshot, loadWorkspaceForCwd(cwd).pipe(Effect.timeout("20 seconds")), ]).pipe( - Effect.map(([machineSnapshot, { skills, commands }]) => ({ + Effect.map(([machineSnapshot, { skills, commands, slashCommandsPending }]) => ({ ...machineSnapshot, skills: openCodeSkillsToServerProviderSkills(skills), slashCommands: openCodeCommandsToServerProviderSlashCommands(commands), + slashCommandsPending, })), Effect.mapError( (cause) => diff --git a/packages/provider-pi/src/server/driver.ts b/packages/provider-pi/src/server/driver.ts index 25d390911781..a8e9d383a70e 100644 --- a/packages/provider-pi/src/server/driver.ts +++ b/packages/provider-pi/src/server/driver.ts @@ -201,7 +201,11 @@ export const PiDriver: ProviderDriver = { ), ), ]).pipe( - Effect.map(([machineSnapshot, commands]) => ({ ...machineSnapshot, ...commands })), + Effect.map(([machineSnapshot, commands]) => ({ + ...machineSnapshot, + ...commands, + slashCommandsPending: false, + })), ), orchestrationAdapter, textGeneration,