Skip to content
Closed
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
51 changes: 51 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import * as NodeServices from "@effect/platform-node/NodeServices";
import { it, vi } from "@effect/vitest";

import * as Context from "effect/Context";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Fiber from "effect/Fiber";
Expand Down Expand Up @@ -2562,6 +2563,56 @@ const scopedLifecycleLayer = it.layer(
);

scopedLifecycleLayer("CodexAdapterLive scoped lifecycle", (it) => {
it.effect(
"releases a stalled startup so the next session can start",
() =>
Effect.gen(function* () {
const starting = yield* Deferred.make<void>();
const runtimeFactory = makeScopedRuntimeFactory();
const stalledThreadId = asThreadId("thread-stalled-startup");
const nextThreadId = asThreadId("thread-after-stalled-startup");
const adapter = yield* makeCodexAdapter(decodeCodexSettings({}), {
makeRuntime: (options) =>
runtimeFactory.factory(options).pipe(
Effect.map((runtime) => {
if (options.threadId === stalledThreadId) {
runtime.start = () =>
Deferred.succeed(starting, undefined).pipe(Effect.andThen(Effect.never));
}
return runtime;
}),
),
});
const recovery = yield* Effect.gen(function* () {
const failed = yield* adapter
.startSession({ threadId: stalledThreadId, runtimeMode: "full-access" })
.pipe(Effect.result);
const started = yield* adapter.startSession({
threadId: nextThreadId,
runtimeMode: "full-access",
});
return { failed, started };
}).pipe(Effect.forkChild);

yield* Deferred.await(starting);
yield* TestClock.adjust("30 seconds");
const { failed, started } = yield* Fiber.join(recovery);

NodeAssert.equal(failed._tag, "Failure");
NodeAssert.equal(failed.failure._tag, "ProviderAdapterProcessError");
NodeAssert.match(
failed.failure.message,
/Codex session startup timed out after 30 seconds/,
);
NodeAssert.deepStrictEqual(runtimeFactory.releasedThreadIds, [stalledThreadId]);
NodeAssert.equal(yield* adapter.hasSession(stalledThreadId), false);
NodeAssert.equal(started.threadId, nextThreadId);
NodeAssert.equal(started.status, "ready");
yield* adapter.stopSession(nextThreadId);
}),
{ timeout: 5_000 },
);

it.effect("closes the externally owned session scope on stopSession", () =>
Effect.gen(function* () {
scopedLifecycleRuntimeFactory.releasedThreadIds.length = 0;
Expand Down
13 changes: 13 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2453,6 +2453,19 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
cause,
}),
),
// A stalled handshake must release the shared provider command worker.
Effect.timeoutOrElse({
duration: "30 seconds",
orElse: () =>
Effect.fail(
new ProviderAdapterProcessError({
provider: PROVIDER,
threadId: input.threadId,
detail:
"Codex session startup timed out after 30 seconds. Try starting the session again.",
}),
),
}),
Effect.onError(() =>
runtime.close.pipe(
Effect.andThen(Effect.ignore(Scope.close(sessionScope, Exit.void))),
Expand Down
Loading