diff --git a/apps/server/scripts/acp-replay-agent.test.ts b/apps/server/scripts/acp-replay-agent.test.ts deleted file mode 100644 index 2a662a0ff808..000000000000 --- a/apps/server/scripts/acp-replay-agent.test.ts +++ /dev/null @@ -1,78 +0,0 @@ -// @effect-diagnostics nodeBuiltinImport:off -import * as NodeChildProcess from "node:child_process"; -import * as NodeFS from "node:fs"; -import * as NodeOS from "node:os"; -import * as NodePath from "node:path"; -import * as NodeURL from "node:url"; - -import { expect, it } from "vite-plus/test"; - -import { acpReplayAgentArgs } from "../src/orchestration-v2/Adapters/AcpAdapterV2.testkit.ts"; - -const replayAgentPath = NodeURL.fileURLToPath(new URL("./acp-replay-agent.ts", import.meta.url)); - -async function runRuntimeExit(status: "success" | "error" | "cancelled") { - const scratch = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-acp-replay-agent-")); - const statusPath = NodePath.join(scratch, "status.json"); - const transcript = { - scenario: `${status}-runtime-exit`, - entries: [{ type: "runtime_exit", status }], - }; - const args = acpReplayAgentArgs(replayAgentPath); - const child = NodeChildProcess.spawn(process.execPath, args, { - env: { - ...process.env, - T3_ACP_REPLAY_STATUS_PATH: statusPath, - T3_ACP_REPLAY_TRANSCRIPT: Buffer.from(JSON.stringify(transcript), "utf8").toString("base64"), - T3_ACP_REPLAY_WORKSPACE: scratch, - }, - stdio: ["pipe", "pipe", "pipe"], - }); - let stderr = ""; - child.stderr.setEncoding("utf8"); - child.stderr.on("data", (chunk: string) => { - stderr += chunk; - }); - child.stdin.end(); - - try { - const code = await new Promise((resolve, reject) => { - child.once("error", reject); - child.once("exit", resolve); - }); - return { - code, - status: JSON.parse(NodeFS.readFileSync(statusPath, "utf8")) as unknown, - stderr, - transcript, - args, - }; - } finally { - NodeFS.rmSync(scratch, { recursive: true, force: true }); - } -} - -it("accepts a recorded cancelled runtime exit", async () => { - const result = await runRuntimeExit("cancelled"); - expect(result.code, result.stderr).toBe(0); - expect(result.args).toEqual(["--experimental-strip-types", replayAgentPath]); - expect(result.status).toEqual({ - scenario: result.transcript.scenario, - cursor: 1, - total: 1, - }); -}); - -it("still rejects a recorded error runtime exit", async () => { - const result = await runRuntimeExit("error"); - expect(result.code).toBe(1); - expect(result.status).toMatchObject({ - scenario: result.transcript.scenario, - cursor: 0, - total: 1, - failure: { - detail: "Recorded runtime exit was error", - cursor: 0, - }, - }); -}); diff --git a/apps/server/scripts/claudeReplayRecordingConfig.test.ts b/apps/server/scripts/claudeReplayRecordingConfig.test.ts deleted file mode 100644 index 4df53b9fa21f..000000000000 --- a/apps/server/scripts/claudeReplayRecordingConfig.test.ts +++ /dev/null @@ -1,43 +0,0 @@ -import { assert, it } from "@effect/vitest"; - -import { validateClaudeReplayRecordingSelection } from "./claudeReplayRecordingConfig.ts"; - -it("rejects a query mode that is incompatible with the selected scenario", () => { - assert.throws( - () => - validateClaudeReplayRecordingSelection({ - scenario: "simple", - configuredQueryMode: "streaming", - selectedQueryMode: "fork_session_siblings", - configuredPromptCount: 1, - selectedPromptCount: 1, - }), - "Select a scenario configured for the requested query mode", - ); -}); - -it("rejects prompt overrides that do not preserve the scenario shape", () => { - assert.throws( - () => - validateClaudeReplayRecordingSelection({ - scenario: "thread_fork_native_siblings", - configuredQueryMode: "fork_session_siblings", - selectedQueryMode: "fork_session_siblings", - configuredPromptCount: 3, - selectedPromptCount: 1, - }), - "requires exactly 3 prompt(s)", - ); -}); - -it("accepts the configured scenario mode and prompt shape", () => { - assert.doesNotThrow(() => - validateClaudeReplayRecordingSelection({ - scenario: "thread_fork_native_siblings", - configuredQueryMode: "fork_session_siblings", - selectedQueryMode: "fork_session_siblings", - configuredPromptCount: 3, - selectedPromptCount: 3, - }), - ); -}); diff --git a/apps/server/scripts/codexReplayRecordingRecords.test.ts b/apps/server/scripts/codexReplayRecordingRecords.test.ts deleted file mode 100644 index 660590ca6df9..000000000000 --- a/apps/server/scripts/codexReplayRecordingRecords.test.ts +++ /dev/null @@ -1,59 +0,0 @@ -import { assert, it } from "@effect/vitest"; - -import { codexReplayRecordingOutputRecords } from "./codexReplayRecordingRecords.ts"; - -it("preserves monotonic request ids across native Codex forks", () => { - const records = [ - { - type: "expect_outbound", - label: "initialize", - frame: { id: 1, method: "initialize" }, - }, - { - type: "expect_outbound", - label: "thread/fork", - frame: { id: 4, method: "thread/fork" }, - }, - { - type: "expect_outbound", - label: "turn/start", - frame: { id: 5, method: "turn/start" }, - }, - ]; - - const output = codexReplayRecordingOutputRecords(records, { workspace: "/recording" }); - - assert.deepEqual(output, records); - assert.deepEqual( - output.map((record) => record.label), - ["initialize", "thread/fork", "turn/start"], - ); - assert.deepEqual( - output.map((record) => (record.frame as { readonly id: number }).id), - [1, 4, 5], - ); -}); - -it("names the recording cwd in outbound frames only", () => { - const output = codexReplayRecordingOutputRecords( - [ - { - type: "expect_outbound", - frame: { id: 3, method: "turn/start", params: { cwd: "/recording", input: [] } }, - }, - { - type: "emit_inbound", - frame: { method: "thread/started", params: { thread: { cwd: "/recording" } } }, - }, - ], - { workspace: "/recording" }, - ); - - assert.deepEqual( - output.map((record) => record.frame), - [ - { id: 3, method: "turn/start", params: { cwd: "", input: [] } }, - { method: "thread/started", params: { thread: { cwd: "/recording" } } }, - ], - ); -}); diff --git a/apps/server/scripts/cursorReplayRecordingWorkspace.test.ts b/apps/server/scripts/cursorReplayRecordingWorkspace.test.ts deleted file mode 100644 index d14f48326d29..000000000000 --- a/apps/server/scripts/cursorReplayRecordingWorkspace.test.ts +++ /dev/null @@ -1,44 +0,0 @@ -import { assert, it } from "@effect/vitest"; - -import { - cursorReplayPromptsForWorkspace, - cursorReplayTranscriptCwd, - shouldSeedCursorReplayWorkspace, -} from "./cursorReplayRecordingWorkspace.ts"; - -it("points the read-only prompt at the unique recording workspace", () => { - assert.deepEqual( - cursorReplayPromptsForWorkspace({ - scenario: "tool_call_read_only", - configuredPrompts: ["stale global path"], - packageJsonPath: "/tmp/t3-cursor-owned-abc123/package.json", - tsconfigPath: "/tmp/t3-cursor-owned-abc123/tsconfig.json", - }), - [ - "Read /tmp/t3-cursor-owned-abc123/package.json and /tmp/t3-cursor-owned-abc123/tsconfig.json, then answer exactly: read only tool fixture complete", - ], - ); -}); - -it("never seeds files into an externally supplied workspace", () => { - assert.isFalse( - shouldSeedCursorReplayWorkspace({ - scenario: "tool_call_read_only", - owned: false, - }), - ); - assert.isTrue( - shouldSeedCursorReplayWorkspace({ - scenario: "tool_call_read_only", - owned: true, - }), - ); -}); - -it("canonicalizes read-only tool paths to the stable fixture workspace", () => { - assert.equal( - cursorReplayTranscriptCwd("tool_call_read_only"), - "/tmp/claude-replay-tool_call_read_only", - ); - assert.equal(cursorReplayTranscriptCwd("simple"), undefined); -}); diff --git a/apps/server/scripts/replayRecorderDeferredRegistry.test.ts b/apps/server/scripts/replayRecorderDeferredRegistry.test.ts deleted file mode 100644 index 0aa0949f09c1..000000000000 --- a/apps/server/scripts/replayRecorderDeferredRegistry.test.ts +++ /dev/null @@ -1,26 +0,0 @@ -import { assert, it } from "@effect/vitest"; -import * as Deferred from "effect/Deferred"; -import * as Effect from "effect/Effect"; - -import { makeReplayRecorderDeferredRegistry } from "./replayRecorderDeferredRegistry.ts"; - -it.effect("atomically shares one deferred across concurrent get-or-create calls", () => - Effect.gen(function* () { - const registry = yield* makeReplayRecorderDeferredRegistry(); - const deferreds = yield* Effect.all( - Array.from({ length: 100 }, () => registry.getOrCreate("turn-1")), - { concurrency: "unbounded" }, - ); - - for (const deferred of deferreds) { - assert.strictEqual(deferred, deferreds[0]); - } - - yield* registry.succeed("turn-1"); - assert.isTrue(yield* Deferred.isDone(deferreds[0]!)); - yield* Effect.all(deferreds.map(Deferred.await), { - concurrency: "unbounded", - discard: true, - }); - }), -); diff --git a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.testkit.ts b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.testkit.ts index 3483b88fa75c..3d620d68f405 100644 --- a/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.testkit.ts +++ b/apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.testkit.ts @@ -146,7 +146,7 @@ export function makeAcpReplayCompletenessAssertion( ); } -export function acpReplayAgentArgs(scriptPath: string): ReadonlyArray { +function acpReplayAgentArgs(scriptPath: string): ReadonlyArray { return ["--experimental-strip-types", scriptPath]; } diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.testkit.test.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.testkit.test.ts deleted file mode 100644 index b4aee0a1bfae..000000000000 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.testkit.test.ts +++ /dev/null @@ -1,920 +0,0 @@ -import * as NodeServices from "@effect/platform-node/NodeServices"; -import { assert, describe, it } from "@effect/vitest"; -import type { SDKUserMessage } from "@anthropic-ai/claude-agent-sdk"; -import { HostProcessPlatform } from "@t3tools/shared/hostProcess"; -import { SpawnExecutableResolution } from "@t3tools/shared/shell"; -import { - ProviderInstanceId, - ProviderSessionId, - ThreadId, - type ProviderReplayEntry, -} from "@t3tools/contracts"; -import * as Cause from "effect/Cause"; -import * as Effect from "effect/Effect"; -import * as Exit from "effect/Exit"; -import * as Fiber from "effect/Fiber"; -import * as Schema from "effect/Schema"; -import * as Stream from "effect/Stream"; -import { vi } from "vite-plus/test"; - -import { ClaudeExecutableFileCheck } from "../../provider/Drivers/ClaudeExecutable.ts"; -import { - CLAUDE_PROVIDER, - ClaudeAgentSdkQueryRunnerError, - makeClaudeUserMessage, - permissionResultFromDecision, - type ClaudeAgentSdkQueryOptions, -} from "./ClaudeAdapterV2.ts"; -import { - CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - ClaudeOrchestratorReplayHarness, - ClaudeReplayFrameMismatchError, - ClaudeReplayIncompleteError, - ClaudeReplayRuntimeExitError, - ClaudeReplayUnexpectedOutboundError, - makeReplayQueryRunner, - recordInterruptedClaudeQuery, - recordMessagesUntilTurnResultAndFinalize, - resolveClaudeRecordingExecutablePath, -} from "./ClaudeAdapterV2.testkit.ts"; -import { makeProviderReplayGate } from "../testkit/ProviderReplayGate.testkit.ts"; -import { readProviderReplayTranscript } from "../testkit/ReplayTranscriptNdjson.ts"; - -const isClaudeAgentSdkQueryRunnerError = Schema.is(ClaudeAgentSdkQueryRunnerError); -const isClaudeReplayFrameMismatchError = Schema.is(ClaudeReplayFrameMismatchError); -const isClaudeReplayIncompleteError = Schema.is(ClaudeReplayIncompleteError); -const isClaudeReplayRuntimeExitError = Schema.is(ClaudeReplayRuntimeExitError); -const isClaudeReplayUnexpectedOutboundError = Schema.is(ClaudeReplayUnexpectedOutboundError); - -const claudeSdkMock = vi.hoisted(() => { - const close = vi.fn(); - const interrupt = vi.fn(async () => {}); - let prompt: AsyncIterable | undefined; - const query = vi.fn((input: { readonly prompt: AsyncIterable }) => { - prompt = input.prompt; - return { - [Symbol.asyncIterator]: () => ({ - next: async () => ({ done: true as const, value: undefined }), - }), - close, - interrupt, - }; - }); - - return { - close, - interrupt, - getPrompt: () => prompt, - query, - reset: () => { - close.mockClear(); - interrupt.mockClear(); - query.mockClear(); - prompt = undefined; - }, - }; -}); - -vi.mock("@anthropic-ai/claude-agent-sdk", () => ({ - forkSession: vi.fn(), - query: claudeSdkMock.query, -})); - -function makeGatedReplaySession() { - const options = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-gate", - } satisfies ClaudeAgentSdkQueryOptions; - const label = "background_tasks_changed:empty"; - const replayGate = makeProviderReplayGate([label]); - const runner = makeReplayQueryRunner( - { - provider: CLAUDE_PROVIDER, - protocol: CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "labeled-inbound-replay-gate", - entries: [ - { - type: "expect_outbound", - frame: { - type: "query.open", - options, - }, - }, - { - type: "emit_inbound", - label, - frame: { - type: "system", - subtype: "background_tasks_changed", - tasks: [], - uuid: "replay-gate-message", - session_id: options.sessionId, - }, - }, - { - type: "runtime_exit", - status: "success", - }, - ], - }, - { replayGate }, - ); - const session = runner.open({ - options, - threadId: ThreadId.make("thread-replay-gate"), - providerSessionId: ProviderSessionId.make("provider-session-replay-gate"), - }); - return { label, replayGate, runner, session }; -} - -const yieldToReplayStream = Effect.promise( - () => - new Promise((resolve) => { - setImmediate(resolve); - }), -); - -describe("ClaudeAdapterV2 replay testkit", () => { - it.effect("holds a labeled inbound frame until its replay gate is released", () => - Effect.gen(function* () { - const { label, replayGate, runner, session } = makeGatedReplaySession(); - const streamFiber = yield* Stream.runDrain(session.messages).pipe(Effect.forkChild); - - yield* yieldToReplayStream; - assert.isTrue(replayGate.hasReached(label)); - assert.isUndefined(streamFiber.pollUnsafe()); - - assert.isTrue(replayGate.release(label)); - yield* Fiber.join(streamFiber); - runner.assertComplete(); - }), - ); - - it.effect("releases a replay gate when its stream consumer is interrupted", () => - Effect.gen(function* () { - const { label, replayGate, session } = makeGatedReplaySession(); - const streamFiber = yield* Stream.runDrain(session.messages).pipe(Effect.forkChild); - - yield* yieldToReplayStream; - assert.isTrue(replayGate.hasReached(label)); - - yield* Fiber.interrupt(streamFiber); - assert.isFalse(replayGate.release(label)); - }), - ); - - it.effect("interrupts a delayed inbound frame without emitting or advancing it", () => - Effect.gen(function* () { - const options = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-delayed-interrupt", - } satisfies ClaudeAgentSdkQueryOptions; - const label = "delayed-inbound"; - const replayGate = makeProviderReplayGate([label]); - const runner = makeReplayQueryRunner( - { - provider: CLAUDE_PROVIDER, - protocol: CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "delayed-inbound-interruption", - entries: [ - { - type: "expect_outbound", - frame: { - type: "query.open", - options, - }, - }, - { - type: "emit_inbound", - label, - afterMs: 30_000, - frame: { - type: "assistant", - message: { content: [] }, - parent_tool_use_id: null, - session_id: options.sessionId, - uuid: "delayed-inbound-message", - }, - }, - { - type: "runtime_exit", - status: "success", - }, - ], - }, - { replayGate }, - ); - const session = runner.open({ - options, - threadId: ThreadId.make("thread-replay-delayed-interrupt"), - providerSessionId: ProviderSessionId.make("provider-session-replay-delayed-interrupt"), - }); - const emittedMessages: Array = []; - const streamFiber = yield* session.messages.pipe( - Stream.tap((message) => Effect.sync(() => emittedMessages.push(message))), - Stream.runDrain, - Effect.forkChild, - ); - - yield* yieldToReplayStream; - assert.isTrue(replayGate.hasReached(label)); - assert.isTrue(replayGate.release(label)); - yield* yieldToReplayStream; - - yield* Fiber.interrupt(streamFiber); - - assert.deepEqual(emittedMessages, []); - const incomplete = assert.throws(() => runner.assertComplete()); - assert.isTrue(isClaudeReplayIncompleteError(incomplete)); - if (isClaudeReplayIncompleteError(incomplete)) { - assert.equal(incomplete.cursor, 1); - assert.equal(incomplete.remaining, 2); - } - }), - ); - - it.effect("wakes a gated replay stream when an outbound frame mismatches", () => - Effect.gen(function* () { - const { label, replayGate, session } = makeGatedReplaySession(); - const emittedMessages: Array = []; - const streamFiber = yield* session.messages.pipe( - Stream.tap((message) => Effect.sync(() => emittedMessages.push(message))), - Stream.runDrain, - Effect.exit, - Effect.forkChild, - ); - - yield* yieldToReplayStream; - assert.isTrue(replayGate.hasReached(label)); - - const offerExit = yield* Effect.exit( - session.offer(makeClaudeUserMessage({ text: "unexpected gated prompt" })), - ); - assert.isTrue(Exit.isFailure(offerExit)); - - const streamExit = yield* Fiber.join(streamFiber); - assert.isTrue(Exit.isFailure(streamExit)); - if (Exit.isFailure(streamExit)) { - const error = Cause.squash(streamExit.cause); - assert.isTrue(isClaudeAgentSdkQueryRunnerError(error)); - if (isClaudeAgentSdkQueryRunnerError(error)) { - assert.isTrue(isClaudeReplayUnexpectedOutboundError(error.cause)); - } - } - assert.deepEqual(emittedMessages, []); - }), - ); - - it.effect("surfaces an outbound mismatch during a delay without inbound callbacks", () => - Effect.gen(function* () { - const stableOptions = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-gated-permission", - } satisfies ClaudeAgentSdkQueryOptions; - const canUseTool = vi.fn>( - async (_toolName, input, options) => ({ - behavior: "allow", - updatedInput: input, - toolUseID: options.toolUseID, - }), - ); - const label = "permission_request"; - const replayGate = makeProviderReplayGate([label]); - const runner = makeReplayQueryRunner( - { - provider: CLAUDE_PROVIDER, - protocol: CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "gated-permission-stops-after-failure", - entries: [ - { - type: "expect_outbound", - frame: { - type: "query.open", - options: stableOptions, - }, - }, - { - type: "emit_inbound", - label, - afterMs: 30_000, - frame: { - type: "permission.request", - toolName: "Read", - input: { file_path: "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/workspace/gated.ts" }, - options: { toolUseID: "tool-use-gated" }, - }, - }, - { - type: "runtime_exit", - status: "success", - }, - ], - }, - { replayGate }, - ); - const session = runner.open({ - options: { ...stableOptions, canUseTool }, - threadId: ThreadId.make("thread-replay-gated-permission"), - providerSessionId: ProviderSessionId.make("provider-session-replay-gated-permission"), - }); - const emittedMessages: Array = []; - const streamFiber = yield* session.messages.pipe( - Stream.tap((message) => Effect.sync(() => emittedMessages.push(message))), - Stream.runDrain, - Effect.exit, - Effect.forkChild, - ); - - yield* yieldToReplayStream; - assert.isTrue(replayGate.hasReached(label)); - assert.isTrue(replayGate.release(label)); - yield* yieldToReplayStream; - - const offerExit = yield* Effect.exit( - session.offer(makeClaudeUserMessage({ text: "unexpected gated prompt" })), - ); - assert.isTrue(Exit.isFailure(offerExit)); - - const streamExit = yield* Fiber.join(streamFiber); - assert.isTrue(Exit.isFailure(streamExit)); - if (Exit.isFailure(streamExit)) { - const error = Cause.squash(streamExit.cause); - assert.isTrue(isClaudeAgentSdkQueryRunnerError(error)); - if (isClaudeAgentSdkQueryRunnerError(error)) { - assert.isTrue(isClaudeReplayUnexpectedOutboundError(error.cause)); - } - } - assert.deepEqual(emittedMessages, []); - assert.equal(canUseTool.mock.calls.length, 0); - }), - ); - - it.effect("wakes a replay stream waiting on an outbound frame when that frame mismatches", () => - Effect.gen(function* () { - const options = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-mismatch", - } satisfies ClaudeAgentSdkQueryOptions; - const expectedMessage = makeClaudeUserMessage({ text: "expected prompt" }); - const actualMessage = makeClaudeUserMessage({ text: "unexpected prompt" }); - const runner = makeReplayQueryRunner({ - provider: CLAUDE_PROVIDER, - protocol: CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "outbound-mismatch-wakes-stream", - entries: [ - { - type: "expect_outbound", - frame: { - type: "query.open", - options, - }, - }, - { - type: "expect_outbound", - frame: { - type: "prompt.offer", - message: expectedMessage, - }, - }, - ], - }); - const session = runner.open({ - options, - threadId: ThreadId.make("thread-replay-mismatch"), - providerSessionId: ProviderSessionId.make("provider-session-replay-mismatch"), - }); - const streamFiber = yield* Stream.runDrain(session.messages).pipe( - Effect.exit, - Effect.forkChild, - ); - - yield* Effect.yieldNow; - const offerExit = yield* Effect.exit(session.offer(actualMessage)); - assert.isTrue(Exit.isFailure(offerExit)); - - yield* Effect.yieldNow; - assert.isDefined( - streamFiber.pollUnsafe(), - "expected the replay stream to finish after the mismatch", - ); - }), - ); - - it.effect("interrupts a replay stream waiting on its next outbound frame", () => - Effect.gen(function* () { - const options = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-outbound-interrupt", - } satisfies ClaudeAgentSdkQueryOptions; - const runner = makeReplayQueryRunner({ - provider: CLAUDE_PROVIDER, - protocol: CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "outbound-wait-interruption", - entries: [ - { - type: "expect_outbound", - frame: { - type: "query.open", - options, - }, - }, - { - type: "expect_outbound", - frame: { - type: "prompt.offer", - message: makeClaudeUserMessage({ text: "next prompt" }), - }, - }, - ], - }); - const session = runner.open({ - options, - threadId: ThreadId.make("thread-replay-outbound-interrupt"), - providerSessionId: ProviderSessionId.make("provider-session-replay-outbound-interrupt"), - }); - const streamFiber = yield* Stream.runDrain(session.messages).pipe(Effect.forkChild); - - yield* yieldToReplayStream; - assert.isUndefined(streamFiber.pollUnsafe()); - - yield* Fiber.interrupt(streamFiber); - - const incomplete = assert.throws(() => runner.assertComplete()); - assert.isTrue(isClaudeReplayIncompleteError(incomplete)); - if (isClaudeReplayIncompleteError(incomplete)) { - assert.equal(incomplete.cursor, 1); - assert.equal(incomplete.remaining, 1); - } - }), - ); - - it.effect("aborts a pending permission callback when replay fails", () => - Effect.gen(function* () { - const stableOptions = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-permission-abort", - } satisfies ClaudeAgentSdkQueryOptions; - let callbackSignal: AbortSignal | undefined; - let notifyCallbackStarted = () => {}; - const callbackStarted = new Promise((resolve) => { - notifyCallbackStarted = resolve; - }); - const canUseTool = vi.fn>( - async (_toolName, _input, options) => { - callbackSignal = options.signal; - notifyCallbackStarted(); - await new Promise((_resolve, reject) => { - const rejectOnAbort = () => reject(new Error("permission callback aborted")); - options.signal.addEventListener("abort", rejectOnAbort, { once: true }); - if (options.signal.aborted) { - rejectOnAbort(); - } - }); - return { - behavior: "allow", - updatedInput: _input, - toolUseID: options.toolUseID, - }; - }, - ); - const runner = makeReplayQueryRunner({ - provider: CLAUDE_PROVIDER, - protocol: CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "pending-permission-aborts-on-replay-failure", - entries: [ - { - type: "expect_outbound", - frame: { - type: "query.open", - options: stableOptions, - }, - }, - { - type: "emit_inbound", - frame: { - type: "permission.request", - toolName: "Read", - input: { file_path: "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/workspace/pending.ts" }, - options: { toolUseID: "tool-use-pending" }, - }, - }, - { - type: "runtime_exit", - status: "success", - }, - ], - }); - const session = runner.open({ - options: { ...stableOptions, canUseTool }, - threadId: ThreadId.make("thread-replay-permission-abort"), - providerSessionId: ProviderSessionId.make("provider-session-replay-permission-abort"), - }); - const streamFiber = yield* Stream.runDrain(session.messages).pipe( - Effect.exit, - Effect.forkChild, - ); - - yield* Effect.promise(() => callbackStarted); - const offerExit = yield* Effect.exit( - session.offer(makeClaudeUserMessage({ text: "unexpected prompt" })), - ); - - assert.isTrue(Exit.isFailure(offerExit)); - const streamExit = yield* Fiber.join(streamFiber); - assert.isTrue(Exit.isFailure(streamExit)); - if (Exit.isFailure(streamExit)) { - const error = Cause.squash(streamExit.cause); - assert.isTrue(isClaudeAgentSdkQueryRunnerError(error)); - if (isClaudeAgentSdkQueryRunnerError(error)) { - assert.isTrue(isClaudeReplayUnexpectedOutboundError(error.cause)); - } - } - assert.equal(canUseTool.mock.calls.length, 1); - assert.isTrue(callbackSignal?.aborted); - }), - ); - - it.effect("keeps permission callbacks scoped to the query that opened the stream", () => - Effect.gen(function* () { - const firstCanUseTool = vi.fn>( - async (_toolName, _input, options) => ({ - behavior: "allow", - updatedInput: { query: "first" }, - toolUseID: options.toolUseID, - decisionClassification: "user_temporary", - }), - ); - const secondCanUseTool = vi.fn>( - async (_toolName, _input, options) => ({ - behavior: "allow", - updatedInput: { query: "second" }, - toolUseID: options.toolUseID, - decisionClassification: "user_temporary", - }), - ); - const firstStableOptions = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-permission-first", - } satisfies ClaudeAgentSdkQueryOptions; - const secondStableOptions = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-permission-second", - } satisfies ClaudeAgentSdkQueryOptions; - const runner = makeReplayQueryRunner({ - provider: CLAUDE_PROVIDER, - protocol: CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "permission-callbacks-are-query-scoped", - entries: [ - { - type: "expect_outbound", - frame: { - type: "query.open", - options: firstStableOptions, - }, - }, - { - type: "expect_outbound", - frame: { - type: "query.open", - options: secondStableOptions, - }, - }, - { - type: "emit_inbound", - frame: { - type: "permission.request", - toolName: "Read", - input: { file_path: "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/workspace/first.ts" }, - options: { - toolUseID: "tool-use-first", - }, - }, - }, - { - type: "expect_outbound", - frame: { - type: "permission.response", - result: { - behavior: "allow", - updatedInput: { query: "first" }, - toolUseID: "tool-use-first", - decisionClassification: "user_temporary", - }, - }, - }, - { - type: "runtime_exit", - status: "success", - }, - { - type: "runtime_exit", - status: "success", - }, - ], - }); - const firstSession = runner.open({ - options: { - ...firstStableOptions, - canUseTool: firstCanUseTool, - }, - threadId: ThreadId.make("thread-replay-permission-first"), - providerSessionId: ProviderSessionId.make("provider-session-replay-permission-first"), - }); - const secondSession = runner.open({ - options: { - ...secondStableOptions, - canUseTool: secondCanUseTool, - }, - threadId: ThreadId.make("thread-replay-permission-second"), - providerSessionId: ProviderSessionId.make("provider-session-replay-permission-second"), - }); - - yield* Stream.runDrain(firstSession.messages); - yield* Stream.runDrain(secondSession.messages); - runner.assertComplete(); - - assert.equal(firstCanUseTool.mock.calls.length, 1); - assert.equal(secondCanUseTool.mock.calls.length, 0); - }), - ); - - it.effect("fails a replay stream when the recorded runtime was cancelled", () => - Effect.gen(function* () { - const options = { - model: "claude-sonnet-4-6", - tools: [], - permissionMode: "default", - sessionId: "session-replay-cancelled", - } satisfies ClaudeAgentSdkQueryOptions; - const runner = makeReplayQueryRunner({ - provider: CLAUDE_PROVIDER, - protocol: CLAUDE_AGENT_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "cancelled-runtime-is-not-success", - entries: [ - { - type: "expect_outbound", - frame: { - type: "query.open", - options, - }, - }, - { - type: "runtime_exit", - status: "cancelled", - }, - ], - }); - const session = runner.open({ - options, - threadId: ThreadId.make("thread-replay-cancelled"), - providerSessionId: ProviderSessionId.make("provider-session-replay-cancelled"), - }); - - const exit = yield* Effect.exit(Stream.runDrain(session.messages)); - - assert.isTrue(Exit.isFailure(exit)); - if (Exit.isFailure(exit)) { - const error = Cause.squash(exit.cause); - assert.isTrue(isClaudeAgentSdkQueryRunnerError(error)); - if (isClaudeAgentSdkQueryRunnerError(error)) { - assert.isTrue(isClaudeReplayRuntimeExitError(error.cause)); - if (isClaudeReplayRuntimeExitError(error.cause)) { - assert.equal(error.cause.status, "cancelled"); - } - } - } - }), - ); - - it.effect("closes an interrupted recording and emits runtime_exit when no tool use arrives", () => - Effect.gen(function* () { - claudeSdkMock.reset(); - const entries: Array = []; - - const recordingExit = yield* Effect.exit( - Effect.tryPromise(() => - recordInterruptedClaudeQuery({ - scenario: "interrupt-without-tool-use", - prompt: "use a tool", - modelSelection: { - instanceId: ProviderInstanceId.make("claudeAgent"), - model: "claude-sonnet-4-6", - }, - cwd: "/workspace", - sessionId: "session-interrupt-without-tool-use", - resume: false, - entries, - queryOpenLabel: "query.open", - promptOfferLabel: "prompt.offer", - interruptLabel: "query.interrupt", - interruptAfter: "tool_use", - }), - ), - ); - assert.isTrue(Exit.isFailure(recordingExit)); - - assert.equal(claudeSdkMock.close.mock.calls.length, 1); - assert.deepInclude(entries.at(-1), { - type: "runtime_exit", - status: "error", - }); - - const prompt = claudeSdkMock.getPrompt(); - assert.isDefined(prompt); - const promptIterator = prompt[Symbol.asyncIterator](); - assert.isFalse((yield* Effect.promise(() => promptIterator.next())).done); - assert.isTrue((yield* Effect.promise(() => promptIterator.next())).done); - }), - ); - - it("finalizes a resumed recording as soon as its terminal result arrives", async () => { - const next = vi.fn().mockResolvedValue({ - done: false, - value: { - type: "result", - subtype: "success", - }, - }); - const finalize = vi.fn(); - const entries: Array = []; - - const completed = await recordMessagesUntilTurnResultAndFinalize({ - iterator: { next }, - entries, - scenario: "resume-at-cursor-terminal", - finalize, - }); - - assert.isTrue(completed); - assert.equal(next.mock.calls.length, 1); - assert.equal(finalize.mock.calls.length, 1); - const terminalEntry = entries.at(-1); - assert.equal(terminalEntry?.type, "emit_inbound"); - if (terminalEntry?.type === "emit_inbound") { - assert.equal(terminalEntry.label, "result"); - } - }); - - it.effect("rejects the denied-write transcript when an allow response is substituted", () => - Effect.gen(function* () { - const rawTranscript = yield* readProviderReplayTranscript( - new URL( - "../testkit/fixtures/tool_call_denied_write/claude_transcript.ndjson", - import.meta.url, - ), - ).pipe(Effect.provide(NodeServices.layer)); - const transcript = yield* ClaudeOrchestratorReplayHarness.decodeTranscript(rawTranscript); - - const expectedOutboundFrame = (label: string): unknown => { - const entry = transcript.entries.find( - (candidate): candidate is Extract => - candidate.type === "expect_outbound" && candidate.label === label, - ); - if (entry === undefined) { - throw new Error(`denied-write transcript is missing outbound frame ${label}.`); - } - return entry.frame; - }; - const openFrame = expectedOutboundFrame("query.open") as { - readonly options: ClaudeAgentSdkQueryOptions; - }; - const offerFrame = expectedOutboundFrame("prompt.offer:1") as { - readonly message: SDKUserMessage; - }; - - const canUseTool: NonNullable = ( - toolName, - toolInput, - options, - ) => - Promise.resolve( - permissionResultFromDecision({ - toolName, - decision: "accept", - toolInput, - toolUseID: options.toolUseID, - ...(options.suggestions === undefined ? {} : { suggestions: options.suggestions }), - }), - ); - - const runner = makeReplayQueryRunner(transcript); - const session = runner.open({ - options: { ...openFrame.options, canUseTool }, - threadId: ThreadId.make("thread-denied-write-allow-substitution"), - providerSessionId: ProviderSessionId.make( - "provider-session-denied-write-allow-substitution", - ), - }); - yield* session.offer(offerFrame.message); - - const exit = yield* Effect.exit(Stream.runDrain(session.messages)); - - assert.isTrue(Exit.isFailure(exit)); - if (Exit.isFailure(exit)) { - const error = Cause.squash(exit.cause); - assert.isTrue(isClaudeAgentSdkQueryRunnerError(error)); - if (isClaudeAgentSdkQueryRunnerError(error)) { - assert.isTrue(isClaudeReplayFrameMismatchError(error.cause)); - if (isClaudeReplayFrameMismatchError(error.cause)) { - assert.equal(error.cause.label, "permission.response:Write"); - } - } - } - }), - ); -}); - -describe("resolveClaudeRecordingExecutablePath", () => { - const NPM_DIR = "C:\\Users\\dev\\AppData\\Roaming\\npm"; - const NPM_SHIM = `${NPM_DIR}\\claude.cmd`; - const NPM_PACKAGE_EXE = `${NPM_DIR}\\node_modules\\@anthropic-ai\\claude-code\\bin\\claude.exe`; - - function withRecordingResolution(input: { - readonly platform: "win32" | "darwin"; - readonly resolvedCommand: string | undefined; - readonly existingFiles?: ReadonlyArray; - }) { - const existing = new Set(input.existingFiles ?? []); - return (effect: Effect.Effect) => - effect.pipe( - Effect.provideService(HostProcessPlatform, input.platform), - Effect.provideService(SpawnExecutableResolution, () => input.resolvedCommand), - Effect.provideService(ClaudeExecutableFileCheck, (filePath) => existing.has(filePath)), - ); - } - - it.effect("returns the resolved path on non-Windows platforms", () => - Effect.gen(function* () { - assert.equal( - yield* resolveClaudeRecordingExecutablePath({}).pipe( - withRecordingResolution({ - platform: "darwin", - resolvedCommand: "/usr/local/bin/claude", - }), - ), - "/usr/local/bin/claude", - ); - }), - ); - - it.effect("leaves SDK discovery in place when nothing resolves on PATH", () => - Effect.gen(function* () { - assert.isUndefined( - yield* resolveClaudeRecordingExecutablePath({}).pipe( - withRecordingResolution({ platform: "darwin", resolvedCommand: undefined }), - ), - ); - }), - ); - - it.effect("follows a Windows launcher shim to its package entry", () => - Effect.gen(function* () { - assert.equal( - yield* resolveClaudeRecordingExecutablePath({}).pipe( - withRecordingResolution({ - platform: "win32", - resolvedCommand: NPM_SHIM, - existingFiles: [NPM_PACKAGE_EXE], - }), - ), - NPM_PACKAGE_EXE, - ); - }), - ); - - it.effect( - "leaves SDK discovery in place for a Windows launcher shim without a package entry", - () => - Effect.gen(function* () { - assert.isUndefined( - yield* resolveClaudeRecordingExecutablePath({}).pipe( - withRecordingResolution({ platform: "win32", resolvedCommand: NPM_SHIM }), - ), - ); - }), - ); -}); diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.testkit.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.testkit.ts index 05a76a0a3ebd..a072dcede3e4 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.testkit.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.testkit.ts @@ -455,7 +455,7 @@ function makeClaudeSessionForkFrame( }; } -export function makeReplayQueryRunner( +function makeReplayQueryRunner( transcript: ClaudeAgentSdkReplayTranscript, replayOptions: { readonly replayGate?: ProviderReplayGate } = {}, ): ClaudeQueryRunner { @@ -1119,7 +1119,7 @@ async function recordMessagesUntilTurnResult(input: { } } -export async function recordMessagesUntilTurnResultAndFinalize(input: { +async function recordMessagesUntilTurnResultAndFinalize(input: { readonly iterator: AsyncIterator; readonly entries: Array; readonly scenario: string; @@ -1269,21 +1269,21 @@ async function recordMessagesUntilFirstToolUse(input: { // executable discovery in place. A Windows launcher shim (`claude.cmd` and // friends) is not directly spawnable, so it only counts when // resolveClaudeSdkExecutablePath can follow it to a real package entry. -export const resolveClaudeRecordingExecutablePath = Effect.fn( - "resolveClaudeRecordingExecutablePath", -)(function* (environment: NodeJS.ProcessEnv) { - const resolveExecutable = yield* SpawnExecutableResolution; - const platform = yield* HostProcessPlatform; - const resolved = resolveExecutable("claude", platform, environment); - if (resolved === undefined) { - return undefined; - } - const executablePath = yield* resolveClaudeSdkExecutablePath(resolved, environment); - if (platform === "win32" && isWindowsClaudeLauncherShimPath(executablePath)) { - return undefined; - } - return executablePath; -}); +const resolveClaudeRecordingExecutablePath = Effect.fn("resolveClaudeRecordingExecutablePath")( + function* (environment: NodeJS.ProcessEnv) { + const resolveExecutable = yield* SpawnExecutableResolution; + const platform = yield* HostProcessPlatform; + const resolved = resolveExecutable("claude", platform, environment); + if (resolved === undefined) { + return undefined; + } + const executablePath = yield* resolveClaudeSdkExecutablePath(resolved, environment); + if (platform === "win32" && isWindowsClaudeLauncherShimPath(executablePath)) { + return undefined; + } + return executablePath; + }, +); async function openRecordingQuery(input: Parameters[0]) { const executablePath = await Effect.runPromise(resolveClaudeRecordingExecutablePath(process.env)); @@ -2092,7 +2092,7 @@ async function recordClaudeForkSessionQuery(input: { } } -export async function recordInterruptedClaudeQuery(input: { +async function recordInterruptedClaudeQuery(input: { readonly scenario: string; readonly prompt: string; readonly modelSelection: ModelSelection; diff --git a/apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.testkit.test.ts b/apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.testkit.test.ts deleted file mode 100644 index 039a033295ee..000000000000 --- a/apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.testkit.test.ts +++ /dev/null @@ -1,222 +0,0 @@ -import { assert, describe, it } from "@effect/vitest"; -import { ProviderSessionId, ThreadId } from "@t3tools/contracts"; -import * as Cause from "effect/Cause"; -import * as Effect from "effect/Effect"; -import * as Exit from "effect/Exit"; -import * as Fiber from "effect/Fiber"; -import * as Schema from "effect/Schema"; - -import { buildRuntimeInstructions } from "../../provider/RuntimeInstructions.ts"; -import { materializeReplayTranscriptRuntimeInstructions } from "../testkit/ReplayTranscriptNdjson.ts"; -import { - CURSOR_AGENT_SDK_PROTOCOL, - CURSOR_PROVIDER, - CursorAgentSdkRunnerError, -} from "./CursorAgentSdk.ts"; -import { - CursorReplayFrameMismatchError, - CursorOrchestratorReplayHarness, - makeCursorAgentSdkReplayRunner, -} from "./CursorAdapterV2.testkit.ts"; - -const isCursorAgentSdkRunnerError = Schema.is(CursorAgentSdkRunnerError); -const isCursorReplayFrameMismatchError = Schema.is(CursorReplayFrameMismatchError); - -describe("CursorAdapterV2 replay testkit", () => { - const runtimeInstructions = buildRuntimeInstructions({ - harness: "Cursor", - model: "composer-2.5", - }); - for (const [name, message] of [ - ["changed user text", `different prompt\n\n${runtimeInstructions}`], - [ - "changed runtime context", - `hello\n\n${buildRuntimeInstructions({ harness: "Cursor", model: "different-model" })}`, - ], - ] as const) { - it.effect(`rejects ${name} after materializing runtime instructions`, () => - Effect.gen(function* () { - const transcript = yield* CursorOrchestratorReplayHarness.decodeTranscript( - materializeReplayTranscriptRuntimeInstructions( - { - provider: CURSOR_PROVIDER, - protocol: CURSOR_AGENT_SDK_PROTOCOL, - version: "test", - scenario: "runtime-instructions-prompt-match", - entries: [ - { - type: "expect_outbound", - frame: { type: "agent.open", operation: "create", options: {} }, - }, - { - type: "emit_inbound", - frame: { type: "agent.opened", agentId: "agent-runtime-instructions" }, - }, - { - type: "expect_outbound", - frame: { type: "run.start", message: "hello", options: {} }, - }, - ], - }, - { driver: CURSOR_PROVIDER, model: "composer-2.5" }, - ), - ); - const runner = makeCursorAgentSdkReplayRunner(transcript); - const session = yield* runner.open({ - operation: "create", - options: {}, - threadId: ThreadId.make("thread-runtime-instructions"), - providerSessionId: ProviderSessionId.make("provider-session-runtime-instructions"), - }); - const error = yield* session.send({ message }).pipe(Effect.flip); - - assert.isTrue(isCursorAgentSdkRunnerError(error)); - if (isCursorAgentSdkRunnerError(error)) { - assert.isTrue(isCursorReplayFrameMismatchError(error.cause)); - if (isCursorReplayFrameMismatchError(error.cause)) { - assert.equal(error.cause.cursor, 2); - assert.deepEqual(error.cause.expected, { - type: "run.start", - message: `hello\n\n${runtimeInstructions}`, - options: {}, - }); - assert.deepEqual(error.cause.actual, { type: "run.start", message, options: {} }); - } - } - }), - ); - } - - it.effect("fails promptly when a run is followed by an unrelated outbound frame", () => - Effect.gen(function* () { - const agentId = "agent-replay-outbound-mismatch"; - const runId = "run-replay-outbound-mismatch"; - const runner = makeCursorAgentSdkReplayRunner({ - provider: CURSOR_PROVIDER, - protocol: CURSOR_AGENT_SDK_PROTOCOL, - version: "test", - scenario: "unrelated-outbound-after-run-started", - entries: [ - { - type: "expect_outbound", - frame: { - type: "agent.open", - operation: "create", - options: {}, - }, - }, - { - type: "emit_inbound", - frame: { - type: "agent.opened", - agentId, - }, - }, - { - type: "expect_outbound", - frame: { - type: "run.start", - message: "hello", - options: {}, - }, - }, - { - type: "emit_inbound", - frame: { - type: "run.started", - runId, - agentId, - }, - }, - { - type: "expect_outbound", - frame: { - type: "agent.close", - agentId, - }, - }, - ], - }); - const session = yield* runner.open({ - operation: "create", - options: {}, - threadId: ThreadId.make("thread-replay-outbound-mismatch"), - providerSessionId: ProviderSessionId.make("provider-session-replay-outbound-mismatch"), - }); - const run = yield* session.send({ message: "hello" }); - - const exit = yield* Effect.exit(run.wait.pipe(Effect.timeout("100 millis"))); - - assert.isTrue(Exit.isFailure(exit)); - if (Exit.isFailure(exit)) { - const error = Cause.squash(exit.cause); - assert.isTrue(isCursorAgentSdkRunnerError(error)); - if (isCursorAgentSdkRunnerError(error)) { - assert.isTrue(isCursorReplayFrameMismatchError(error.cause)); - if (isCursorReplayFrameMismatchError(error.cause)) { - assert.deepEqual(error.cause.expected, { - type: "run.cancel", - runId, - }); - assert.deepEqual(error.cause.actual, { - type: "agent.close", - agentId, - }); - } - } - } - }), - ); - - it.effect("wakes a paused run waiter when another operation mismatches", () => - Effect.gen(function* () { - const agentId = "agent-paused-mismatch"; - const runId = "run-paused-mismatch"; - const runner = makeCursorAgentSdkReplayRunner({ - provider: CURSOR_PROVIDER, - protocol: CURSOR_AGENT_SDK_PROTOCOL, - version: "test", - scenario: "paused-run-mismatch", - entries: [ - { - type: "expect_outbound", - frame: { type: "agent.open", operation: "create", options: {} }, - }, - { type: "emit_inbound", frame: { type: "agent.opened", agentId } }, - { type: "expect_outbound", frame: { type: "run.start", message: "hello", options: {} } }, - { type: "emit_inbound", frame: { type: "run.started", runId, agentId } }, - { type: "expect_outbound", frame: { type: "run.cancel", runId } }, - ], - }); - const session = yield* runner.open({ - operation: "create", - options: {}, - threadId: ThreadId.make("thread-paused-mismatch"), - providerSessionId: ProviderSessionId.make("provider-session-paused-mismatch"), - }); - const run = yield* session.send({ message: "hello" }); - const waiter = yield* run.wait.pipe(Effect.forkChild({ startImmediately: true })); - - const closeError = yield* session.close.pipe(Effect.flip); - const waitError = yield* Fiber.join(waiter).pipe(Effect.flip); - - assert.isTrue(isCursorAgentSdkRunnerError(closeError)); - assert.isTrue(isCursorAgentSdkRunnerError(waitError)); - if (!isCursorAgentSdkRunnerError(closeError) || !isCursorAgentSdkRunnerError(waitError)) { - return; - } - assert.isTrue(isCursorReplayFrameMismatchError(closeError.cause)); - assert.isTrue(isCursorReplayFrameMismatchError(waitError.cause)); - if ( - !isCursorReplayFrameMismatchError(closeError.cause) || - !isCursorReplayFrameMismatchError(waitError.cause) - ) { - return; - } - assert.strictEqual(closeError.cause.cursor, 4); - assert.deepEqual(closeError.cause.expected, { type: "run.cancel", runId }); - assert.deepEqual(closeError.cause.actual, { type: "agent.close", agentId }); - assert.strictEqual(closeError.cause, waitError.cause); - }), - ); -}); diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.testkit.test.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.testkit.test.ts deleted file mode 100644 index fd4ff7eedaf3..000000000000 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.testkit.test.ts +++ /dev/null @@ -1,167 +0,0 @@ -import { assert, describe, it } from "@effect/vitest"; -import * as Effect from "effect/Effect"; - -import { OPENCODE_PROVIDER } from "./OpenCodeAdapterV2.ts"; -import { - OPENCODE_SDK_REPLAY_PROTOCOL, - OpenCodeReplayController, - OpenCodeReplayMismatchError, -} from "./OpenCodeAdapterV2.testkit.ts"; - -describe("OpenCodeAdapterV2 replay testkit", () => { - const metadata = { - provider: OPENCODE_PROVIDER, - protocol: OPENCODE_SDK_REPLAY_PROTOCOL, - version: "test", - }; - const promptFrame = (messageID: string, text = "recorded-user") => ({ - type: "session.promptAsync", - input: { sessionID: "native-session", messageID, parts: [{ type: "text", text }] }, - }); - - it.effect("correlates a prompt ID without adopting unrelated messages or rewriting text", () => - Effect.gen(function* () { - const userEvent = (id: string) => ({ - type: "message.updated", - properties: { sessionID: "native-session", info: { id, role: "user" } }, - }); - const controller = new OpenCodeReplayController({ - ...metadata, - scenario: "prompt-message-identity", - entries: [ - { type: "expect_outbound", frame: promptFrame("recorded-user") }, - { - type: "emit_inbound", - frame: { type: "sdk.event", event: userEvent("unrelated-user") }, - }, - { - type: "emit_inbound", - frame: { type: "sdk.event", event: userEvent("recorded-user") }, - }, - { - type: "emit_inbound", - frame: { - type: "sdk.event", - event: { - type: "message.updated", - properties: { - sessionID: "native-session", - info: { id: "assistant-message", role: "assistant", parentID: "recorded-user" }, - }, - }, - }, - }, - { - type: "expect_outbound", - frame: { type: "session.revert", input: { messageID: "recorded-user" } }, - }, - { - type: "emit_inbound", - frame: { - type: "sdk.response", - operation: "session.revert", - data: { id: "native-session", revert: { messageID: "recorded-user" } }, - }, - }, - { type: "runtime_exit", status: "success" }, - ], - }); - yield* Effect.promise(() => controller.expectOutbound(promptFrame("generated-user"))); - const iterator = controller.events()[Symbol.asyncIterator](); - const unrelated = yield* Effect.promise(() => iterator.next()); - const admitted = yield* Effect.promise(() => iterator.next()); - const assistant = yield* Effect.promise(() => iterator.next()); - assert.deepEqual(unrelated.value, userEvent("unrelated-user")); - assert.deepEqual(admitted.value, userEvent("generated-user")); - assert.deepEqual(assistant.value, { - type: "message.updated", - properties: { - sessionID: "native-session", - info: { id: "assistant-message", role: "assistant", parentID: "generated-user" }, - }, - }); - yield* Effect.promise(() => - controller.expectOutbound({ - type: "session.revert", - input: { messageID: "generated-user" }, - }), - ); - assert.deepEqual(yield* Effect.promise(() => controller.response("session.revert")), { - id: "native-session", - revert: { messageID: "generated-user" }, - }); - assert.isTrue((yield* Effect.promise(() => iterator.next())).done); - controller.assertComplete(); - }), - ); - - it.effect("still rejects changed prompt text when binding a generated message ID", () => - Effect.gen(function* () { - const controller = new OpenCodeReplayController({ - ...metadata, - scenario: "prompt-text-mismatch", - entries: [{ type: "expect_outbound", frame: promptFrame("recorded-user") }], - }); - const error = yield* Effect.tryPromise(() => - controller.expectOutbound(promptFrame("generated-user", "different prompt")), - ).pipe(Effect.flip); - assert.instanceOf(error.cause, OpenCodeReplayMismatchError); - }), - ); - - for (const [name, recordedId, actualId] of [ - ["rebinding an existing recorded ID", "recorded-user", "another-generated-user"], - [ - "reusing one generated ID for distinct recorded users", - "another-recorded-user", - "generated-user", - ], - ] as const) { - it.effect(`rejects ${name}`, () => - Effect.gen(function* () { - const controller = new OpenCodeReplayController({ - ...metadata, - scenario: "conflicting-message-identity", - entries: [ - { type: "expect_outbound", frame: promptFrame("recorded-user") }, - { type: "expect_outbound", frame: promptFrame(recordedId) }, - ], - }); - yield* Effect.promise(() => controller.expectOutbound(promptFrame("generated-user"))); - const error = yield* Effect.tryPromise(() => - controller.expectOutbound(promptFrame(actualId)), - ).pipe(Effect.flip); - assert.instanceOf(error.cause, OpenCodeReplayMismatchError); - }), - ); - } - - it.effect("stops an event stream when abort races with listener registration", () => - Effect.gen(function* () { - let aborted = false; - const signal = { - get aborted() { - return aborted; - }, - addEventListener: () => { - aborted = true; - }, - removeEventListener: () => {}, - } as unknown as AbortSignal; - const controller = new OpenCodeReplayController({ - provider: OPENCODE_PROVIDER, - protocol: OPENCODE_SDK_REPLAY_PROTOCOL, - version: "test", - scenario: "abort-during-listener-registration", - entries: [], - }); - const iterator = controller.events(signal)[Symbol.asyncIterator](); - - const result = yield* Effect.promise(() => iterator.next()).pipe( - Effect.timeout("100 millis"), - ); - - assert.isTrue(result.done); - }), - ); -}); diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.testkit.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.testkit.ts index 9b8a09c17dc0..74be69433fbe 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.testkit.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.testkit.ts @@ -30,7 +30,7 @@ import { OpenCodeAdapterV2Driver, } from "./OpenCodeAdapterV2.ts"; -export const OPENCODE_SDK_REPLAY_PROTOCOL = OPENCODE_SDK_PROTOCOL; +const OPENCODE_SDK_REPLAY_PROTOCOL = OPENCODE_SDK_PROTOCOL; const OpenCodeSdkReplayTranscript = Schema.Struct({ provider: Schema.Literal(OPENCODE_PROVIDER), @@ -133,7 +133,7 @@ function materializeMessageIds(value: unknown, messageIds: ReadonlyMap void>(); private failure: unknown = null; diff --git a/apps/server/src/orchestration-v2/testkit/CodexReplayFixtures.integration.test.ts b/apps/server/src/orchestration-v2/testkit/CodexReplayFixtures.integration.test.ts index 947fde9b4361..483cee2fa590 100644 --- a/apps/server/src/orchestration-v2/testkit/CodexReplayFixtures.integration.test.ts +++ b/apps/server/src/orchestration-v2/testkit/CodexReplayFixtures.integration.test.ts @@ -742,26 +742,4 @@ describe("Codex replay fixtures", () => { Object.keys(scenarioExpectations).toSorted(), ); }); - - it("rejects conflicting recorded scenarios for one canonical transcript", () => { - const transcriptFile = new URL( - "./fixtures/queued_turn/codex_transcript.ndjson", - import.meta.url, - ); - - assert.throws(() => - uniqueCanonicalTranscripts([ - { - registrationScenario: "queued_turn", - recordedScenario: "queued_turn", - transcriptFile, - }, - { - registrationScenario: "conflicting_alias", - recordedScenario: "different_recording", - transcriptFile, - }, - ]), - ); - }); }); diff --git a/apps/server/src/orchestration-v2/testkit/OrchestratorReplayFixtures.contract.test.ts b/apps/server/src/orchestration-v2/testkit/OrchestratorReplayFixtures.contract.test.ts index 749e219beb8e..71e79dfae601 100644 --- a/apps/server/src/orchestration-v2/testkit/OrchestratorReplayFixtures.contract.test.ts +++ b/apps/server/src/orchestration-v2/testkit/OrchestratorReplayFixtures.contract.test.ts @@ -1,17 +1,13 @@ import { assert, describe, it } from "@effect/vitest"; import * as NodeServices from "@effect/platform-node/NodeServices"; -import { OrchestrationV2Command, ProviderDriverKind, ProviderInstanceId } from "@t3tools/contracts"; +import { OrchestrationV2Command, ProviderInstanceId } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; import * as Schema from "effect/Schema"; -import { IdAllocatorV2, layer as idAllocatorLayer } from "../IdAllocator.ts"; +import { layer as idAllocatorLayer } from "../IdAllocator.ts"; import { provideDeterministicTestRuntime } from "./DeterministicRuntime.ts"; import { ORCHESTRATOR_REPLAY_FIXTURES } from "./fixtures/index.ts"; -import { - CODEX_MODEL_SELECTION, - materializeFixtureInput, - type OrchestratorFixtureInput, -} from "./fixtures/shared.ts"; +import { materializeFixtureInput } from "./fixtures/shared.ts"; import { readProviderReplayTranscript } from "./ReplayTranscriptNdjson.ts"; const decodeCommand = Schema.decodeUnknownEffect(OrchestrationV2Command); @@ -24,259 +20,6 @@ function assertUnique(values: ReadonlyArray, label: string) { } describe("orchestrator replay fixture contract", () => { - it.effect("materializes queued fixture messages as queue-after-active dispatches", () => - Effect.gen(function* () { - const materialized = yield* materializeFixtureInput({ - scenario: "queued-message-dispatch-mode", - fixtureInput: { - steps: [ - { type: "message", text: "active run" }, - { type: "queue_message", text: "queued run" }, - ], - }, - driver: ProviderDriverKind.make("codex"), - modelSelection: CODEX_MODEL_SELECTION, - }); - const queuedCommand = materialized.commands.find( - (command) => command.type === "message.dispatch" && command.text === "queued run", - ); - - assert.isDefined(queuedCommand); - assert.deepInclude(queuedCommand, { - dispatchMode: { type: "queue_after_active" }, - }); - }).pipe(Effect.provide(idAllocatorLayer), provideDeterministicTestRuntime), - ); - - it.effect( - "keeps consecutive queue_message inputs queued without an intermediate idle barrier", - () => - Effect.gen(function* () { - const materialized = yield* materializeFixtureInput({ - scenario: "consecutive-queue-message-ordering", - fixtureInput: { - steps: [ - { type: "message", text: "active run" }, - { type: "queue_message", text: "queued run 1" }, - { type: "queue_message", text: "queued run 2" }, - ], - }, - driver: ProviderDriverKind.make("codex"), - modelSelection: CODEX_MODEL_SELECTION, - }); - const queueCommands = materialized.commands.filter( - (command) => - command.type === "message.dispatch" && - (command.text === "queued run 1" || command.text === "queued run 2"), - ); - assert.lengthOf(queueCommands, 2); - for (const command of queueCommands) { - assert.deepInclude(command, { - dispatchMode: { type: "queue_after_active" }, - }); - } - - const firstQueueDispatchIndex = materialized.steps.findIndex( - (step) => - step.type === "dispatch" && - step.command.type === "message.dispatch" && - step.command.text === "queued run 1", - ); - const secondQueueDispatchIndex = materialized.steps.findIndex( - (step) => - step.type === "dispatch" && - step.command.type === "message.dispatch" && - step.command.text === "queued run 2", - ); - assert.isAtLeast(firstQueueDispatchIndex, 0); - assert.isAbove(secondQueueDispatchIndex, firstQueueDispatchIndex); - - const betweenQueueSteps = materialized.steps.slice( - firstQueueDispatchIndex + 1, - secondQueueDispatchIndex, - ); - assert.isFalse( - betweenQueueSteps.some( - (step) => step.type === "await" || step.type === "await_thread_idle", - ), - "consecutive queue_message steps must not await the active run or thread idle between queues", - ); - - const afterSecondQueue = materialized.steps.slice(secondQueueDispatchIndex + 1); - const barrierAwaitIndex = afterSecondQueue.findIndex((step) => step.type === "await"); - const barrierIdleIndex = afterSecondQueue.findIndex( - (step) => step.type === "await_thread_idle", - ); - assert.isAtLeast( - barrierAwaitIndex, - 0, - "the final queue_message still inserts the post-queue await barrier", - ); - assert.deepEqual(afterSecondQueue[barrierAwaitIndex], { - type: "await", - key: "run:1", - }); - assert.isAbove( - barrierIdleIndex, - barrierAwaitIndex, - "the final queue_message still inserts await_thread_idle after its await", - ); - }).pipe(Effect.provide(idAllocatorLayer), provideDeterministicTestRuntime), - ); - - it.effect( - "does not await thread idle between steer/restart and answer_next_user_input_request", - () => - Effect.gen(function* () { - for (const steeringType of ["steer", "restart"] as const) { - const answers = { "question-0": "answer" }; - const materialized = yield* materializeFixtureInput({ - scenario: `${steeringType}-then-answer-user-input-ordering`, - fixtureInput: { - steps: [ - { type: "message", text: "active run" }, - { - type: steeringType, - text: `${steeringType} active run`, - targetRunIndex: 1, - }, - { - type: "answer_next_user_input_request", - answers, - }, - ], - }, - driver: ProviderDriverKind.make("codex"), - modelSelection: CODEX_MODEL_SELECTION, - }); - - const steeringDispatchIndex = materialized.steps.findIndex( - (step) => - step.type === "dispatch" && - step.command.type === "message.dispatch" && - step.command.text === `${steeringType} active run`, - ); - const answerStepIndex = materialized.steps.findIndex( - (step) => - step.type === "respond_to_next_runtime_request" && - step.answers !== undefined && - Object.keys(step.answers).includes("question-0"), - ); - assert.isAtLeast(steeringDispatchIndex, 0, `${steeringType} dispatch must materialize`); - assert.isAbove( - answerStepIndex, - steeringDispatchIndex, - `${steeringType} answer must follow the steering dispatch`, - ); - - const betweenSteps = materialized.steps.slice(steeringDispatchIndex + 1, answerStepIndex); - assert.isFalse( - betweenSteps.some((step) => step.type === "await" || step.type === "await_thread_idle"), - `${steeringType} followed by answer_next_user_input_request must not await idle before answering`, - ); - } - }).pipe(Effect.provide(idAllocatorLayer), provideDeterministicTestRuntime), - ); - - it.effect("keeps the active dispatch barrier through queued-run replacement", () => - Effect.gen(function* () { - const materialized = yield* materializeFixtureInput({ - scenario: "queued-run-replacement-ordering", - fixtureInput: { - steps: [ - { type: "message", text: "active run" }, - { type: "queue_message", text: "queued run" }, - { type: "cancel_queued_run", targetRunIndex: 2 }, - { type: "queue_message", text: "replacement queued run" }, - ], - }, - driver: ProviderDriverKind.make("codex"), - modelSelection: CODEX_MODEL_SELECTION, - }); - const replacementQueueDispatchIndex = materialized.steps.findIndex( - (step) => - step.type === "dispatch" && - step.command.type === "message.dispatch" && - step.command.text === "replacement queued run", - ); - assert.isAtLeast(replacementQueueDispatchIndex, 0); - - const afterReplacementQueue = materialized.steps.slice(replacementQueueDispatchIndex + 1); - const barrierAwait = afterReplacementQueue.find((step) => step.type === "await"); - assert.deepEqual(barrierAwait, { - type: "await", - key: "run:1", - }); - }).pipe(Effect.provide(idAllocatorLayer), provideDeterministicTestRuntime), - ); - - it.effect("materializes fixture run-status waits against derived run IDs", () => - Effect.gen(function* () { - const idAllocator = yield* IdAllocatorV2; - const materialized = yield* materializeFixtureInput({ - scenario: "await-fixture-run-status", - fixtureInput: { - steps: [{ type: "await_run_status", targetRunIndex: 1, status: "running" }], - }, - driver: ProviderDriverKind.make("codex"), - modelSelection: CODEX_MODEL_SELECTION, - }); - const threadId = materialized.projectionThreadIds[0]; - assert.isDefined(threadId); - const runStatusWait = materialized.steps.find((step) => step.type === "await_run_status"); - - assert.deepEqual(runStatusWait, { - type: "await_run_status", - threadId, - runId: idAllocator.derive.run({ threadId, ordinal: 1 }), - status: "running", - }); - }).pipe(Effect.provide(idAllocatorLayer), provideDeterministicTestRuntime), - ); - - it.effect("keeps message ordinals separate from app run ordinals after steering", () => - Effect.gen(function* () { - const idAllocator = yield* IdAllocatorV2; - for (const steeringType of ["steer", "restart"] as const) { - const fixtureInput: OrchestratorFixtureInput = { - steps: [ - { type: "message", text: "first run" }, - { type: steeringType, text: "steer first run", targetRunIndex: 1 }, - { type: "message", text: "second run" }, - { type: "interrupt", targetRunIndex: 2 }, - ], - }; - const materialized = yield* materializeFixtureInput({ - scenario: `run-index-after-${steeringType}`, - fixtureInput, - driver: ProviderDriverKind.make("codex"), - modelSelection: CODEX_MODEL_SELECTION, - }); - const secondRunCommand = materialized.commands.find( - (command) => command.type === "message.dispatch" && command.text === "second run", - ); - assert.isDefined(secondRunCommand); - const secondRunDispatch = materialized.steps.find( - (step) => step.type === "dispatch" && step.command === secondRunCommand, - ); - assert.deepInclude(secondRunDispatch, { - type: "dispatch", - await: false, - key: "run:2", - }); - const expectedSecondRunId = idAllocator.derive.run({ - threadId: materialized.projectionThreadIds[0]!, - ordinal: 2, - }); - assert.isTrue( - materialized.steps.some( - (step) => step.type === "await_run_steerable" && step.runId === expectedSecondRunId, - ), - ); - } - }).pipe(Effect.provide(idAllocatorLayer), provideDeterministicTestRuntime), - ); - it.effect( "defines one stable input and provider-specific replay/output contracts per scenario", () => diff --git a/apps/server/src/orchestration-v2/testkit/ProviderReplayGate.testkit.test.ts b/apps/server/src/orchestration-v2/testkit/ProviderReplayGate.testkit.test.ts deleted file mode 100644 index d574ff74f800..000000000000 --- a/apps/server/src/orchestration-v2/testkit/ProviderReplayGate.testkit.test.ts +++ /dev/null @@ -1,35 +0,0 @@ -import { describe, expect, it } from "vite-plus/test"; - -import { makeProviderReplayGate } from "./ProviderReplayGate.testkit.ts"; - -describe("ProviderReplayGate", () => { - it("signals arrival before the held frame is released", async () => { - const label = "held-frame"; - const gate = makeProviderReplayGate([label]); - const reached = gate.waitForReached(label); - let emitted = false; - const emission = gate.beforeEmit(label).then(() => { - emitted = true; - }); - - expect(await reached).toBe(true); - expect(emitted).toBe(false); - gate.release(label); - await emission; - expect(emitted).toBe(true); - expect(await gate.waitForReached(label)).toBe(true); - expect(await gate.waitForReached("unknown-frame")).toBe(false); - }); - - it("stops waiting when the replay consumer is interrupted", async () => { - const label = "held-frame"; - const gate = makeProviderReplayGate([label]); - const controller = new AbortController(); - const waiting = gate.beforeEmit(label, controller.signal); - - expect(gate.hasReached(label)).toBe(true); - controller.abort(); - await waiting; - expect(gate.release(label)).toBe(true); - }); -}); diff --git a/apps/server/src/orchestration-v2/testkit/ReplayTranscriptNdjson.test.ts b/apps/server/src/orchestration-v2/testkit/ReplayTranscriptNdjson.test.ts deleted file mode 100644 index 8217ebb9a4cc..000000000000 --- a/apps/server/src/orchestration-v2/testkit/ReplayTranscriptNdjson.test.ts +++ /dev/null @@ -1,175 +0,0 @@ -import { ProviderReplayNdjsonParseError } from "./ReplayTranscriptNdjson.ts"; -import * as NodePath from "@effect/platform-node/NodePath"; -import * as NodeServices from "@effect/platform-node/NodeServices"; -import { assert, describe, it } from "@effect/vitest"; -import * as Effect from "effect/Effect"; -import * as FileSystem from "effect/FileSystem"; -import * as Path from "effect/Path"; -import * as Schema from "effect/Schema"; - -import { - decodeProviderReplayNdjson, - materializeReplayTranscriptWorkspace, - readProviderReplayTranscript, -} from "./ReplayTranscriptNdjson.ts"; - -const encodeParseError = Schema.encodeUnknownEffect(ProviderReplayNdjsonParseError); - -describe("decodeProviderReplayNdjson", () => { - it.effect("decodes a self-describing provider replay fixture", () => - Effect.gen(function* () { - const transcript = yield* decodeProviderReplayNdjson(` - {"type":"transcript_start","provider":"codex","protocol":"codex.app-server","version":"0.120.0","scenario":"simple"} - {"type":"expect_outbound","label":"initialize","frame":{"id":1,"method":"initialize","params":{}}} - {"type":"emit_inbound","label":"initialized","afterMs":5,"frame":{"id":1,"result":{"ok":true}}} - {"type":"runtime_exit","status":"success"} - `); - - assert.equal(transcript.provider, "codex"); - assert.equal(transcript.protocol, "codex.app-server"); - assert.equal(transcript.scenario, "simple"); - assert.deepEqual( - transcript.entries.map((entry) => entry.type), - ["expect_outbound", "emit_inbound", "runtime_exit"], - ); - }), - ); - - it.effect("decodes entry-only fixtures when metadata is supplied by the test", () => - Effect.gen(function* () { - const transcript = yield* decodeProviderReplayNdjson( - ` - {"type":"emit_inbound","frame":{"method":"thread/created","params":{"id":"native-thread"}}} - {"type":"runtime_exit","status":"success"} - `, - { - provider: "claudeAgent", - protocol: "claude-agent-sdk", - version: "0.2.111", - scenario: "entry-only", - }, - ); - - assert.equal(transcript.provider, "claudeAgent"); - assert.equal(transcript.entries.length, 2); - }), - ); - - it.effect("returns a schema-serializable typed parse error", () => - Effect.gen(function* () { - const error = yield* decodeProviderReplayNdjson(`{"type":`).pipe(Effect.flip); - const encoded = yield* encodeParseError(error); - - assert.equal(error._tag, "ProviderReplayNdjsonLineParseError"); - assert.equal(encoded._tag, "ProviderReplayNdjsonLineParseError"); - if (encoded._tag !== "ProviderReplayNdjsonLineParseError") { - throw new Error("Expected line parse error encoding."); - } - assert.equal(encoded.lineNumber, 1); - assert.equal(encoded.line, '{"type":'); - const cause = encoded.cause; - if (typeof cause !== "object" || cause === null || Array.isArray(cause)) { - throw new Error("Expected encoded parse cause to be a JSON object."); - } - const causeRecord = cause as Record; - assert.equal(causeRecord.name, "SchemaError"); - assert.doesNotThrow(() => JSON.stringify(encoded)); - }), - ); - - it.effect("materializes outbound workspace placeholders without weakening replay frames", () => - Effect.gen(function* () { - const transcript = yield* decodeProviderReplayNdjson(` - {"type":"transcript_start","provider":"codex","protocol":"codex.app-server","version":"0.120.0","scenario":"workspace"} - {"type":"expect_outbound","frame":{"method":"turn/start","params":{"cwd":"","nested":[""]}}} - {"type":"emit_inbound","frame":{"method":"item/completed","params":{"text":""}}} - `); - - const materialized = materializeReplayTranscriptWorkspace(transcript, "/tmp/workspace"); - - assert.deepEqual(materialized.entries[0], { - type: "expect_outbound", - frame: { - method: "turn/start", - params: { cwd: "/tmp/workspace", nested: ["/tmp/workspace"] }, - }, - }); - assert.deepEqual(materialized.entries[1], { - type: "emit_inbound", - frame: { method: "item/completed", params: { text: "" } }, - }); - }), - ); -}); - -const FILE_URL_TRANSCRIPT = `{"type":"transcript_start","provider":"codex","protocol":"codex.app-server","version":"0.120.0","scenario":"file-url-read"} -{"type":"runtime_exit","status":"success"} -`; - -it.layer(NodeServices.layer)("readProviderReplayTranscript", (it) => { - it.effect("loads a transcript behind a file URL containing an encoded space", () => - Effect.gen(function* () { - const fs = yield* FileSystem.FileSystem; - const path = yield* Path.Path; - const dir = yield* fs.makeTempDirectoryScoped({ prefix: "t3 replay fixture " }); - const filePath = path.join(dir, "transcript.ndjson"); - yield* fs.writeFileString(filePath, FILE_URL_TRANSCRIPT); - - const fileUrl = yield* path.toFileUrl(filePath); - assert.isTrue(fileUrl.pathname.includes("%20")); - - const transcript = yield* readProviderReplayTranscript(fileUrl); - assert.equal(transcript.scenario, "file-url-read"); - }), - ); - - it.effect("passes drive-letter and UNC file URLs through the Windows path service", () => - Effect.gen(function* () { - const fs = yield* FileSystem.FileSystem; - const requestedPaths: Array = []; - const recordingFs = FileSystem.FileSystem.of({ - ...fs, - readFileString: (filePath: string) => { - requestedPaths.push(filePath); - return Effect.succeed(FILE_URL_TRANSCRIPT); - }, - }); - const readAsWindows = (file: URL) => - readProviderReplayTranscript(file).pipe( - Effect.provideService(FileSystem.FileSystem, recordingFs), - Effect.provide(NodePath.layerWin32), - ); - - const driveLetter = yield* readAsWindows( - new URL("file:///C:/Users/dev/t3%20worktree/transcript.ndjson"), - ); - const unc = yield* readAsWindows(new URL("file://fileserver/shared/transcript.ndjson")); - - assert.equal(driveLetter.scenario, "file-url-read"); - assert.equal(unc.scenario, "file-url-read"); - assert.deepEqual(requestedPaths, [ - "C:\\Users\\dev\\t3 worktree\\transcript.ndjson", - "\\\\fileserver\\shared\\transcript.ndjson", - ]); - }), - ); - - it.effect("rejects non-file URLs instead of decoding their pathname", () => - Effect.gen(function* () { - const fs = yield* FileSystem.FileSystem; - const recordingFs = FileSystem.FileSystem.of({ - ...fs, - readFileString: () => Effect.succeed(FILE_URL_TRANSCRIPT), - }); - const error = yield* readProviderReplayTranscript( - new URL("https://example.com/transcript.ndjson"), - ).pipe( - Effect.provideService(FileSystem.FileSystem, recordingFs), - Effect.provide(NodePath.layerPosix), - Effect.flip, - ); - - assert.equal(error._tag, "BadArgument"); - }), - ); -}); diff --git a/packages/effect-codex-app-server/src/replay.test.ts b/packages/effect-codex-app-server/src/replay.test.ts deleted file mode 100644 index 804513d64190..000000000000 --- a/packages/effect-codex-app-server/src/replay.test.ts +++ /dev/null @@ -1,146 +0,0 @@ -import * as Exit from "effect/Exit"; -import * as Effect from "effect/Effect"; -import * as Layer from "effect/Layer"; -import * as Schema from "effect/Schema"; -import * as Scope from "effect/Scope"; - -import { assert, it } from "@effect/vitest"; - -import * as CodexClient from "./client.ts"; -import * as CodexError from "./errors.ts"; -import * as CodexReplay from "./replay.ts"; - -const isCodexAppServerTransportError = Schema.is(CodexError.CodexAppServerTransportError); -const decodeReplayError = Schema.decodeUnknownEffect(CodexReplay.CodexAppServerReplayError); -const encodeReplayError = Schema.encodeUnknownEffect(CodexReplay.CodexAppServerReplayError); - -const initializeParams = { - clientInfo: { - name: "effect-codex-app-server-test", - title: "Effect Codex App Server Test", - version: "0.0.0", - }, - capabilities: { - experimentalApi: true, - optOutNotificationMethods: null, - }, -} as const; - -const initializeResponse = { - userAgent: "replay-codex-app-server", - codexHome: "/tmp/codex-home", - platformFamily: "unix", - platformOs: "macos", -} as const; - -function buildContext(transcript: CodexReplay.CodexAppServerReplayTranscript) { - return Effect.gen(function* () { - const scope = yield* Scope.make(); - const context = yield* Layer.buildWithScope(CodexReplay.layerReplay(transcript), scope); - return { context, scope }; - }); -} - -it.effect("replays Codex app-server frames through the real client protocol", () => - Effect.gen(function* () { - const { context, scope } = yield* buildContext({ - provider: "codex", - protocol: "codex.app-server", - version: "test", - scenario: "initialize", - entries: [ - { - type: "expect_outbound", - label: "initialize", - frame: { - id: 1, - method: "initialize", - params: { - ...initializeParams, - clientInfo: { - ...initializeParams.clientInfo, - version: "older-fixture-version", - }, - }, - }, - }, - { - type: "emit_inbound", - label: "initialize", - frame: { - id: 1, - result: initializeResponse, - }, - }, - { - type: "expect_outbound", - label: "initialized", - frame: { - method: "initialized", - }, - }, - { - type: "runtime_exit", - status: "success", - }, - ], - }); - - yield* Effect.gen(function* () { - const client = yield* CodexClient.CodexAppServerClient; - assert.deepEqual(yield* client.request("initialize", initializeParams), initializeResponse); - yield* client.notify("initialized", undefined); - }).pipe(Effect.provide(context), Effect.ensuring(Scope.close(scope, Exit.void))); - }), -); - -it.effect("fails pending client requests with a schema-serializable replay mismatch", () => - Effect.gen(function* () { - const { context, scope } = yield* buildContext({ - provider: "codex", - protocol: "codex.app-server", - version: "test", - scenario: "mismatch", - entries: [ - { - type: "expect_outbound", - label: "initialize", - frame: { - id: 1, - method: "initialize", - params: initializeParams, - }, - }, - { - type: "runtime_exit", - status: "success", - }, - ], - }); - - const error = yield* Effect.gen(function* () { - const client = yield* CodexClient.CodexAppServerClient; - return yield* client.request("account/read", {}); - }).pipe(Effect.provide(context), Effect.flip, Effect.ensuring(Scope.close(scope, Exit.void))); - - assert.equal(error._tag, "CodexAppServerTransportError"); - if (!isCodexAppServerTransportError(error)) { - throw new Error("Expected transport error."); - } - - const replayError = yield* decodeReplayError(error.cause); - const encoded = yield* encodeReplayError(replayError); - - assert.equal(encoded._tag, "CodexAppServerReplayFrameMismatchError"); - if (encoded._tag !== "CodexAppServerReplayFrameMismatchError") { - throw new Error("Expected frame mismatch error."); - } - assert.equal(encoded.scenario, "mismatch"); - assert.equal(encoded.cursor, 0); - assert.deepEqual(encoded.actual, { - id: 1, - method: "account/read", - params: {}, - }); - }), -);