Skip to content
Draft
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
200 changes: 183 additions & 17 deletions apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,16 +15,21 @@ import {
} from "@t3tools/contracts";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Deferred from "effect/Deferred";
import * as Exit from "effect/Exit";
import * as Scope from "effect/Scope";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";

import * as ServerConfig from "../../config.ts";
import * as McpProviderSession from "../../mcp/McpProviderSession.ts";
import * as IdAllocator from "../IdAllocator.ts";
import type { ProviderAdapterV2Event } from "../ProviderAdapter.ts";
import { ProviderAdapterV2RuntimePolicy } from "../ProviderAdapter.ts";
import {
cursorMcpServers,
Expand All @@ -45,10 +50,18 @@ describe("CursorAdapterV2", () => {
{ status: "error", model: "custom-fable", lateModel: undefined },
{ status: "finished", model: undefined, lateModel: "gpt-6-sol" },
{ status: "finished", model: "gpt-6-sol", lateModel: null },
{ status: "interrupted", model: undefined, lateModel: undefined },
{ status: "teardown", model: undefined, lateModel: undefined },
] as const)(
"projects Cursor tasks: $status, late model $lateModel",
({ status, model, lateModel }) =>
Effect.gen(function* () {
const deltas = Array.from({ length: 64 }, (_, i) => `${i}:ą🙂\n`);
const text = deltas.join("");
const tailReady = yield* Deferred.make<void>();
const cancelled = yield* Deferred.make<void>();
const sessionScope = yield* Scope.make();
yield* Effect.addFinalizer(() => Scope.close(sessionScope, Exit.void));
const fileSystem = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const workspace = yield* fileSystem.makeTempDirectoryScoped({
Expand Down Expand Up @@ -87,6 +100,14 @@ describe("CursorAdapterV2", () => {
// partial-tool-call before tool-call-started with the same args.
// The run then ends without tool-call-completed, which a live
// run cannot produce on demand.
for (const delta of deltas) {
yield* input.onDelta!({ type: "text-delta", text: delta }).pipe(Effect.orDie);
}
for (const delta of deltas) {
yield* input.onDelta!({ type: "thinking-delta", text: delta }).pipe(
Effect.orDie,
);
}
const taskToolCall = {
type: "task" as const,
args: {
Expand Down Expand Up @@ -118,28 +139,67 @@ describe("CursorAdapterV2", () => {
},
}).pipe(Effect.orDie);
}
yield* input.onDelta!({
type: "tool-call-started",
modelCallId: "model-call",
callId: "shell-burst",
toolCall: { type: "shell", args: { command: "echo burst" } },
}).pipe(Effect.orDie);
for (const delta of deltas) {
yield* input.onDelta!({
type: "shell-output-delta",
event: { stdout: delta },
}).pipe(Effect.orDie);
}
return {
agentId: "native-cursor-lifecycle",
runId: "native-cursor-run",
wait: Effect.succeed({
id: "native-cursor-run",
requestId: "native-request",
status,
model: { id: "composer-2.5" },
durationMs: 1,
wait: Effect.gen(function* () {
yield* TestClock.adjust(50);
for (const delta of deltas) {
yield* input.onDelta!({
type: "shell-output-delta",
event: { stdout: delta },
}).pipe(Effect.orDie);
}
if (status === "finished") {
yield* input.onDelta!({
type: "tool-call-completed",
modelCallId: "model-call",
callId: "shell-burst",
toolCall: { type: "shell", args: { command: "echo burst" } },
}).pipe(Effect.orDie);
}
for (const delta of deltas) {
yield* input.onDelta!({ type: "text-delta", text: delta }).pipe(
Effect.orDie,
);
}
yield* Deferred.succeed(tailReady, undefined);
if (status === "teardown") return yield* Effect.never;
if (status === "interrupted") yield* Deferred.await(cancelled);
return {
id: "native-cursor-run",
requestId: "native-request",
status: status === "interrupted" ? "cancelled" : status,
model: { id: "composer-2.5" },
durationMs: 1,
};
}),
cancel: Effect.void,
cancel: Deferred.succeed(cancelled, undefined).pipe(Effect.asVoid),
};
}),
}),
},
});
const runtime = yield* adapter.openSession({
threadId,
providerSessionId: ProviderSessionId.make("cursor-lifecycle-session"),
modelSelection,
runtimePolicy,
});
const runtime = yield* adapter
.openSession({
threadId,
providerSessionId: ProviderSessionId.make("cursor-lifecycle-session"),
modelSelection,
runtimePolicy,
})
.pipe(Effect.provideService(Scope.Scope, sessionScope));
const providerThread = yield* runtime.ensureThread({
threadId,
modelSelection,
Expand Down Expand Up @@ -187,10 +247,114 @@ describe("CursorAdapterV2", () => {
attachments: [],
},
});
const events = yield* runtime.events.pipe(
Stream.takeUntil((event) => event.type === "turn.terminal"),
Stream.runCollect,
// Collected from the start so a session that closes mid-response is
// observed on the same terms as one the provider settles. Teardown has
// no terminal to stop at, so it waits for the flush the close owes.
const captured: Array<ProviderAdapterV2Event> = [];
const terminalArrived = yield* Deferred.make<void>();
// A close settles nothing, so the second completed text segment — the
// one only a flush can finish — is the only thing to wait for.
const flushed = yield* Deferred.make<void>();
let completedSegments = 0;
yield* runtime.events.pipe(
Stream.tap((event) => {
captured.push(event);
if (event.type === "turn.terminal") return Deferred.succeed(terminalArrived, void 0);
if (
status === "teardown" &&
event.type === "turn_item.updated" &&
event.turnItem.type === "assistant_message" &&
!event.turnItem.streaming &&
++completedSegments === 2
) {
return Deferred.succeed(flushed, void 0);
}
return Effect.void;
}),
Stream.runDrain,
Effect.forkScoped,
);
yield* Deferred.await(tailReady);
if (status === "teardown") yield* Scope.close(sessionScope, Exit.void);
if (status === "interrupted") {
yield* runtime.interruptTurn({
providerThread,
providerTurnId: (yield* IdAllocator.IdAllocatorV2).derive.providerTurn({
driver: runtime.driver,
nativeTurnId: "native-cursor-run",
}),
});
}
yield* Deferred.await(status === "teardown" ? flushed : terminalArrived);
const events = captured;
const items = events
.filter((event) => event.type === "turn_item.updated")
.map((event) => event.turnItem);
const texts = items.filter((event) => event.type === "assistant_message");
const completedTexts = texts.filter((event) => !event.streaming);
assert.deepEqual(
completedTexts.map((event) => event.text),
[text, text],
);
assert.isAtMost(texts.filter((item) => item.streaming).length, 2);
assert.isAtMost(
events.filter((event) => event.type === "message.updated" && event.message.streaming)
.length,
4,
);
const reasoning = items.filter((event) => event.type === "reasoning");
assert.deepEqual(
reasoning.map((event) => event.text),
[text],
);
const shell = items.filter((event) => event.type === "command_execution");
assert.equal(shell.at(-1)?.output, text + text);
assert.isAtMost(shell.length, 5);
assert.equal(items[0]?.type, "assistant_message");
// `ordinal` is allocated when a stream segment first buffers a delta, so it says
// nothing about when an event was queued. A client only sees emission order,
// so these assert on the event array; every index below is an index into it.
const emitted = events.flatMap((event, index) =>
event.type === "turn_item.updated" ? [{ index, item: event.turnItem }] : [],
);
const types = emitted.map(({ item }) => item.type);
const firstTool = emitted.findIndex(
({ item }) => item.type === "command_execution" || item.type === "subagent",
);
const reasoningItem = emitted.findIndex(({ item }) => item.type === "reasoning");
const lastAssistant = types.lastIndexOf("assistant_message");
// Reasoning stays its own item and reaches the client before the tool it precedes.
assert.isAbove(firstTool, 0);
assert.isAbove(reasoningItem, 0);
assert.isBelow(emitted[reasoningItem]!.index, emitted[firstTool]!.index);
// Text that follows the tool is projected after it, not pulled forward.
assert.isAbove(emitted[lastAssistant]!.index, emitted[firstTool]!.index);
// Shell output is a growing prefix of what arrived: coalescing never truncates
// or reorders it, and the terminal update carries the whole burst.
const shellOutputs = emitted.flatMap(({ item }) =>
item.type === "command_execution" ? [item.output ?? ""] : [],
);
assert.isTrue(shellOutputs.every((output) => (text + text).startsWith(output)));
assert.deepEqual(
shellOutputs,
[...shellOutputs].sort((left, right) => left.length - right.length),
);
// The coalescer's tail is projected on every path, including a close
// that settles nothing.
assert.deepEqual(
completedTexts.map((event) => event.text),
[text, text],
);
if (status === "teardown") {
// Closing the session flushes, it does not settle: no terminal, and no
// tool or subagent row ends with the turn.
assert.isUndefined(events.at(-1)?.type === "turn.terminal" ? "terminal" : undefined);
const rows = events.filter((event) => event.type === "subagent.updated");
assert.equal(rows.at(-1)?.subagent.status, "running");
assert.isNull(rows.at(-1)?.subagent.completedAt);
return;
}
assert.equal(events.at(-1)?.type, "turn.terminal");
const rows = events.filter((event) => event.type === "subagent.updated");
assert.equal(rows[0]?.subagent.status, "running");
assert.equal(rows[0]?.subagent.model, model ?? null);
Expand All @@ -203,7 +367,9 @@ describe("CursorAdapterV2", () => {
? "idle"
: status === "cancelled"
? "cancelled"
: "failed",
: status === "interrupted"
? "interrupted"
: "failed",
);
assert.isNotNull(rows.at(-1)?.subagent.completedAt);
}).pipe(Effect.scoped, Effect.provide(Layer.merge(NodeServices.layer, IdAllocator.layer))),
Expand Down
Loading
Loading