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
51 changes: 51 additions & 0 deletions apps/desktop/src/backend/DesktopServerExposure.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,9 @@ import * as NodeFileSystem from "@effect/platform-node/NodeFileSystem";
import * as NodeHttpClient from "@effect/platform-node/NodeHttpClient";
import * as NodeServices from "@effect/platform-node/NodeServices";
import { assert, describe, it } from "@effect/vitest";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Sink from "effect/Sink";
Expand Down Expand Up @@ -305,6 +307,55 @@ describe("DesktopServerExposure", () => {
);
});

it.effect("keeps a Tailscale Serve change made while a mode change is saving", () =>
Effect.gen(function* () {
const modeWriteStarted = yield* Deferred.make<void>();
const releaseModeWrite = yield* Deferred.make<void>();
const settingsLayer = Layer.effect(
DesktopAppSettings.DesktopAppSettings,
Effect.gen(function* () {
const settings = yield* DesktopAppSettings.DesktopAppSettings;
return DesktopAppSettings.DesktopAppSettings.of({
...settings,
// Hold the mode write the way a slow disk would.
setServerExposureMode: (mode) =>
Deferred.succeed(modeWriteStarted, undefined).pipe(
Effect.andThen(Deferred.await(releaseModeWrite)),
Effect.andThen(settings.setServerExposureMode(mode)),
),
});
}),
).pipe(Layer.provide(DesktopAppSettings.layerTest()));

return yield* withHarness(
lanNetworkInterfaces,
Effect.gen(function* () {
const serverExposure = yield* DesktopServerExposure.DesktopServerExposure;
yield* serverExposure.configureFromSettings({ port: 4173 });

const modeChange = yield* serverExposure
.setMode("network-accessible")
.pipe(Effect.forkChild({ startImmediately: true }));
yield* Deferred.await(modeWriteStarted);
const tailscaleChange = yield* serverExposure
.setTailscaleServeEnabled({ enabled: true, port: 8443 })
.pipe(Effect.forkChild({ startImmediately: true }));
yield* Deferred.succeed(releaseModeWrite, undefined);
yield* Fiber.join(modeChange);
yield* Fiber.join(tailscaleChange);

const state = yield* serverExposure.getState;
assert.equal(state.mode, "network-accessible");
assert.equal(state.tailscaleServeEnabled, true);
assert.equal(state.tailscaleServePort, 8443);
}),
{},
undefined,
settingsLayer,
);
}),
);

it.effect("keeps LAN and Tailscale endpoints distinct when Tailscale is enumerated first", () =>
withHarness(
{ ...tailnetNetworkInterfaces, ...lanNetworkInterfaces },
Expand Down
8 changes: 7 additions & 1 deletion apps/desktop/src/backend/DesktopServerExposure.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
import * as Semaphore from "effect/Semaphore";
import * as HttpClient from "effect/http/HttpClient";
import * as ChildProcessSpawner from "effect/process/ChildProcessSpawner";

Expand Down Expand Up @@ -419,6 +420,10 @@ export const make = Effect.gen(function* () {
const httpClient = yield* HttpClient.HttpClient;
const desktopSettings = yield* DesktopAppSettings.DesktopAppSettings;
const stateRef = yield* Ref.make(initialRuntimeState());
// Each change reads the runtime state, persists settings, then writes the
// state back. Run them one at a time so a change never writes over another
// with values it read before that change persisted.
const changePermit = yield* Semaphore.make(1);

// Cache the `tailscale status` spawn for the TTL. On macOS, the Mac App
// Store Tailscale CLI lives inside Tailscale's sandbox container, so each
Expand Down Expand Up @@ -492,7 +497,7 @@ export const make = Effect.gen(function* () {
state: toContractState(resolved.state),
requiresRelaunch: change.changed || requiresBackendRelaunch(previous, resolved.state),
};
});
}, changePermit.withPermit);

const setTailscaleServeEnabled = Effect.fn("desktop.serverExposure.setTailscaleServeEnabled")(
function* (input: { readonly enabled: boolean; readonly port?: number }) {
Expand Down Expand Up @@ -527,6 +532,7 @@ export const make = Effect.gen(function* () {
requiresRelaunch: result.changed,
};
},
changePermit.withPermit,
);

const getAdvertisedEndpoints = Effect.gen(function* () {
Expand Down
187 changes: 187 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1656,6 +1656,193 @@ function makeCodexReplayTranscript(input: {
};
}

function withReplayRequestId(
entry: CodexReplay.CodexAppServerReplayEntry,
id: number,
): CodexReplay.CodexAppServerReplayEntry {
return entry.type !== "runtime_exit" && Predicate.isObject(entry.frame)
? { ...entry, frame: { ...entry.frame, id } }
: entry;
}

describe("CodexAdapterV2 session initialize", () => {
const openReplaySession = (
transcript: CodexReplay.CodexAppServerReplayTranscript,
beforeEmitInbound?: CodexReplay.CodexAppServerReplayDriver["beforeEmitInbound"],
) =>
Effect.gen(function* () {
const driver = yield* CodexReplay.makeReplayDriver(
transcript,
beforeEmitInbound === undefined ? {} : { beforeEmitInbound },
);
let initializeRequests = 0;
const adapter = CodexAdapterV2.makeCodexAdapterV2({
instanceId: CodexAdapterV2.CODEX_DEFAULT_INSTANCE_ID,
settings: DEFAULT_CODEX_SETTINGS,
environment: {},
clientFactory: {
open: (openInput) =>
Layer.build(CodexReplay.layerReplayWithDriver(driver)).pipe(
Effect.flatMap((context) =>
Effect.service(CodexClient.CodexAppServerClient).pipe(Effect.provide(context)),
),
Effect.map(
(client) =>
({
...client,
request: (method, params) =>
Effect.sync(() => {
if (method === "initialize") initializeRequests++;
}).pipe(Effect.andThen(client.request(method, params))),
}) satisfies CodexClient.CodexAppServerClient["Service"],
),
Effect.mapError(
(cause) =>
new ProviderAdapterOpenSessionError({
driver: CodexAdapterV2.CODEX_DRIVER_KIND,
providerSessionId: openInput.providerSessionId,
cause,
}),
),
),
},
fileSystem: yield* FileSystem.FileSystem,
idAllocator: yield* IdAllocator.IdAllocatorV2,
serverConfig: yield* makeReplayServerConfig(transcript.scenario).pipe(Effect.orDie),
});
const runtime = yield* adapter.openSession({
threadId: ThreadId.make(`thread-${transcript.scenario}`),
providerSessionId: ProviderSessionId.make(`provider-session-${transcript.scenario}`),
modelSelection: CODEX_TEST_MODEL_SELECTION,
runtimePolicy: CODEX_TEST_RUNTIME_POLICY,
});
return {
ensureThread: (threadId: string) =>
runtime.ensureThread({
threadId: ThreadId.make(threadId),
modelSelection: CODEX_TEST_MODEL_SELECTION,
runtimePolicy: CODEX_TEST_RUNTIME_POLICY,
}),
initializeRequests: () => initializeRequests,
};
});

const replayPreamble = (nativeThreadId: string) =>
codexReplayPreamble({ nativeThreadId, nativeTurnId: "unused", prompt: "unused" });

it.effect("sends one initialize when two threads start on a fresh session at once", () =>
Effect.gen(function* () {
const initializeAwaitingResponse = yield* Deferred.make<void>();
const releaseInitialize = yield* Deferred.make<void>();
// The transcript allows exactly one handshake: a second `initialize`
// frame fails the replay, as Codex rejects it with "Already initialized".
const session = yield* openReplaySession(
makeCodexReplayTranscript({
scenario: "concurrent-initialize",
entries: [
...replayPreamble("concurrent-first").slice(0, 5),
...replayPreamble("concurrent-second")
.slice(3, 5)
.map((entry) => withReplayRequestId(entry, 3)),
],
}),
(entry) =>
entry.label === "initialize"
? Deferred.succeed(initializeAwaitingResponse, undefined).pipe(
Effect.andThen(Deferred.await(releaseInitialize)),
)
: Effect.void,
);

const first = yield* session
.ensureThread("thread-concurrent-first")
.pipe(Effect.forkChild({ startImmediately: true }));
yield* Deferred.await(initializeAwaitingResponse);
// The second thread arrives while the handshake is still unanswered.
const second = yield* session
.ensureThread("thread-concurrent-second")
.pipe(Effect.forkChild({ startImmediately: true }));
yield* Deferred.succeed(releaseInitialize, undefined);

const providerThreads = [yield* Fiber.join(first), yield* Fiber.join(second)];
assert.equal(session.initializeRequests(), 1);
assert.sameMembers(
providerThreads.map((providerThread) => providerThread.nativeThreadRef?.nativeId),
["concurrent-first", "concurrent-second"],
);
}).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))),
);

it.effect("retries initialize after a failed handshake", () =>
Effect.gen(function* () {
const preamble = replayPreamble("initialize-retry");
const entries: Array<CodexReplay.CodexAppServerReplayEntry> = [
...preamble.slice(0, 1),
{
type: "emit_inbound",
label: "initialize",
frame: { id: 1, error: { code: -32603, message: "Codex is not ready." } },
},
...preamble.slice(0, 2).map((entry) => withReplayRequestId(entry, 2)),
...preamble.slice(2, 3),
...preamble.slice(3, 5).map((entry) => withReplayRequestId(entry, 3)),
];
const session = yield* openReplaySession(
makeCodexReplayTranscript({ scenario: "initialize-retry", entries }),
);

const failure = yield* session.ensureThread("thread-initialize-retry").pipe(Effect.flip);
assert.equal(failure._tag, "ProviderAdapterEnsureThreadError");
const providerThread = yield* session.ensureThread("thread-initialize-retry");
assert.equal(providerThread.nativeThreadRef?.nativeId, "initialize-retry");
assert.equal(session.initializeRequests(), 2);
}).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))),
);

it.effect("completes the handshake after a caller is interrupted mid-initialize", () =>
Effect.gen(function* () {
const initializeAwaitingResponse = yield* Deferred.make<void>();
const releaseInitialize = yield* Deferred.make<void>();
const preamble = replayPreamble("initialize-interrupted");
const session = yield* openReplaySession(
makeCodexReplayTranscript({
scenario: "initialize-interrupted",
entries: [
...preamble.slice(0, 2),
// Codex handled the interrupted caller's `initialize`, so it
// rejects the next one.
...preamble.slice(0, 1).map((entry) => withReplayRequestId(entry, 2)),
{
type: "emit_inbound",
label: "initialize-rejected",
frame: { id: 2, error: { code: -32600, message: "Already initialized" } },
},
...preamble.slice(2, 3),
...preamble.slice(3, 5).map((entry) => withReplayRequestId(entry, 3)),
],
}),
(entry) =>
entry.label === "initialize"
? Deferred.succeed(initializeAwaitingResponse, undefined).pipe(
Effect.andThen(Deferred.await(releaseInitialize)),
)
: Effect.void,
);

const interrupted = yield* session
.ensureThread("thread-initialize-interrupted")
.pipe(Effect.forkChild({ startImmediately: true }));
yield* Deferred.await(initializeAwaitingResponse);
yield* Fiber.interrupt(interrupted);
yield* Deferred.succeed(releaseInitialize, undefined);

const providerThread = yield* session.ensureThread("thread-initialize-interrupted");
assert.equal(providerThread.nativeThreadRef?.nativeId, "initialize-interrupted");
assert.equal(session.initializeRequests(), 2);
}).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))),
);
});

describe("CodexAdapterV2 post-settle continuation", () => {
const awaitUntil = (predicate: () => boolean, label: string): Effect.Effect<void> =>
Effect.gen(function* () {
Expand Down
46 changes: 32 additions & 14 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1687,21 +1687,39 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
),
);
const initialized = yield* Ref.make(false);
const ensureInitialized = Effect.gen(function* () {
const alreadyInitialized = yield* Ref.get(initialized);
if (alreadyInitialized) {
return;
}
// Threads share this app-server, and Codex rejects a second
// `initialize`. Callers wait for an in-flight handshake instead of
// starting their own; a failed handshake leaves the flag unset so the
// next caller retries.
const initializePermit = yield* Semaphore.make(1);
const ensureInitialized = initializePermit.withPermit(
Effect.gen(function* () {
const alreadyInitialized = yield* Ref.get(initialized);
if (alreadyInitialized) {
return;
}

yield* client.request("initialize", {
// Codex uses the client name as the request originator, so sessions
// identify themselves exactly like the provider probe.
clientInfo: buildCodexInitializeParams().clientInfo,
capabilities: CODEX_CLIENT_CAPABILITIES,
});
yield* client.notify("initialized", undefined);
yield* Ref.set(initialized, true);
});
yield* client
.request("initialize", {
// Codex uses the client name as the request originator, so sessions
// identify themselves exactly like the provider probe.
clientInfo: buildCodexInitializeParams().clientInfo,
capabilities: CODEX_CLIENT_CAPABILITIES,
})
.pipe(
Effect.catchTags({
// A caller interrupted after its `initialize` reached Codex
// leaves the app-server initialized but the flag unset.
CodexAppServerRequestError: (error) =>
error.code === -32600 && error.errorMessage === "Already initialized"
? Effect.void
: Effect.fail(error),
}),
);
yield* client.notify("initialized", undefined);
yield* Ref.set(initialized, true);
}),
);
const now = yield* DateTime.now;
const session = providerSession({
providerSessionId: input.providerSessionId,
Expand Down
Loading