Skip to content
301 changes: 300 additions & 1 deletion apps/server/src/provider/Layers/GrokAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,11 @@ import {
} from "@t3tools/contracts";

import { ServerConfig } from "../../config.ts";
import { grokPromptSettlementBelongsToContext, makeGrokAdapter } from "./GrokAdapter.ts";
import {
grokPromptSettlementBelongsToContext,
grokTurnCompletionForPromptEpoch,
makeGrokAdapter,
} from "./GrokAdapter.ts";
const decodeGrokSettings = Schema.decodeSync(GrokSettings);

const __dirname = NodePath.dirname(NodeURL.fileURLToPath(import.meta.url));
Expand Down Expand Up @@ -122,6 +126,58 @@ it("requires a settlement to match the live Grok turn", () => {
);
});

it("emits the current epoch result when the cancelled prompt drains first", () => {
const completed = { completedStopReason: "end_turn" as const };
const cancelled = { completedStopReason: "cancelled" as const };
const afterSuperseded = grokTurnCompletionForPromptEpoch({
promptEpoch: 1,
discardBeforeEpoch: 2,
remainingPrompts: 1,
stored: undefined,
incoming: cancelled,
emitTurnCompletion: false,
});
assert.isUndefined(afterSuperseded.stored);
assert.isUndefined(afterSuperseded.emit);

const afterCurrent = grokTurnCompletionForPromptEpoch({
promptEpoch: 2,
discardBeforeEpoch: 2,
remainingPrompts: 0,
stored: afterSuperseded.stored,
incoming: completed,
emitTurnCompletion: true,
});
assert.deepEqual(afterCurrent.stored, completed);
assert.deepEqual(afterCurrent.emit, completed);
});

it("emits the current epoch result when the cancelled prompt drains last", () => {
const completed = { completedStopReason: "end_turn" as const };
const cancelled = { completedStopReason: "cancelled" as const };
const afterCurrent = grokTurnCompletionForPromptEpoch({
promptEpoch: 2,
discardBeforeEpoch: 2,
remainingPrompts: 1,
stored: undefined,
incoming: completed,
emitTurnCompletion: true,
});
assert.deepEqual(afterCurrent.stored, completed);
assert.isUndefined(afterCurrent.emit);

const afterSuperseded = grokTurnCompletionForPromptEpoch({
promptEpoch: 1,
discardBeforeEpoch: 2,
remainingPrompts: 0,
stored: afterCurrent.stored,
incoming: cancelled,
emitTurnCompletion: false,
});
assert.deepEqual(afterSuperseded.stored, completed);
assert.deepEqual(afterSuperseded.emit, completed);
});

it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => {
it.effect("starts a session and maps mock ACP prompt flow to runtime events", () =>
Effect.gen(function* () {
Expand Down Expand Up @@ -1264,4 +1320,247 @@ it.layer(grokAdapterTestLayer)("GrokAdapterLive", (it) => {
// hang until the suite timeout instead of failing here.
}).pipe(TestClock.withLive),
);

it.effect("cancels an in-flight prompt when a mid-turn sendTurn steers", () =>
Effect.gen(function* () {
const threadId = ThreadId.make("grok-steer-cancels-in-flight");
const tempDir = yield* Effect.promise(() =>
NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "grok-acp-steer-")),
);
const requestLogPath = NodePath.join(tempDir, "requests.ndjson");
const wrapperPath = yield* Effect.promise(() =>
makeMockGrokWrapper({
T3_ACP_HANG_FIRST_PROMPT_FOREVER: "1",
T3_ACP_REQUEST_LOG_PATH: requestLogPath,
}),
);
const adapter = yield* makeTestAdapter(wrapperPath);

const runtimeEvents: ProviderRuntimeEvent[] = [];
const firstTurnStarted = yield* Deferred.make<TurnId>();
const turnCompleted = yield* Deferred.make<void>();
const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) =>
Effect.gen(function* () {
runtimeEvents.push(event);
if (String(event.threadId) !== String(threadId)) {
return;
}
if (event.type === "turn.started" && event.turnId !== undefined) {
yield* Deferred.succeed(firstTurnStarted, event.turnId).pipe(Effect.ignore);
return;
}
if (event.type === "turn.completed") {
yield* Deferred.succeed(turnCompleted, undefined).pipe(Effect.ignore);
}
}),
).pipe(Effect.forkChild);

yield* adapter.startSession({
threadId,
provider: ProviderDriverKind.make("grok"),
cwd: process.cwd(),
runtimeMode: "full-access",
});

const firstSendTurnFiber = yield* adapter
.sendTurn({ threadId, input: "hang until steered", attachments: [] })
.pipe(Effect.forkChild);
const firstTurnId = yield* Deferred.await(firstTurnStarted).pipe(Effect.timeout("2 seconds"));
yield* waitForFileContent(requestLogPath, 80, '"method":"session/prompt"');

const steered = yield* adapter
.sendTurn({ threadId, input: "take this instead", attachments: [] })
.pipe(Effect.timeout("3 seconds"));
yield* Deferred.await(turnCompleted).pipe(Effect.timeout("3 seconds"));
yield* Fiber.join(firstSendTurnFiber).pipe(Effect.timeout("3 seconds"));

const requestLog = yield* Effect.promise(() => readJsonLines(requestLogPath));
const methods = requestLog.flatMap((entry) =>
typeof entry.method === "string" ? [entry.method] : [],
);
const turnStartedEvents = runtimeEvents.filter(
(event) => event.type === "turn.started" && String(event.threadId) === String(threadId),
);
const turnCompletedEvents = runtimeEvents.filter(
(event): event is Extract<ProviderRuntimeEvent, { type: "turn.completed" }> =>
event.type === "turn.completed" && String(event.threadId) === String(threadId),
);
const readySessions = yield* adapter.listSessions();
const readySession = readySessions.find((session) => session.threadId === threadId);

assert.equal(String(steered.turnId), String(firstTurnId));
assert.isTrue(methods.includes("session/cancel"));
assert.isAtLeast(methods.filter((method) => method === "session/prompt").length, 2);
assert.lengthOf(turnStartedEvents, 1);
assert.lengthOf(turnCompletedEvents, 1);
assert.equal(turnCompletedEvents[0]?.payload.state, "completed");
assert.equal(readySession?.status, "ready");
assert.isUndefined(readySession?.activeTurnId);

yield* Fiber.interrupt(runtimeEventsFiber);
yield* adapter.stopSession(threadId);
}).pipe(TestClock.withLive),
);

it.effect("keeps a steered turn completed when the cancelled prompt settles first", () =>
Effect.gen(function* () {
const threadId = ThreadId.make("grok-steer-cancelled-prompt-settles-first");
const tempDir = yield* Effect.promise(() =>
NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "grok-acp-steer-old-first-")),
);
const requestLogPath = NodePath.join(tempDir, "requests.ndjson");
const wrapperPath = yield* Effect.promise(() =>
makeMockGrokWrapper({
T3_ACP_HANG_FIRST_PROMPT_FOREVER: "1",
T3_ACP_REQUEST_LOG_PATH: requestLogPath,
}),
);
const adapter = yield* makeTestAdapter(wrapperPath);

const runtimeEvents: ProviderRuntimeEvent[] = [];
const firstTurnStarted = yield* Deferred.make<TurnId>();
const turnCompleted = yield* Deferred.make<void>();
const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) =>
Effect.gen(function* () {
runtimeEvents.push(event);
if (String(event.threadId) !== String(threadId)) {
return;
}
if (event.type === "turn.started" && event.turnId !== undefined) {
yield* Deferred.succeed(firstTurnStarted, event.turnId).pipe(Effect.ignore);
return;
}
if (event.type === "turn.completed") {
yield* Deferred.succeed(turnCompleted, undefined).pipe(Effect.ignore);
}
}),
).pipe(Effect.forkChild);

yield* adapter.startSession({
threadId,
provider: ProviderDriverKind.make("grok"),
cwd: process.cwd(),
runtimeMode: "full-access",
});

const firstSendTurnFiber = yield* adapter
.sendTurn({ threadId, input: "hang until steered", attachments: [] })
.pipe(Effect.forkChild);
yield* Deferred.await(firstTurnStarted).pipe(Effect.timeout("2 seconds"));
yield* waitForFileContent(requestLogPath, 80, '"method":"session/prompt"');

const steered = yield* adapter
.sendTurn({ threadId, input: "take this instead", attachments: [] })
.pipe(Effect.forkChild);
yield* Fiber.join(firstSendTurnFiber).pipe(Effect.timeout("3 seconds"));
const steeredResult = yield* Fiber.join(steered).pipe(Effect.timeout("3 seconds"));
yield* Deferred.await(turnCompleted).pipe(Effect.timeout("3 seconds"));

const turnCompletedEvents = runtimeEvents.filter(
(event): event is Extract<ProviderRuntimeEvent, { type: "turn.completed" }> =>
event.type === "turn.completed" && String(event.threadId) === String(threadId),
);
const readySessions = yield* adapter.listSessions();
const readySession = readySessions.find((session) => session.threadId === threadId);

assert.lengthOf(turnCompletedEvents, 1);
assert.equal(String(steeredResult.turnId), String(turnCompletedEvents[0]?.turnId));
assert.equal(turnCompletedEvents[0]?.payload.state, "completed");
assert.equal(readySession?.status, "ready");
assert.isUndefined(readySession?.activeTurnId);

yield* Fiber.interrupt(runtimeEventsFiber);
yield* adapter.stopSession(threadId);
}).pipe(TestClock.withLive),
);

it.effect("keeps the original prompt running when a steer fails during preparation", () =>
Effect.gen(function* () {
const threadId = ThreadId.make("grok-failed-steer-keeps-original-prompt");
const tempDir = yield* Effect.promise(() =>
NodeFSP.mkdtemp(NodePath.join(NodeOS.tmpdir(), "grok-acp-failed-steer-")),
);
const requestLogPath = NodePath.join(tempDir, "requests.ndjson");
const wrapperPath = yield* Effect.promise(() =>
makeMockGrokWrapper({
T3_ACP_HANG_FIRST_PROMPT_FOREVER: "1",
T3_ACP_REQUEST_LOG_PATH: requestLogPath,
}),
);
const adapter = yield* makeTestAdapter(wrapperPath);

const runtimeEvents: ProviderRuntimeEvent[] = [];
const firstTurnStarted = yield* Deferred.make<TurnId>();
const turnCompleted = yield* Deferred.make<void>();
const runtimeEventsFiber = yield* Stream.runForEach(adapter.streamEvents, (event) =>
Effect.gen(function* () {
runtimeEvents.push(event);
if (String(event.threadId) !== String(threadId)) {
return;
}
if (event.type === "turn.started" && event.turnId !== undefined) {
yield* Deferred.succeed(firstTurnStarted, event.turnId).pipe(Effect.ignore);
return;
}
if (event.type === "turn.completed") {
yield* Deferred.succeed(turnCompleted, undefined).pipe(Effect.ignore);
}
}),
).pipe(Effect.forkChild);

yield* adapter.startSession({
threadId,
provider: ProviderDriverKind.make("grok"),
cwd: process.cwd(),
runtimeMode: "full-access",
});

const firstSendTurnFiber = yield* adapter
.sendTurn({ threadId, input: "hang until a failed steer", attachments: [] })
.pipe(Effect.forkChild);
const firstTurnId = yield* Deferred.await(firstTurnStarted).pipe(Effect.timeout("2 seconds"));

const steerError = yield* Effect.flip(
adapter.sendTurn({
threadId,
input: " ",
attachments: [],
}),
);
yield* waitForFileContent(requestLogPath, 80, '"method":"session/prompt"');

const sessionsAfterFailedSteer = yield* adapter.listSessions();
const sessionAfterFailedSteer = sessionsAfterFailedSteer.find(
(session) => session.threadId === threadId,
);
const completedBeforeInterrupt = runtimeEvents.filter(
(event): event is Extract<ProviderRuntimeEvent, { type: "turn.completed" }> =>
event.type === "turn.completed" && String(event.threadId) === String(threadId),
);

yield* adapter.interruptTurn(threadId, firstTurnId).pipe(Effect.timeout("2 seconds"));
yield* Fiber.join(firstSendTurnFiber).pipe(Effect.timeout("3 seconds"));
yield* Deferred.await(turnCompleted).pipe(Effect.timeout("3 seconds"));

const turnCompletedEvents = runtimeEvents.filter(
(event): event is Extract<ProviderRuntimeEvent, { type: "turn.completed" }> =>
event.type === "turn.completed" && String(event.threadId) === String(threadId),
);
const readySessions = yield* adapter.listSessions();
const readySession = readySessions.find((session) => session.threadId === threadId);

assert.equal(steerError._tag, "ProviderAdapterValidationError");
assert.equal(sessionAfterFailedSteer?.status, "running");
assert.equal(String(sessionAfterFailedSteer?.activeTurnId), String(firstTurnId));
assert.lengthOf(completedBeforeInterrupt, 0);
assert.lengthOf(turnCompletedEvents, 1);
assert.equal(String(turnCompletedEvents[0]?.turnId), String(firstTurnId));
assert.equal(turnCompletedEvents[0]?.payload.state, "cancelled");
assert.equal(readySession?.status, "ready");
assert.isUndefined(readySession?.activeTurnId);

yield* Fiber.interrupt(runtimeEventsFiber);
yield* adapter.stopSession(threadId);
}).pipe(TestClock.withLive),
);
});
Loading
Loading