diff --git a/apps/server/integration/lifecycle.integration.test.ts b/apps/server/integration/lifecycle.integration.test.ts new file mode 100644 index 000000000000..e17fb375965a --- /dev/null +++ b/apps/server/integration/lifecycle.integration.test.ts @@ -0,0 +1,375 @@ +import type { ProviderRuntimeEvent } from "@t3tools/contracts"; +import { ThreadId } from "@t3tools/contracts"; +import { DEFAULT_SERVER_SETTINGS } from "@t3tools/contracts/settings"; +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { it, assert } from "@effect/vitest"; +import { Effect, FileSystem, Layer, Path, Queue, Stream } from "effect"; + +import { ProviderUnsupportedError } from "../src/provider/Errors.ts"; +import { ProviderAdapterRegistry } from "../src/provider/Services/ProviderAdapterRegistry.ts"; +import { McpConfigServiceLive } from "../src/provider/Layers/McpConfig.ts"; +import { McpConfigService } from "../src/provider/Services/McpConfig.ts"; +import { ProviderSessionDirectoryLive } from "../src/provider/Layers/ProviderSessionDirectory.ts"; +import { makeProviderServiceLive } from "../src/provider/Layers/ProviderService.ts"; +import { ProviderService } from "../src/provider/Services/ProviderService.ts"; +import { ServerConfig } from "../src/config.ts"; +import { ServerSettingsService } from "../src/serverSettings.ts"; +import { AnalyticsService } from "../src/telemetry/Services/AnalyticsService.ts"; +import { SqlitePersistenceMemory } from "../src/persistence/Layers/Sqlite.ts"; +import { ProviderSessionRuntimeRepositoryLive } from "../src/persistence/Layers/ProviderSessionRuntime.ts"; +import { ProviderSessionRuntimeRepository } from "../src/persistence/Services/ProviderSessionRuntime.ts"; + +import { makeTestProviderAdapterHarness } from "./TestProviderAdapter.integration.ts"; +import { codexTurnTextFixture } from "./fixtures/providerRuntime.ts"; + +// ── Helpers ────────────────────────────────────────────────────── + +const makeWorkspaceDirectory = Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const pathService = yield* Path.Path; + const cwd = yield* fs.makeTempDirectory(); + yield* fs.writeFileString(pathService.join(cwd, "README.md"), "v1\n"); + return cwd; +}).pipe(Effect.provide(NodeServices.layer)); + +/** + * Lifecycle fixture uses the **real** McpConfigServiceLive (not a mock) + * so that config resolution, snapshot persistence, and rehydration are + * exercised through the same code paths as production. + */ +const makeLifecycleFixture = Effect.gen(function* () { + const cwd = yield* makeWorkspaceDirectory; + const harness = yield* makeTestProviderAdapterHarness(); + + const registry: typeof ProviderAdapterRegistry.Service = { + getByProvider: (provider) => + provider === "codex" + ? Effect.succeed(harness.adapter) + : Effect.fail(new ProviderUnsupportedError({ provider })), + listProviders: () => Effect.succeed(["codex"]), + }; + + const runtimeRepositoryLayer = ProviderSessionRuntimeRepositoryLive.pipe( + Layer.provide(SqlitePersistenceMemory), + ); + const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); + + // shared includes McpConfigServiceLive (real) — NOT McpConfigService.layerTest(). + const shared = Layer.mergeAll( + runtimeRepositoryLayer, + directoryLayer, + Layer.succeed(ProviderAdapterRegistry, registry), + ServerSettingsService.layerTest(DEFAULT_SERVER_SETTINGS), + AnalyticsService.layerTest, + McpConfigServiceLive, + ); + + const layer = Layer.merge(shared, makeProviderServiceLive().pipe(Layer.provide(shared))).pipe( + Layer.provideMerge(ServerConfig.layerTest(cwd, { prefix: "lifecycle-int-" })), + Layer.provideMerge(NodeServices.layer), + ); + + return { cwd, harness, layer }; +}); + +const collectEventsDuring = ( + stream: Stream.Stream, + count: number, + action: Effect.Effect, +) => + Effect.gen(function* () { + const queue = yield* Queue.unbounded(); + yield* Stream.runForEach(stream, (event) => Queue.offer(queue, event).pipe(Effect.asVoid)).pipe( + Effect.forkScoped, + ); + + yield* action; + + return yield* Effect.forEach( + Array.from({ length: count }, () => undefined), + () => Queue.take(queue), + { discard: false }, + ); + }); + +// ── Resume cursor ──────────────────────────────────────────────── + +it.effect("recovers session with persisted resume cursor after adapter death", () => + Effect.gen(function* () { + const fixture = yield* makeLifecycleFixture; + + yield* Effect.gen(function* () { + const provider = yield* ProviderService; + const threadId = ThreadId.makeUnsafe("lifecycle-resume-cursor"); + + // 1. Start session — adapter generates a resume cursor. + const session = yield* provider.startSession(threadId, { + threadId, + provider: "codex", + cwd: fixture.cwd, + runtimeMode: "full-access", + }); + const originalCursor = session.resumeCursor; + assert.isDefined(originalCursor); + + // 2. Run a turn so the binding is fully exercised. + yield* fixture.harness.queueTurnResponse(threadId, { + events: codexTurnTextFixture, + }); + yield* collectEventsDuring( + provider.streamEvents, + codexTurnTextFixture.length, + provider.sendTurn({ threadId, input: "hello", attachments: [] }), + ); + + // 3. Kill adapter sessions directly (simulates process crash — + // ProviderService still has the persisted binding). + yield* fixture.harness.adapter.stopAll(); + assert.equal(fixture.harness.listActiveSessionIds().length, 0); + + // 4. Queue a response for the session that recovery will create. + yield* fixture.harness.queueTurnResponseForNextSession({ + events: codexTurnTextFixture, + }); + + // 5. sendTurn triggers automatic recovery: the ProviderService reads + // the persisted binding, extracts the resume cursor, and calls + // adapter.startSession(resumeCursor) before forwarding the turn. + const recoveryEvents = yield* collectEventsDuring( + provider.streamEvents, + codexTurnTextFixture.length, + provider.sendTurn({ + threadId, + input: "continue after crash", + attachments: [], + }), + ); + assert.equal(recoveryEvents.length, codexTurnTextFixture.length); + + // 6. Verify: the recovered session carries the original resume cursor. + const sessions = yield* provider.listSessions(); + const recovered = sessions.find((s) => String(s.threadId) === String(threadId)); + assert.isDefined(recovered); + assert.deepEqual(recovered!.resumeCursor, originalCursor); + + // 7. Verify: adapter was started exactly twice (original + recovery). + assert.equal(fixture.harness.getStartCount(), 2); + }).pipe(Effect.provide(fixture.layer)); + }).pipe(Effect.provide(NodeServices.layer)), +); + +// ── MCP snapshot lifecycle ─────────────────────────────────────── + +it.effect("MCP snapshot survives adapter death and is reused on recovery", () => + Effect.gen(function* () { + const fixture = yield* makeLifecycleFixture; + + yield* Effect.gen(function* () { + const provider = yield* ProviderService; + const mcpConfig = yield* McpConfigService; + const runtimeRepo = yield* ProviderSessionRuntimeRepository; + const fs = yield* FileSystem.FileSystem; + const { join } = yield* Path.Path; + const threadId = ThreadId.makeUnsafe("lifecycle-mcp-snapshot"); + + // 1. Write project MCP config with one server. + yield* fs.makeDirectory(join(fixture.cwd, ".t3"), { recursive: true }); + yield* fs.writeFileString( + join(fixture.cwd, ".t3", "mcp.json"), + JSON.stringify({ + servers: { + "lifecycle-server": { + command: "node", + args: ["mcp-server.js"], + transport: "stdio", + enabled: true, + }, + }, + }), + ); + + // 2. Start session — real McpConfigServiceLive resolves config from + // disk and persists a snapshot (both in-memory cache + disk). + yield* provider.startSession(threadId, { + threadId, + provider: "codex", + cwd: fixture.cwd, + runtimeMode: "full-access", + }); + + // 3. Verify: MCP config ref persisted in the session binding. + const binding = yield* runtimeRepo.getByThreadId({ threadId }); + assert.equal(binding._tag, "Some"); + if (binding._tag !== "Some") return; + const payload = binding.value.runtimePayload as Record; + const mcpRef = payload.mcpConfigRef as Record | undefined; + assert.isDefined(mcpRef); + assert.equal(mcpRef!.serverCount, 1); + + // 4. Verify: snapshot is in the MCP service cache. + const snapshot = yield* mcpConfig.getSnapshot(threadId); + assert.isNotNull(snapshot); + assert.equal(snapshot!.servers.length, 1); + assert.equal(snapshot!.servers[0]!.name, "lifecycle-server"); + + // 5. Kill adapter + delete the config file. + // Without the snapshot, re-resolution would return 0 servers. + yield* fixture.harness.adapter.stopAll(); + yield* fs.remove(join(fixture.cwd, ".t3", "mcp.json")); + + // 6. Queue a response and trigger recovery via sendTurn. + yield* fixture.harness.queueTurnResponseForNextSession({ + events: codexTurnTextFixture, + }); + yield* collectEventsDuring( + provider.streamEvents, + codexTurnTextFixture.length, + provider.sendTurn({ + threadId, + input: "continue with mcp", + attachments: [], + }), + ); + + // 7. Verify: snapshot survived recovery (not cleared — only + // stopSession clears it; adapter crash does not). + const snapshotAfterRecovery = yield* mcpConfig.getSnapshot(threadId); + assert.isNotNull(snapshotAfterRecovery); + assert.equal(snapshotAfterRecovery!.servers.length, 1); + assert.equal(snapshotAfterRecovery!.servers[0]!.name, "lifecycle-server"); + }).pipe(Effect.provide(fixture.layer)); + }).pipe(Effect.provide(NodeServices.layer)), +); + +// ── Clean teardown ─────────────────────────────────────────────── + +it.effect("stopSession cleans up binding and MCP snapshot", () => + Effect.gen(function* () { + const fixture = yield* makeLifecycleFixture; + + yield* Effect.gen(function* () { + const provider = yield* ProviderService; + const mcpConfig = yield* McpConfigService; + const runtimeRepo = yield* ProviderSessionRuntimeRepository; + const fs = yield* FileSystem.FileSystem; + const { join } = yield* Path.Path; + const threadId = ThreadId.makeUnsafe("lifecycle-stop-cleanup"); + + // Write MCP config and start session. + yield* fs.makeDirectory(join(fixture.cwd, ".t3"), { recursive: true }); + yield* fs.writeFileString( + join(fixture.cwd, ".t3", "mcp.json"), + JSON.stringify({ + servers: { + "cleanup-server": { + command: "echo", + transport: "stdio", + enabled: true, + }, + }, + }), + ); + + yield* provider.startSession(threadId, { + threadId, + provider: "codex", + cwd: fixture.cwd, + runtimeMode: "full-access", + }); + + // Pre-condition: binding and snapshot exist. + const bindingBefore = yield* runtimeRepo.getByThreadId({ threadId }); + assert.equal(bindingBefore._tag, "Some"); + const snapshotBefore = yield* mcpConfig.getSnapshot(threadId); + assert.isNotNull(snapshotBefore); + + // Stop session through the ProviderService public API. + yield* provider.stopSession({ threadId }); + + // Binding removed. + const bindingAfter = yield* runtimeRepo.getByThreadId({ threadId }); + assert.equal(bindingAfter._tag, "None"); + + // MCP snapshot cleared. + const snapshotAfter = yield* mcpConfig.getSnapshot(threadId); + assert.isNull(snapshotAfter); + }).pipe(Effect.provide(fixture.layer)); + }).pipe(Effect.provide(NodeServices.layer)), +); + +// ── Session isolation ──────────────────────────────────────────── + +it.effect("concurrent sessions on separate threads remain isolated", () => + Effect.gen(function* () { + const fixture = yield* makeLifecycleFixture; + + yield* Effect.gen(function* () { + const provider = yield* ProviderService; + const threadA = ThreadId.makeUnsafe("lifecycle-concurrent-a"); + const threadB = ThreadId.makeUnsafe("lifecycle-concurrent-b"); + + // Start two sessions on the same provider. + const sessionA = yield* provider.startSession(threadA, { + threadId: threadA, + provider: "codex", + cwd: fixture.cwd, + runtimeMode: "full-access", + }); + const sessionB = yield* provider.startSession(threadB, { + threadId: threadB, + provider: "codex", + cwd: fixture.cwd, + runtimeMode: "full-access", + }); + assert.notEqual(String(sessionA.threadId), String(sessionB.threadId)); + + // Run a turn on each. + yield* fixture.harness.queueTurnResponse(threadA, { + events: codexTurnTextFixture, + }); + yield* collectEventsDuring( + provider.streamEvents, + codexTurnTextFixture.length, + provider.sendTurn({ + threadId: threadA, + input: "turn for A", + attachments: [], + }), + ); + + yield* fixture.harness.queueTurnResponse(threadB, { + events: codexTurnTextFixture, + }); + yield* collectEventsDuring( + provider.streamEvents, + codexTurnTextFixture.length, + provider.sendTurn({ + threadId: threadB, + input: "turn for B", + attachments: [], + }), + ); + + // Stop A — B must not be affected. + yield* provider.stopSession({ threadId: threadA }); + + const sessions = yield* provider.listSessions(); + assert.equal(sessions.length, 1); + assert.equal(String(sessions[0]!.threadId), String(threadB)); + + // B is still fully operational. + yield* fixture.harness.queueTurnResponse(threadB, { + events: codexTurnTextFixture, + }); + yield* collectEventsDuring( + provider.streamEvents, + codexTurnTextFixture.length, + provider.sendTurn({ + threadId: threadB, + input: "another turn for B", + attachments: [], + }), + ); + }).pipe(Effect.provide(fixture.layer)); + }).pipe(Effect.provide(NodeServices.layer)), +); diff --git a/apps/server/src/provider/Layers/McpConfig.test.ts b/apps/server/src/provider/Layers/McpConfig.test.ts new file mode 100644 index 000000000000..17858231568b --- /dev/null +++ b/apps/server/src/provider/Layers/McpConfig.test.ts @@ -0,0 +1,344 @@ +import nodeFs from "node:fs"; +import os from "node:os"; +import path from "node:path"; + +import { assert, it } from "@effect/vitest"; +import { ThreadId } from "@t3tools/contracts"; +import type { ResolvedMcpConfig } from "@t3tools/contracts"; +import { Effect, Layer } from "effect"; +import * as NodeServices from "@effect/platform-node/NodeServices"; + +import { ServerConfig } from "../../config.ts"; +import { McpConfigService, toPersistedMcpConfigRef } from "../Services/McpConfig.ts"; +import { McpConfigServiceLive } from "./McpConfig.ts"; + +// ── Helpers ────────────────────────────────────────────────────── + +const asThreadId = (value: string): ThreadId => ThreadId.makeUnsafe(value); + +function makeTempContext() { + const cwd = nodeFs.mkdtempSync(path.join(os.tmpdir(), "mcp-cwd-")); + const baseDir = nodeFs.mkdtempSync(path.join(os.tmpdir(), "mcp-base-")); + const stateDir = path.join(baseDir, "userdata"); + return { cwd, baseDir, stateDir }; +} + +function writeProjectConfig(cwd: string, config: unknown): void { + nodeFs.mkdirSync(path.join(cwd, ".t3"), { recursive: true }); + nodeFs.writeFileSync(path.join(cwd, ".t3", "mcp.json"), JSON.stringify(config)); +} + +function writeGlobalConfig(stateDir: string, config: unknown): void { + nodeFs.mkdirSync(path.join(stateDir, "mcp"), { recursive: true }); + nodeFs.writeFileSync(path.join(stateDir, "mcp", "global.json"), JSON.stringify(config)); +} + +function makeTestLayer(ctx: { cwd: string; baseDir: string }) { + return McpConfigServiceLive.pipe( + Layer.provide(ServerConfig.layerTest(ctx.cwd, ctx.baseDir)), + Layer.provide(NodeServices.layer), + ); +} + +// ── resolveConfig ──────────────────────────────────────────────── + +it.effect("resolveConfig reads project .t3/mcp.json", () => { + const ctx = makeTempContext(); + writeProjectConfig(ctx.cwd, { + servers: { + "file-server": { + command: "node", + args: ["serve.js"], + transport: "stdio", + enabled: true, + }, + }, + }); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const config = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + + assert.equal(config.servers.length, 1); + assert.equal(config.servers[0]!.name, "file-server"); + assert.equal(config.servers[0]!.transport, "stdio"); + assert.include(config.sourcePaths as string[], path.join(ctx.cwd, ".t3", "mcp.json")); + assert.isString(config.version); + assert.match(config.version as string, /^[a-f0-9]{16}$/); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("resolveConfig reads global config from stateDir", () => { + const ctx = makeTempContext(); + writeGlobalConfig(ctx.stateDir, { + servers: { + "global-api": { + url: "https://global.example.com", + transport: "http", + enabled: true, + }, + }, + }); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const config = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + + assert.equal(config.servers.length, 1); + assert.equal(config.servers[0]!.name, "global-api"); + assert.equal(config.servers[0]!.transport, "http"); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("project config overrides global config for same server name", () => { + const ctx = makeTempContext(); + writeGlobalConfig(ctx.stateDir, { + servers: { + "shared-server": { + command: "global-cmd", + transport: "stdio", + enabled: true, + }, + }, + }); + writeProjectConfig(ctx.cwd, { + servers: { + "shared-server": { + command: "project-cmd", + args: ["--local"], + transport: "stdio", + enabled: true, + }, + }, + }); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const config = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + + assert.equal(config.servers.length, 1); + const server = config.servers[0]!; + assert.equal(server.name, "shared-server"); + // Project config wins — global processed first, project overwrites. + if (server.transport === "stdio") { + assert.equal(server.command, "project-cmd"); + } + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("resolveConfig returns empty servers when no config files exist", () => { + const ctx = makeTempContext(); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const config = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + + assert.equal(config.servers.length, 0); + assert.equal(config.sourcePaths.length, 0); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("resolveConfig normalizes mcpServers input key", () => { + const ctx = makeTempContext(); + writeProjectConfig(ctx.cwd, { + mcpServers: { + "alt-format": { + command: "python", + args: ["-m", "server"], + transport: "stdio", + enabled: true, + }, + }, + }); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const config = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + + assert.equal(config.servers.length, 1); + assert.equal(config.servers[0]!.name, "alt-format"); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("resolveConfig normalizes mcp key with array format", () => { + const ctx = makeTempContext(); + writeProjectConfig(ctx.cwd, { + mcp: [{ name: "arr-server", command: "echo", transport: "stdio", enabled: true }], + }); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const config = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + + assert.equal(config.servers.length, 1); + assert.equal(config.servers[0]!.name, "arr-server"); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("resolveConfig normalizes transport type aliases", () => { + const ctx = makeTempContext(); + writeProjectConfig(ctx.cwd, { + servers: { + "local-tool": { type: "local", command: "python", args: ["-m", "tool"] }, + "remote-api": { type: "remote", url: "https://api.example.com" }, + }, + }); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const config = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + + assert.equal(config.servers.length, 2); + const local = config.servers.find((s) => s.name === "local-tool"); + const remote = config.servers.find((s) => s.name === "remote-api"); + assert.equal(local?.transport, "stdio"); + assert.equal(remote?.transport, "http"); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +// ── Snapshots ──────────────────────────────────────────────────── + +it.effect("resolveConfig auto-persists snapshot when threadId is provided", () => { + const ctx = makeTempContext(); + writeProjectConfig(ctx.cwd, { + servers: { + "snap-server": { command: "echo", transport: "stdio", enabled: true }, + }, + }); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const threadId = asThreadId("auto-persist-thread"); + + const resolved = yield* service.resolveConfig({ + provider: "codex", + cwd: ctx.cwd, + threadId, + }); + + const snapshot = yield* service.getSnapshot(threadId); + assert.isNotNull(snapshot); + assert.equal(snapshot!.version, resolved.version); + assert.equal(snapshot!.servers.length, 1); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("setSnapshot and getSnapshot round-trip preserves config", () => { + const ctx = makeTempContext(); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const threadId = asThreadId("round-trip-thread"); + + const config = { + version: "abc12345def67890", + resolvedAt: "2026-03-28T12:00:00.000Z", + sourcePaths: ["/some/path/mcp.json"], + servers: [ + { + name: "snapshot-server", + transport: "stdio" as const, + command: "node", + args: ["index.js"], + enabled: true, + }, + ], + } as ResolvedMcpConfig; + + yield* service.setSnapshot(threadId, config); + const retrieved = yield* service.getSnapshot(threadId); + + assert.isNotNull(retrieved); + assert.equal(retrieved!.version, "abc12345def67890"); + assert.equal(retrieved!.servers.length, 1); + assert.equal(retrieved!.servers[0]!.name, "snapshot-server"); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("clearSnapshot removes snapshot from cache and disk", () => { + const ctx = makeTempContext(); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const threadId = asThreadId("clear-thread"); + + const config = { + version: "to-clear-00000000", + resolvedAt: "2026-03-28T00:00:00.000Z", + sourcePaths: [], + servers: [], + } as ResolvedMcpConfig; + + yield* service.setSnapshot(threadId, config); + yield* service.clearSnapshot(threadId); + + const retrieved = yield* service.getSnapshot(threadId); + assert.isNull(retrieved); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +it.effect("getSnapshot returns null for unknown threadId", () => { + const ctx = makeTempContext(); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const result = yield* service.getSnapshot(asThreadId("nonexistent-thread")); + assert.isNull(result); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +// ── Version ────────────────────────────────────────────────────── + +it.effect("version hash is deterministic for same server config", () => { + const ctx = makeTempContext(); + writeProjectConfig(ctx.cwd, { + servers: { + "stable-server": { command: "node", transport: "stdio", enabled: true }, + }, + }); + + return Effect.gen(function* () { + const service = yield* McpConfigService; + const first = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + const second = yield* service.resolveConfig({ provider: "codex", cwd: ctx.cwd }); + + assert.equal(first.version, second.version); + assert.match(first.version as string, /^[a-f0-9]{16}$/); + }).pipe(Effect.provide(makeTestLayer(ctx))); +}); + +// ── toPersistedMcpConfigRef ────────────────────────────────────── + +it.effect("toPersistedMcpConfigRef returns undefined for empty config", () => + Effect.sync(() => { + const config = { + version: "empty-version-0000", + resolvedAt: "2026-03-28T00:00:00.000Z", + sourcePaths: [], + servers: [], + } as ResolvedMcpConfig; + + const ref = toPersistedMcpConfigRef(config); + assert.isUndefined(ref); + }), +); + +it.effect("toPersistedMcpConfigRef extracts server count from resolved config", () => + Effect.sync(() => { + const config = { + version: "v123456789abcdef", + resolvedAt: "2026-03-28T00:00:00.000Z", + sourcePaths: ["/path/mcp.json"], + servers: [ + { name: "a", transport: "stdio" as const, command: "cmd", enabled: true }, + { name: "b", transport: "http" as const, url: "http://x", enabled: true }, + ], + } as ResolvedMcpConfig; + + const ref = toPersistedMcpConfigRef(config); + assert.isDefined(ref); + assert.equal(ref!.serverCount, 2); + assert.equal(ref!.version, "v123456789abcdef"); + assert.deepEqual([...ref!.sourcePaths], ["/path/mcp.json"]); + }), +); diff --git a/apps/server/src/provider/mcpTranslation.test.ts b/apps/server/src/provider/mcpTranslation.test.ts new file mode 100644 index 000000000000..2ef6c94dd2d7 --- /dev/null +++ b/apps/server/src/provider/mcpTranslation.test.ts @@ -0,0 +1,181 @@ +import { assert, it } from "@effect/vitest"; +import { ThreadId } from "@t3tools/contracts"; +import type { McpServerConfig, ResolvedMcpConfig } from "@t3tools/contracts"; +import { Effect } from "effect"; + +import { + codexTomlFromResolved, + generatedMcpDir, + mergeCodexToml, + openCodeConfigFromResolved, +} from "./mcpTranslation.ts"; + +// ── Fixtures ───────────────────────────────────────────────────── + +function makeConfig(servers: ReadonlyArray): ResolvedMcpConfig { + return { + version: "test-hash-0000", + resolvedAt: "2026-03-28T00:00:00.000Z", + sourcePaths: [], + servers: [...servers], + } as ResolvedMcpConfig; +} + +const stdioServer = { + name: "my-server", + transport: "stdio" as const, + command: "node", + args: ["server.js", "--port=3000"], + // Keys intentionally in reverse-alphabetical order to verify sorting. + env: { NODE_ENV: "production", API_KEY: "secret123" }, + enabled: true, +} as McpServerConfig; + +const httpServer = { + name: "remote-api", + transport: "http" as const, + url: "https://api.example.com/mcp", + enabled: true, +} as McpServerConfig; + +const sseServer = { + name: "event-stream", + transport: "sse" as const, + url: "https://events.example.com/stream", + enabled: false, +} as McpServerConfig; + +// ── codexTomlFromResolved ──────────────────────────────────────── + +it.effect("codexTomlFromResolved generates TOML for stdio server with args and env", () => + Effect.sync(() => { + const result = codexTomlFromResolved(makeConfig([stdioServer])); + + assert.include(result, "[mcp_servers.my-server]"); + assert.include(result, 'command = "node"'); + assert.include(result, 'args = ["server.js", "--port=3000"]'); + assert.include(result, "enabled = true"); + assert.include(result, "[mcp_servers.my-server.env]"); + assert.include(result, '"API_KEY" = "secret123"'); + assert.include(result, '"NODE_ENV" = "production"'); + + // Env keys are sorted alphabetically regardless of insertion order. + const apiKeyPos = result.indexOf('"API_KEY"'); + const nodeEnvPos = result.indexOf('"NODE_ENV"'); + assert.isBelow(apiKeyPos, nodeEnvPos); + }), +); + +it.effect("codexTomlFromResolved generates TOML for HTTP server with url", () => + Effect.sync(() => { + const result = codexTomlFromResolved(makeConfig([httpServer])); + + assert.include(result, "[mcp_servers.remote-api]"); + assert.include(result, 'url = "https://api.example.com/mcp"'); + assert.include(result, "enabled = true"); + assert.notInclude(result, "command"); + }), +); + +it.effect("codexTomlFromResolved sanitizes server names with special characters", () => + Effect.sync(() => { + const server = { ...stdioServer, name: "my.server@v2" } as McpServerConfig; + const result = codexTomlFromResolved(makeConfig([server])); + + assert.include(result, "[mcp_servers.my_server_v2]"); + }), +); + +it.effect("codexTomlFromResolved includes generated header comment", () => + Effect.sync(() => { + const result = codexTomlFromResolved(makeConfig([stdioServer])); + + assert.include(result, "# Generated by T3 Code."); + }), +); + +// ── mergeCodexToml ─────────────────────────────────────────────── + +it.effect("mergeCodexToml appends generated block after existing config", () => + Effect.sync(() => { + const existing = '[model]\nname = "gpt-4"'; + const result = mergeCodexToml(existing, makeConfig([stdioServer])); + + assert.isTrue(result.startsWith('[model]\nname = "gpt-4"')); + assert.include(result, "# Generated by T3 Code."); + assert.include(result, "[mcp_servers.my-server]"); + }), +); + +it.effect("mergeCodexToml replaces previously generated section on re-merge", () => + Effect.sync(() => { + const existing = '[model]\nname = "gpt-4"'; + const firstMerge = mergeCodexToml(existing, makeConfig([stdioServer])); + const secondMerge = mergeCodexToml(firstMerge, makeConfig([httpServer])); + + assert.include(secondMerge, '[model]\nname = "gpt-4"'); + assert.notInclude(secondMerge, "[mcp_servers.my-server]"); + assert.include(secondMerge, "[mcp_servers.remote-api]"); + const markerCount = secondMerge.split("# Generated by T3 Code.").length - 1; + assert.equal(markerCount, 1); + }), +); + +it.effect("mergeCodexToml with empty existing config yields just generated block", () => + Effect.sync(() => { + const result = mergeCodexToml("", makeConfig([stdioServer])); + + assert.isTrue(result.startsWith("# Generated by T3 Code.")); + assert.include(result, "[mcp_servers.my-server]"); + }), +); + +// ── openCodeConfigFromResolved ─────────────────────────────────── + +it.effect("openCodeConfigFromResolved maps stdio server to local type", () => + Effect.sync(() => { + const result = openCodeConfigFromResolved(makeConfig([stdioServer])); + const parsed = JSON.parse(result) as Record; + + assert.deepNestedInclude(parsed, { $schema: "https://opencode.ai/config.json" }); + const mcp = parsed.mcp as Record>; + assert.deepEqual(mcp["my-server"], { + type: "local", + command: "node", + args: ["server.js", "--port=3000"], + env: { NODE_ENV: "production", API_KEY: "secret123" }, + enabled: true, + }); + }), +); + +it.effect("openCodeConfigFromResolved maps SSE server to sse type", () => + Effect.sync(() => { + const result = openCodeConfigFromResolved(makeConfig([sseServer])); + const parsed = JSON.parse(result) as { mcp: Record> }; + + assert.equal(parsed.mcp["event-stream"]!.type, "sse"); + assert.equal(parsed.mcp["event-stream"]!.url, "https://events.example.com/stream"); + assert.equal(parsed.mcp["event-stream"]!.enabled, false); + }), +); + +it.effect("openCodeConfigFromResolved maps HTTP server to remote type", () => + Effect.sync(() => { + const result = openCodeConfigFromResolved(makeConfig([httpServer])); + const parsed = JSON.parse(result) as { mcp: Record> }; + + assert.equal(parsed.mcp["remote-api"]!.type, "remote"); + }), +); + +// ── generatedMcpDir ────────────────────────────────────────────── + +it.effect("generatedMcpDir constructs correct per-session path", () => + Effect.sync(() => { + const threadId = ThreadId.makeUnsafe("thread-42"); + const result = generatedMcpDir("/app/state", "codex", threadId); + + assert.equal(result, "/app/state/mcp/codex/thread-42"); + }), +);