Skip to content
Open
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
314 changes: 314 additions & 0 deletions apps/server/src/provider/Layers/OpenCodeAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ import {
OpenCodeSettings,
ProviderDriverKind,
ProviderInstanceId,
type ProviderRuntimeEvent,
ThreadId,
} from "@t3tools/contracts";
import { createModelSelection } from "@t3tools/shared/model";
Expand Down Expand Up @@ -2021,6 +2022,319 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => {
}),
);

it.effect("admits a busy event as the only evidence just after the probe deadline", () =>
Effect.gen(function* () {
const adapter = yield* OpenCodeAdapter;
const threadId = asThreadId("thread-busy-event-after-probe-deadline");
const sessionId = "ses_busy_event_after_probe_deadline";
const pushEvent = makeOpenCodeEventQueue();
const fifthProbeObserved = promiseWithResolvers<void>();
const terminalObserved = promiseWithResolvers<void>();
const terminalEvents: Array<ProviderRuntimeEvent> = [];
runtimeMock.state.autoPromptEcho = false;
runtimeMock.state.createdSessionIds.push(sessionId);
runtimeMock.state.sessionStatusImplementation = async () => {
if (runtimeMock.state.sessionStatusCalls === 5) {
fifthProbeObserved.resolve(undefined);
}
return { data: {} };
};

const eventsFiber = yield* adapter.streamEvents.pipe(
Stream.filter(
(event) =>
event.threadId === threadId &&
(event.type === "turn.completed" || event.type === "runtime.error"),
),
Stream.runForEach((event) =>
Effect.sync(() => {
terminalEvents.push(event);
if (event.type === "turn.completed") {
terminalObserved.resolve(undefined);
}
}),
),
Effect.forkChild,
);
yield* adapter.startSession({
provider: ProviderDriverKind.make("opencode"),
threadId,
runtimeMode: "full-access",
});
const turn = yield* adapter.sendTurn({
threadId,
input: "Keep this admitted turn alive",
modelSelection: createModelSelection(
ProviderInstanceId.make("opencode"),
"opencode/kimi-k3",
),
});

for (const delayMs of [250, 500, 1_000, 2_000]) {
yield* advanceTestClock(delayMs);
}
yield* Effect.promise(() => fifthProbeObserved.promise);
yield* advanceTestClock(2_000);
const abortCallsBeforeAdmissionEvidence = runtimeMock.state.abortCalls.filter(
(candidate) => candidate === sessionId,
).length;

pushEvent({
id: "evt-busy-after-probe-deadline",
type: "session.status",
properties: {
sessionID: sessionId,
status: { type: "busy" },
},
});
pushEvent({
id: "evt-idle-after-probe-deadline",
type: "session.status",
properties: {
sessionID: sessionId,
status: { type: "idle" },
},
});
yield* Effect.promise(() => terminalObserved.promise);

const turnCompleted = terminalEvents.find((event) => event.type === "turn.completed");
NodeAssert.equal(turnCompleted?.payload.state, "completed");
NodeAssert.equal(
terminalEvents.some((event) => event.type === "runtime.error"),
false,
);
NodeAssert.equal(
runtimeMock.state.abortCalls.filter((candidate) => candidate === sessionId).length,
abortCallsBeforeAdmissionEvidence,
);
NodeAssert.equal(turnCompleted?.turnId, turn.turnId);

yield* Fiber.interrupt(eventsFiber);
yield* adapter.stopSession(threadId);
}),
);

it.effect("fails admission after bounded probes and grace produce no evidence", () =>
Effect.gen(function* () {
const adapter = yield* OpenCodeAdapter;
const threadId = asThreadId("thread-admission-without-evidence");
const sessionId = "ses_admission_without_evidence";
const fifthProbeObserved = promiseWithResolvers<void>();
runtimeMock.state.autoPromptEcho = false;
runtimeMock.state.createdSessionIds.push(sessionId);
runtimeMock.state.sessionStatusImplementation = async () => {
if (runtimeMock.state.sessionStatusCalls === 5) {
fifthProbeObserved.resolve(undefined);
}
return { data: {} };
};

const terminalFiber = yield* adapter.streamEvents.pipe(
Stream.filter(
(event) =>
event.threadId === threadId &&
(event.type === "turn.completed" || event.type === "runtime.error"),
),
Stream.take(2),
Stream.runCollect,
Effect.forkChild,
);
yield* adapter.startSession({
provider: ProviderDriverKind.make("opencode"),
threadId,
runtimeMode: "full-access",
});
yield* adapter.sendTurn({
threadId,
input: "Fail when admission has no evidence",
modelSelection: createModelSelection(
ProviderInstanceId.make("opencode"),
"opencode/kimi-k3",
),
});

for (const delayMs of [250, 500, 1_000, 2_000]) {
yield* advanceTestClock(delayMs);
}
yield* Effect.promise(() => fifthProbeObserved.promise);
yield* advanceTestClock(2_000);
const abortCallsBeforeTerminalFailure = runtimeMock.state.abortCalls.filter(
(candidate) => candidate === sessionId,
).length;
yield* advanceTestClock(1_000);

const terminalEvents = Array.from(yield* Fiber.join(terminalFiber));
NodeAssert.equal(runtimeMock.state.sessionStatusCalls, 5);
NodeAssert.equal(
runtimeMock.state.abortCalls.filter((candidate) => candidate === sessionId).length,
abortCallsBeforeTerminalFailure + 1,
);
NodeAssert.deepEqual(
terminalEvents.map((event) => event.type),
["turn.completed", "runtime.error"],
);
const completed = terminalEvents[0];
NodeAssert.equal(completed?.type, "turn.completed");
if (completed?.type === "turn.completed") {
NodeAssert.equal(completed.payload.state, "failed");
}

yield* adapter.stopSession(threadId);
}),
);

it.effect("does not fail admission when the turn is interrupted during grace", () =>
Effect.gen(function* () {
const adapter = yield* OpenCodeAdapter;
const threadId = asThreadId("thread-interrupt-during-admission-grace");
const sessionId = "ses_interrupt_during_admission_grace";
const fifthProbeObserved = promiseWithResolvers<void>();
const terminalEvents: Array<ProviderRuntimeEvent> = [];
runtimeMock.state.autoPromptEcho = false;
runtimeMock.state.createdSessionIds.push(sessionId);
runtimeMock.state.sessionStatusImplementation = async () => {
if (runtimeMock.state.sessionStatusCalls === 5) {
fifthProbeObserved.resolve(undefined);
}
return { data: {} };
};

const eventsFiber = yield* adapter.streamEvents.pipe(
Stream.filter(
(event) =>
event.threadId === threadId &&
(event.type === "turn.aborted" ||
event.type === "turn.completed" ||
event.type === "runtime.error"),
),
Stream.runForEach((event) => Effect.sync(() => terminalEvents.push(event))),
Effect.forkChild,
);
yield* adapter.startSession({
provider: ProviderDriverKind.make("opencode"),
threadId,
runtimeMode: "full-access",
});
const turn = yield* adapter.sendTurn({
threadId,
input: "Interrupt while admission waits for late evidence",
modelSelection: createModelSelection(
ProviderInstanceId.make("opencode"),
"opencode/kimi-k3",
),
});

for (const delayMs of [250, 500, 1_000, 2_000]) {
yield* advanceTestClock(delayMs);
}
yield* Effect.promise(() => fifthProbeObserved.promise);
yield* advanceTestClock(2_000);
const abortCallsBeforeInterrupt = runtimeMock.state.abortCalls.filter(
(candidate) => candidate === sessionId,
).length;
yield* adapter.interruptTurn(threadId, turn.turnId);
yield* advanceTestClock(1_000);

NodeAssert.equal(
runtimeMock.state.abortCalls.filter((candidate) => candidate === sessionId).length,
abortCallsBeforeInterrupt + 1,
);
NodeAssert.deepEqual(
terminalEvents.map((event) => event.type),
["turn.aborted"],
);

yield* Fiber.interrupt(eventsFiber);
yield* adapter.stopSession(threadId);
}),
);

it.effect("does not fail admission after cleanup abort loses ownership to interruption", () =>
Effect.gen(function* () {
const adapter = yield* OpenCodeAdapter;
const threadId = asThreadId("thread-interrupt-during-admission-cleanup");
const sessionId = "ses_interrupt_during_admission_cleanup";
const fifthProbeObserved = promiseWithResolvers<void>();
const admissionAbortStarted = promiseWithResolvers<void>();
const interruptAbortStarted = promiseWithResolvers<void>();
const releaseAdmissionAbort = promiseWithResolvers<void>();
const terminalEvents: Array<ProviderRuntimeEvent> = [];
let ownedAbortCalls = 0;
runtimeMock.state.autoPromptEcho = false;
runtimeMock.state.createdSessionIds.push(sessionId);
runtimeMock.state.sessionStatusImplementation = async () => {
if (runtimeMock.state.sessionStatusCalls === 5) {
fifthProbeObserved.resolve(undefined);
}
return { data: {} };
};
runtimeMock.state.abortImplementation = async (candidate) => {
if (candidate !== sessionId) return;
ownedAbortCalls += 1;
if (ownedAbortCalls === 1) {
admissionAbortStarted.resolve(undefined);
await releaseAdmissionAbort.promise;
} else if (ownedAbortCalls === 2) {
interruptAbortStarted.resolve(undefined);
}
};

const eventsFiber = yield* adapter.streamEvents.pipe(
Stream.filter(
(event) =>
event.threadId === threadId &&
(event.type === "turn.aborted" ||
event.type === "turn.completed" ||
event.type === "runtime.error"),
),
Stream.runForEach((event) => Effect.sync(() => terminalEvents.push(event))),
Effect.forkChild,
);
yield* adapter.startSession({
provider: ProviderDriverKind.make("opencode"),
threadId,
runtimeMode: "full-access",
});
const turn = yield* adapter.sendTurn({
threadId,
input: "Interrupt while admission cleanup is aborting",
modelSelection: createModelSelection(
ProviderInstanceId.make("opencode"),
"opencode/kimi-k3",
),
});

for (const delayMs of [250, 500, 1_000, 2_000]) {
yield* advanceTestClock(delayMs);
}
yield* Effect.promise(() => fifthProbeObserved.promise);
yield* advanceTestClock(3_000);
yield* Effect.promise(() => admissionAbortStarted.promise);

const interruptFiber = yield* adapter
.interruptTurn(threadId, turn.turnId)
.pipe(Effect.forkChild);
yield* Effect.promise(() => interruptAbortStarted.promise);
releaseAdmissionAbort.resolve(undefined);
yield* Fiber.join(interruptFiber);
yield* Effect.yieldNow;

NodeAssert.equal(ownedAbortCalls, 2);
NodeAssert.deepEqual(
terminalEvents.map((event) => event.type),
["turn.aborted"],
);
const sessions = yield* adapter.listSessions();
const session = sessions.find((candidate) => candidate.threadId === threadId);
NodeAssert.equal(session?.status, "ready");
NodeAssert.equal(session?.activeTurnId, undefined);

runtimeMock.state.abortImplementation = null;
yield* Fiber.interrupt(eventsFiber);
yield* adapter.stopSession(threadId);
}),
);

it.effect("uses polled busy status to admit output after a stopped turn", () =>
Effect.gen(function* () {
const adapter = yield* OpenCodeAdapter;
Expand Down
Loading
Loading