diff --git a/apps/desktop/src/backend/DesktopServerExposure.test.ts b/apps/desktop/src/backend/DesktopServerExposure.test.ts index b9c460825561..b1312e81def9 100644 --- a/apps/desktop/src/backend/DesktopServerExposure.test.ts +++ b/apps/desktop/src/backend/DesktopServerExposure.test.ts @@ -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"; @@ -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(); + const releaseModeWrite = yield* Deferred.make(); + 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 }, diff --git a/apps/desktop/src/backend/DesktopServerExposure.ts b/apps/desktop/src/backend/DesktopServerExposure.ts index 8deafeefacd8..e5bd11035f5c 100644 --- a/apps/desktop/src/backend/DesktopServerExposure.ts +++ b/apps/desktop/src/backend/DesktopServerExposure.ts @@ -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"; @@ -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 @@ -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 }) { @@ -527,6 +532,7 @@ export const make = Effect.gen(function* () { requiresRelaunch: result.changed, }; }, + changePermit.withPermit, ); const getAdvertisedEndpoints = Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts index 89f5b7ae1a26..32a73e410497 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts @@ -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(); + const releaseInitialize = yield* Deferred.make(); + // 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 = [ + ...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(); + const releaseInitialize = yield* Deferred.make(); + 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 => Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 6ef64f685bf3..6a80238561b7 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -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,