Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 0 additions & 72 deletions apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -794,16 +794,6 @@ describe("ClaudeAdapterV2 native protocol logging", () => {
}),
);

it("does not install a protocol logger when native logging is unavailable", () => {
const protocolLogger = makeClaudeAgentSdkProtocolLogger({
nativeEventLogger: undefined,
threadId: ThreadId.make("thread-1"),
providerSessionId: ProviderSessionId.make("provider-session-1"),
});

assert.equal(protocolLogger, undefined);
});

it("logs query options without leaking environment values or callback functions", () => {
const options: ClaudeAgentSdkQueryOptions = {
model: "claude-sonnet-4-6",
Expand Down Expand Up @@ -1398,11 +1388,6 @@ describe("ClaudeAdapterV2 attachments", () => {
});

describe("ClaudeAdapterV2 native fork", () => {
it("advertises Claude Agent SDK session forks", () => {
assert.equal(ClaudeProviderCapabilitiesV2.threads.canForkThread, true);
assert.equal(ClaudeProviderCapabilitiesV2.threads.canForkFromTurn, true);
});

it.effect("forks at the source assistant cursor and resumes the forked session", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down Expand Up @@ -4257,63 +4242,6 @@ describe("ClaudeAdapterV2 background wake turns", () => {
),
);

it.effect("clears the pending task when the wake notification carries no summary", () =>
Effect.scoped(
Effect.gen(function* () {
const harness = yield* makeWakeHarness;
const now = yield* DateTime.now;

yield* harness.runtime.startTurn(
makeClaudeTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now,
attemptId: RunAttemptId.make("attempt-claude-wake-5a"),
text: "Run the build in the background.",
attachments: [],
}),
);
yield* Queue.offer(harness.sdkMessages, wakeTaskStarted);
yield* Queue.offer(harness.sdkMessages, turnOneResult);
yield* awaitUntil(() => harness.terminalEvents().length === 1, "first turn terminal");

yield* Queue.offer(
harness.sdkMessages,
claudeSdkFrame({
type: "system",
subtype: "task_notification",
task_id: WAKE_TASK_ID,
status: "completed",
output_file: "/tmp/task-wake-build.log",
summary: null,
uuid: "00000000-0000-4000-8000-000000000106",
session_id: WAKE_NATIVE_SESSION,
}),
);
yield* Queue.offer(harness.sdkMessages, wakeResult);
yield* awaitUntil(() => harness.continuationRequests.length === 1, "continuation request");
assert.isNull(harness.continuationRequests[0]?.detail);

yield* harness.runtime.startTurn(
makeClaudeTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now,
attemptId: RunAttemptId.make("attempt-claude-wake-5b"),
text: "Background task completed.",
attachments: [],
providerTurnOrdinal: 2,
messageCreatedBy: "agent",
messageCreationSource: "provider",
}),
);
yield* awaitUntil(() => harness.terminalEvents().length === 2, "continuation terminal");
assert.equal(harness.terminalEvents()[1]?.status, "completed");
assert.isFalse(yield* harness.hasPendingBackgroundWork);
}).pipe(Effect.provide(Layer.merge(idAllocatorLayer, NodeServices.layer))),
),
);

it.effect("settles a continuation turn immediately when no wake output is buffered", () =>
Effect.scoped(
Effect.gen(function* () {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -900,19 +900,6 @@ function makeClaudeProviderAdapterRegistryReplayLayer(
);
}

export async function replayClaudeAgentSdkTranscript(input: {
readonly transcript: ClaudeAgentSdkReplayTranscript;
readonly prompts: ReadonlyArray<string>;
readonly modelSelection: ModelSelection;
readonly cwd?: string;
}): Promise<ReadonlyArray<SDKMessage>> {
return input.transcript.entries.flatMap((entry) =>
entry.type === "emit_inbound" && isClaudeSdkReplayMessage(entry.frame)
? [sdkMessageFromReplayFrame(entry.frame)]
: [],
);
}

function serializeReplayError(error: unknown, scenario?: string): unknown {
return error instanceof Error
? {
Expand Down
10 changes: 0 additions & 10 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -984,16 +984,6 @@ describe("CodexAdapterV2 native protocol logging", () => {
assert.nestedPropertyVal(writes[0], "event.payload.params.turnId", "native-turn");
}),
);

it("does not install a protocol logger when native logging is unavailable", () => {
const protocolLogger = makeCodexAppServerProtocolLogger({
nativeEventLogger: undefined,
threadId: ThreadId.make("thread-1"),
providerSessionId: ProviderSessionId.make("provider-session-1"),
});

assert.equal(protocolLogger, undefined);
});
});

describe("CodexAdapterV2 rollback mapping", () => {
Expand Down
15 changes: 0 additions & 15 deletions apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,6 @@ import * as McpProviderSession from "../../mcp/McpProviderSession.ts";
import { IdAllocatorV2, layer as idAllocatorLayer } from "../IdAllocator.ts";
import { ProviderAdapterV2RuntimePolicy } from "../ProviderAdapter.ts";
import {
CursorProviderCapabilitiesV2,
cursorMcpServers,
cursorRuntimeAgentPolicy,
cursorSdkModelSelection,
Expand Down Expand Up @@ -846,20 +845,6 @@ describe("CursorAdapterV2", () => {
);
});

it("advertises only capabilities exposed by the official SDK adapter", () => {
assert.isTrue(CursorProviderCapabilitiesV2.threads.canReadThreadSnapshot);
assert.isFalse(CursorProviderCapabilitiesV2.threads.canForkThread);
assert.isFalse(CursorProviderCapabilitiesV2.threads.canRollbackThread);
assert.isTrue(CursorProviderCapabilitiesV2.turns.supportsInterrupt);
assert.isFalse(CursorProviderCapabilitiesV2.turns.supportsActiveSteering);
assert.isTrue(CursorProviderCapabilitiesV2.turns.supportsSteeringByInterruptRestart);
assert.isTrue(CursorProviderCapabilitiesV2.tools.supportsMcpTools);
assert.isTrue(CursorProviderCapabilitiesV2.subagents.supportsSubagents);
assert.isFalse(CursorProviderCapabilitiesV2.subagents.exposesSubagentThreadIds);
assert.equal(CursorProviderCapabilitiesV2.identity.nativeItemIds, "weak");
assert.isFalse(CursorProviderCapabilitiesV2.approvals.supportsCommandApproval);
});

it("injects thread-scoped MCP credentials without logging them", () => {
const threadId = ThreadId.make("thread-cursor-mcp");
McpProviderSession.setMcpProviderSession({
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,6 @@ import {
makeOpenCodeProtocolLogger,
makeOpenCodeAdapterV2,
OPENCODE_PROVIDER,
OpenCodeProviderCapabilitiesV2,
reconcileOpenCodePromptAdmissionStatus,
} from "./OpenCodeAdapterV2.ts";
import { ProviderAdapterV2RuntimePolicy } from "../ProviderAdapter.ts";
Expand Down Expand Up @@ -456,39 +455,6 @@ describe("OpenCodeAdapterV2", () => {
);
}

it.effect("fails an active turn when the event stream reaches unexpected clean EOF", () =>
Effect.gen(function* () {
const nativeEvents = asyncEventStream();
const harness = yield* makeOpenCodeRuntimeHarness("eof", "root", {
event: { subscribe: async () => ({ stream: nativeEvents.stream }) },
session: {
create: async () => ({ data: { id: "root", time: { created: 1, updated: 1 } } }),
promptAsync: async () => ({ data: true }),
abort: async () => ({ data: true }),
children: async () => ({ data: [] }),
},
});
yield* harness.startTurn();
const received = yield* harness.runtime.events.pipe(
Stream.takeUntil((event) => event.type === "turn.terminal"),
Stream.runCollect,
Effect.forkScoped,
);
nativeEvents.close();
const events = yield* Fiber.join(received);
assert.isTrue(
events.some(
(event) =>
event.type === "provider_session.updated" && event.providerSession.status === "error",
),
);
const terminal = events.find((event) => event.type === "turn.terminal");
assert.equal(terminal?.status, "failed");
assert.equal(terminal?.failure?.class, "transport_error");
assert.equal(terminal?.threadDisposition, "broken");
}).pipe(Effect.provide(idAllocatorLayer), Effect.scoped),
);

it.effect("aborts external root and descendants before closing the event stream", () =>
Effect.gen(function* () {
const scope = yield* Scope.make();
Expand Down Expand Up @@ -1497,6 +1463,7 @@ describe("OpenCodeAdapterV2", () => {
const terminal = received.find((event) => event.type === "turn.terminal");
assert.equal(terminal?.status, "failed");
assert.equal(terminal?.failure?.class, "transport_error");
assert.equal(terminal?.threadDisposition, "broken");
assert.equal((yield* Effect.exit(harness.startTurn()))._tag, "Failure");
}).pipe(Effect.provide(idAllocatorLayer), Effect.scoped),
);
Expand Down Expand Up @@ -2084,17 +2051,6 @@ describe("OpenCodeAdapterV2", () => {
),
);

it("advertises the identity strengths exposed by the SDK boundary", () => {
assert.equal(OpenCodeProviderCapabilitiesV2.identity.nativeThreadIds, "strong");
assert.equal(OpenCodeProviderCapabilitiesV2.identity.nativeTurnIds, "weak");
assert.equal(OpenCodeProviderCapabilitiesV2.identity.nativeItemIds, "strong");
assert.equal(OpenCodeProviderCapabilitiesV2.identity.nativeRequestIds, "strong");
assert.isTrue(OpenCodeProviderCapabilitiesV2.threads.canForkFromTurn);
assert.isTrue(OpenCodeProviderCapabilitiesV2.turns.supportsActiveSteering);
assert.equal(OpenCodeProviderCapabilitiesV2.turns.terminalStatusQuality, "strong");
assert.isFalse(OpenCodeProviderCapabilitiesV2.subagents.canCloseSubagents);
});

it("maps native permission families to orchestration request kinds", () => {
assert.equal(openCodePermissionRequestKind("bash"), "command");
assert.equal(openCodePermissionRequestKind("read"), "file-read");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ import {
} from "./CheckpointRollbackService.ts";
import { EventSinkV2 } from "./EventSink.ts";
import { layer as idAllocatorLayer } from "./IdAllocator.ts";
import { ProjectionStoreReadError, ProjectionStoreV2 } from "./ProjectionStore.ts";
import { ProjectionStoreV2 } from "./ProjectionStore.ts";
import type { ProviderAdapterV2RollbackThreadInput } from "./ProviderAdapter.ts";
import { ProviderSessionManagerV2 } from "./ProviderSessionManager.ts";
import { RuntimePolicyV2 } from "./RuntimePolicy.ts";
Expand Down Expand Up @@ -328,50 +328,6 @@ it.effect("reports a missing provider turn as a structured rollback failure", ()
}).pipe(Effect.provide(testLayer));
});

it.effect("wraps underlying failures with an unexpected-failure reason and cause", () => {
const threadId = ThreadId.make("thread:rollback-unexpected-failure");
const providerThreadId = ProviderThreadId.make("provider-thread:rollback-unexpected-failure");
const checkpointId = CheckpointId.make("checkpoint:rollback-unexpected-failure");
const scopeId = CheckpointScopeId.make("checkpoint-scope:rollback-unexpected-failure");
const projectionError = new ProjectionStoreReadError({
threadId,
cause: new Error("database read failed"),
});
const testLayer = checkpointRollbackServiceLayer.pipe(
Layer.provide(
Layer.mergeAll(
Layer.mock(CheckpointServiceV2)({}),
Layer.mock(EventSinkV2)({}),
idAllocatorLayer,
Layer.mock(ProjectionStoreV2)({
getThreadRecords: () => Effect.fail(projectionError),
}),
Layer.mock(ProviderSessionManagerV2)({}),
Layer.mock(RuntimePolicyV2)({}),
),
),
);

return Effect.gen(function* () {
const service = yield* CheckpointRollbackServiceV2;
const error = yield* service
.execute({
threadId,
providerThreadId,
checkpointId,
scopeId,
})
.pipe(Effect.flip);

assert.equal(error.reason, "unexpected-failure");
assert.equal(
error.message,
`Failed to execute rollback target ${checkpointId} on provider thread ${providerThreadId} for thread ${threadId}.`,
);
assert.strictEqual(error.cause, projectionError);
}).pipe(Effect.provide(testLayer));
});

it.effect.each([
{ restoreFiles: true, shared: "none" },
{ restoreFiles: false, shared: "root" },
Expand Down
52 changes: 1 addition & 51 deletions apps/server/src/orchestration-v2/EffectWorker.test.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
import { assert, it } from "@effect/vitest";
import {
CommandId,
MessageId,
ProviderSessionId,
ProviderThreadId,
ProviderTurnId,
Expand Down Expand Up @@ -149,9 +148,7 @@ function makeExecutorLayer(input: {
),
Layer.succeed(
ThreadTitleRegenerationService,
ThreadTitleRegenerationService.of({
execute: ({ requestId, kind }) => record(`title:${kind.type}:${requestId}`),
}),
ThreadTitleRegenerationService.of({ execute: () => Effect.void }),
),
);
return executorLayer.pipe(
Expand Down Expand Up @@ -722,53 +719,6 @@ it.effect("backs off briefly when a due deadline loses a claim race", () =>
}).pipe(Effect.provide(TestClock.layer())),
);

it.effect("detaches a handed-off session only after the old turn terminalizes", () =>
Effect.gen(function* () {
const now = yield* DateTime.now;
const events = yield* Ref.make<ReadonlyArray<string>>([]);

yield* Effect.gen(function* () {
const executor = yield* OrchestrationEffectExecutorV2;
yield* executor.execute(restartEffect(now, { type: "detach" }));
}).pipe(Effect.provide(makeExecutorLayer({ events })));

assert.deepEqual(yield* Ref.get(events), ["interrupt", "detach", "start"]);
}),
);

it.effect("executes durable thread title generation effects", () =>
Effect.gen(function* () {
const now = DateTime.formatIso(yield* DateTime.now);
const events = yield* Ref.make<ReadonlyArray<string>>([]);
const commandId = CommandId.make("command:title-generation");
const effect: OrchestrationEffectV2 = {
id: "effect:title-generation",
commandId,
threadId,
request: {
type: "thread-title.generate",
kind: { type: "initial", messageId: MessageId.make("message:title-generation") },
},
status: "running",
attemptCount: 1,
availableAt: now,
leaseOwner: "test-worker",
leaseExpiresAt: now,
createdAt: now,
updatedAt: now,
completedAt: null,
lastError: null,
};

yield* Effect.gen(function* () {
const executor = yield* OrchestrationEffectExecutorV2;
yield* executor.execute(effect);
}).pipe(Effect.provide(makeExecutorLayer({ events })));

assert.deepEqual(yield* Ref.get(events), [`title:initial:${commandId}`]);
}),
);

it.effect("safely retries after replacement cleanup succeeds and start fails", () =>
Effect.gen(function* () {
const now = yield* DateTime.now;
Expand Down
18 changes: 0 additions & 18 deletions apps/server/src/orchestration-v2/ThreadLaunchService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,6 @@ import * as Option from "effect/Option";
import * as Ref from "effect/Ref";
import * as Stream from "effect/Stream";
import * as Schema from "effect/Schema";
import * as SqlClient from "effect/unstable/sql/SqlClient";
import * as TestClock from "effect/testing/TestClock";

import * as GitWorkflow from "../git/GitWorkflowService.ts";
Expand Down Expand Up @@ -1676,23 +1675,6 @@ it.effect("creates a strong provider-thread mapping for an imported native sessi
}).pipe(Effect.provide(harness.layer));
});

it.effect("does not depend on the legacy launch workflow table", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
const launches = yield* ThreadLaunch.ThreadLaunchService;
yield* sql`DROP TABLE orchestration_v2_thread_launch_workflows`;
const launched = yield* launches.launch(
launchInput({
command: "command:launch:no-workflow-table",
thread: "thread:launch:no-workflow-table",
message: "No private workflow state",
}),
);
assert.equal(launched.projection.messages[0]?.text, "No private workflow state");
}).pipe(Effect.provide(harness.layer));
});

it.effect("shared intake preserves durable attachment bytes after a lost launch result", () => {
const harness = makeHarness();
const files = ServerConfig.layerTest(process.cwd(), { prefix: "t3-message-intake-" }).pipe(
Expand Down
Loading
Loading