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
41 changes: 41 additions & 0 deletions apps/server/src/provider/Layers/CursorAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,47 @@ async function withAdapter<A>(
}

describe("CursorAdapter SDK", () => {
it("keeps two snapshot records per turn while delivering the complete SDK stream", async () => {
const agent = new FakeAgent("agent-retention", new FakeRun("initial", "agent-retention", []));
await withAdapter(new FakeCursorSdkClient(agent), (adapter) =>
Effect.gen(function* () {
const threadId = asThreadId("thread-retention");
yield* adapter.startSession({ threadId, cwd: process.cwd(), runtimeMode: "full-access" });
for (let turn = 0; turn < 3; turn++) {
const runId = `retention-${turn}`;
const chunks = Array.from({ length: 200 }, (_, index) => `chunk-${index} `);
const result = { id: runId, status: "finished" as const };
agent.nextRun = new FakeRun(
runId,
"agent-retention",
chunks.map((text) =>
makeSdkMessage({
type: "assistant",
run_id: runId,
message: { role: "assistant", content: [{ type: "text", text }] },
}),
),
result,
);
const eventsFiber = yield* collectThroughTurnCompleted(adapter);
yield* adapter.sendTurn({ threadId, input: `prompt-${turn}`, attachments: [] });
const events = yield* Fiber.join(eventsFiber);
const deltas = events.filter((event) => event.type === "content.delta");
expect(deltas).toHaveLength(chunks.length);
expect(deltas.map((event) => event.payload.delta).join("")).toBe(chunks.join(""));
const snapshot = yield* adapter.readThread(threadId);
expect(snapshot.turns).toHaveLength(turn + 1);
expect(snapshot.turns.every((entry) => entry.items.length === 2)).toBe(true);
expect(snapshot.turns.at(-1)?.items[1]).toEqual({
assistantText: chunks.join(""),
reasoningText: "",
result,
});
}
}),
);
});

it("rejects rollback without changing the retained conversation", async () => {
const run = new FakeRun("run-rollback", "agent-rollback", [
makeSdkMessage({
Expand Down
50 changes: 29 additions & 21 deletions apps/server/src/provider/Layers/CursorAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,13 +79,18 @@ interface CursorTurnSnapshot {
readonly items: Array<unknown>;
}

interface CursorTurnSummary {
assistantText: string;
reasoningText: string;
result?: CursorSdkRunResult;
}

interface CursorTurnStreamState {
readonly turnId: TurnId;
readonly runId: string;
assistantItemStarted: boolean;
assistantText: string;
readonly summary: CursorTurnSummary;
reasoningItemStarted: boolean;
reasoningText: string;
readonly seenToolItemIds: Set<string>;
readonly startedTaskIds: Set<string>;
}
Expand All @@ -103,15 +108,6 @@ interface CursorSessionContext {
stopped: boolean;
}

function appendTurnItem(context: CursorSessionContext, turnId: TurnId, item: unknown): void {
const existing = context.turns.find((turn) => turn.id === turnId);
if (existing) {
existing.items.push(item);
return;
}
context.turns.push({ id: turnId, items: [item] });
}

function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === "object" && value !== null && !Array.isArray(value);
}
Expand Down Expand Up @@ -459,7 +455,7 @@ export function makeCursorAdapter(
},
});
}
state.assistantText += delta;
state.summary.assistantText += delta;
yield* emit({
...(yield* makeEventBase({
threadId: context.threadId,
Expand Down Expand Up @@ -502,7 +498,7 @@ export function makeCursorAdapter(
},
});
}
state.reasoningText += delta;
state.summary.reasoningText += delta;
yield* emit({
...(yield* makeEventBase({
threadId: context.threadId,
Expand Down Expand Up @@ -577,7 +573,6 @@ export function makeCursorAdapter(
state: CursorTurnStreamState,
message: CursorSdkMessage,
) {
appendTurnItem(context, state.turnId, message);
yield* writeNativeEvent(context, message);

switch (message.type) {
Expand Down Expand Up @@ -743,7 +738,9 @@ export function makeCursorAdapter(
itemType: "reasoning",
status: "completed",
title: "Reasoning",
...(state.reasoningText ? { detail: state.reasoningText.slice(0, 2_000) } : {}),
...(state.summary.reasoningText
? { detail: state.summary.reasoningText.slice(0, 2_000) }
: {}),
},
});
}
Expand All @@ -761,24 +758,30 @@ export function makeCursorAdapter(
itemType: "assistant_message",
status: "completed",
title: "Assistant message",
...(state.assistantText ? { detail: state.assistantText.slice(0, 2_000) } : {}),
...(state.summary.assistantText
? { detail: state.summary.assistantText.slice(0, 2_000) }
: {}),
},
});
}
});

const drainCursorRun = Effect.fn("drainCursorRun")(function* (
context: CursorSessionContext,
turnId: TurnId,
turn: CursorTurnSnapshot,
run: CursorSdkRun,
) {
const turnId = turn.id;
const summary: CursorTurnSummary = { assistantText: "", reasoningText: "" };
// Canonical events and optional native logs own the full stream. Keep only
// one summary here, including partial text when a run fails or is cancelled.
turn.items.push(summary);
const state: CursorTurnStreamState = {
turnId,
runId: run.id,
assistantItemStarted: false,
assistantText: "",
summary,
reasoningItemStarted: false,
reasoningText: "",
seenToolItemIds: new Set(),
startedTaskIds: new Set(),
};
Expand All @@ -805,6 +808,7 @@ export function makeCursorAdapter(
}

const result = yield* resolveRunResult(run);
summary.result = result;
if (!state.assistantItemStarted && result.result) {
yield* emitAssistantText(
context,
Expand Down Expand Up @@ -1171,8 +1175,12 @@ export function makeCursorAdapter(
);

context.activeRun = run;
appendTurnItem(context, turnId, { prompt: message, runId: run.id, model: sdkModel });
const drainFiber = yield* drainCursorRun(context, turnId, run).pipe(
const turn: CursorTurnSnapshot = {
id: turnId,
items: [{ prompt: message, runId: run.id, model: sdkModel }],
};
context.turns.push(turn);
const drainFiber = yield* drainCursorRun(context, turn, run).pipe(
Effect.catch((cause) => handleRunDrainFailure(context, turnId, cause)),
Effect.catchCause((cause) => handleRunDrainFailure(context, turnId, Cause.squash(cause))),
Effect.forkIn(context.scope),
Expand Down
Loading