Skip to content
Closed
Show file tree
Hide file tree
Changes from 2 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
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,8 @@ const startupDependencies = Layer.mergeAll(
getInstanceInfo: () => Effect.die("unused"),
rollbackConversation: () => Effect.die("unused"),
uploadFeedback: () => Effect.die("unused"),
startRealtimeVoice: () => Effect.die("unused"),
stopRealtimeVoice: () => Effect.die("unused"),
streamEvents: Stream.empty,
}),
);
Expand Down
9 changes: 9 additions & 0 deletions apps/server/src/auth/RpcAuthorization.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,15 @@ describe("RPC authorization scopes", () => {
);
});

it("requires permission to operate on a thread for realtime voice", () => {
expect(requiredScopeForRpcMethod(WS_METHODS.providerRealtimeVoiceStart)).toBe(
AuthOrchestrationOperateScope,
);
expect(requiredScopeForRpcMethod(WS_METHODS.providerRealtimeVoiceStop)).toBe(
AuthOrchestrationOperateScope,
);
});

it("reads the reviewer menu under the same scope as the pull request it belongs to", () => {
// The candidate list is a read like the detail beside it, and asking somebody for a review is
// a write like every other pull request operation.
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/auth/RpcAuthorization.ts
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,8 @@ export const RPC_REQUIRED_SCOPES = {
[WS_METHODS.attachmentsCreateUploadUrl]: AuthOrchestrationOperateScope,
[WS_METHODS.attachmentsDelete]: AuthOrchestrationOperateScope,
[WS_METHODS.providerUploadFeedback]: AuthOrchestrationOperateScope,
[WS_METHODS.providerRealtimeVoiceStart]: AuthOrchestrationOperateScope,
[WS_METHODS.providerRealtimeVoiceStop]: AuthOrchestrationOperateScope,
[WS_METHODS.subscribeVcsStatus]: AuthOrchestrationReadScope,
[WS_METHODS.subscribeResourceTelemetry]: AuthOrchestrationReadScope,
[WS_METHODS.vcsRefreshStatus]: AuthOrchestrationReadScope,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,8 @@ function createProviderServiceHarness(
}),
rollbackConversation,
uploadFeedback: () => unsupported(),
startRealtimeVoice: () => unsupported(),
stopRealtimeVoice: () => unsupported(),
get streamEvents() {
return Stream.fromPubSub(runtimeEventPubSub);
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -352,6 +352,8 @@ describe("ProviderCommandReactor", () => {
},
rollbackConversation: () => unsupported(),
uploadFeedback: () => unsupported(),
startRealtimeVoice: () => unsupported(),
stopRealtimeVoice: () => unsupported(),
get streamEvents() {
return Stream.fromPubSub(runtimeEventPubSub);
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,8 @@ function createProviderServiceHarness() {
},
rollbackConversation: () => unsupported(),
uploadFeedback: () => unsupported(),
startRealtimeVoice: () => unsupported(),
stopRealtimeVoice: () => unsupported(),
get streamEvents() {
return Stream.fromPubSub(runtimeEventPubSub);
},
Expand Down
31 changes: 31 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,10 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape {
Promise.resolve({ threadId: "provider-thread-1" }),
);

public readonly startRealtimeVoiceImpl = vi.fn((sdp: string) => Promise.resolve(sdp));

public readonly stopRealtimeVoiceImpl = vi.fn(() => Promise.resolve(undefined));

public readonly respondToRequestImpl = vi.fn(
(_requestId: ApprovalRequestId, _decision: ProviderApprovalDecision): Promise<void> =>
Promise.resolve(undefined),
Expand Down Expand Up @@ -150,6 +154,12 @@ class FakeCodexRuntime implements CodexSessionRuntimeShape {
return Effect.promise(() => this.uploadFeedbackImpl(reason));
}

startRealtimeVoice(sdp: string) {
return Effect.promise(() => this.startRealtimeVoiceImpl(sdp));
}

stopRealtimeVoice = Effect.promise(() => this.stopRealtimeVoiceImpl());

respondToRequest(requestId: ApprovalRequestId, decision: ProviderApprovalDecision) {
return Effect.promise(() => this.respondToRequestImpl(requestId, decision));
}
Expand Down Expand Up @@ -372,6 +382,27 @@ sessionErrorLayer("CodexAdapterLive session errors", (it) => {
}),
);

it.effect("routes realtime voice signaling through the active Codex runtime", () =>
Effect.gen(function* () {
const adapter = yield* CodexAdapter;
const threadId = asThreadId("thread-voice");
yield* adapter.startSession({
provider: ProviderDriverKind.make("codex"),
threadId,
runtimeMode: "full-access",
});
const runtime = sessionRuntimeFactory.lastRuntime;
NodeAssert.ok(runtime);

const result = yield* adapter.startRealtimeVoice({ threadId, sdp: "offer-sdp" });
yield* adapter.stopRealtimeVoice({ threadId });

NodeAssert.deepStrictEqual(result, { sdp: "offer-sdp" });
NodeAssert.deepStrictEqual(runtime.startRealtimeVoiceImpl.mock.calls, [["offer-sdp"]]);
NodeAssert.equal(runtime.stopRealtimeVoiceImpl.mock.calls.length, 1);
}),
);

it.effect("maps codex model options before sending a turn", () =>
Effect.gen(function* () {
const adapter = yield* CodexAdapter;
Expand Down
23 changes: 23 additions & 0 deletions apps/server/src/provider/Layers/CodexAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1870,6 +1870,27 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
),
);

const startRealtimeVoice: CodexAdapterShape["startRealtimeVoice"] = (input) =>
requireSession(input.threadId).pipe(
Effect.flatMap((session) => session.runtime.startRealtimeVoice(input.sdp)),
Effect.map((sdp) => ({ sdp })),
Effect.mapError((cause) =>
cause._tag === "ProviderAdapterSessionNotFoundError"
? cause
: mapCodexRuntimeError(input.threadId, "thread/realtime/start", cause),
),
);

const stopRealtimeVoice: CodexAdapterShape["stopRealtimeVoice"] = (input) =>
requireSession(input.threadId).pipe(
Effect.flatMap((session) => session.runtime.stopRealtimeVoice),
Effect.mapError((cause) =>
cause._tag === "ProviderAdapterSessionNotFoundError"
? cause
: mapCodexRuntimeError(input.threadId, "thread/realtime/stop", cause),
),
);

const readThread: CodexAdapterShape["readThread"] = (threadId) =>
requireSession(threadId).pipe(
Effect.flatMap((session) => session.runtime.readThread),
Expand Down Expand Up @@ -2005,6 +2026,8 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
startSession,
sendTurn,
interruptTurn,
startRealtimeVoice,
stopRealtimeVoice,
readThread,
rollbackThread,
uploadFeedback,
Expand Down
12 changes: 12 additions & 0 deletions apps/server/src/provider/Layers/CodexSessionRuntime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import {
} from "../CodexDeveloperInstructions.ts";
import { codexSessionAppServerArgs } from "./codexLaunchArgs.ts";
import {
buildCodexRealtimeStartParams,
buildTurnStartParams,
describeMcpElicitation,
hasConfiguredMcpServer,
Expand All @@ -26,6 +27,17 @@ import {
} from "./CodexSessionRuntime.ts";
const isCodexAppServerRequestError = Schema.is(CodexErrors.CodexAppServerRequestError);

describe("buildCodexRealtimeStartParams", () => {
it("uses Codex's v3 WebRTC audio transport", () => {
NodeAssert.deepStrictEqual(buildCodexRealtimeStartParams("provider-thread-1", "offer-sdp"), {
threadId: "provider-thread-1",
outputModality: "audio",
version: "v3",
transport: { type: "webrtc", sdp: "offer-sdp" },
});
});
});

describe("CodexSessionRuntimeIdentifierGenerationError", () => {
it("retains identifier purpose and the random source failure", () => {
const cause = new Error("random source unavailable");
Expand Down
72 changes: 72 additions & 0 deletions apps/server/src/provider/Layers/CodexSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Queue from "effect/Queue";
import * as Ref from "effect/Ref";
import * as Schema from "effect/Schema";
Expand Down Expand Up @@ -191,6 +192,8 @@ export interface CodexSessionRuntimeShape {
input: CodexSessionRuntimeSendTurnInput,
) => Effect.Effect<ProviderTurnStartResult, CodexSessionRuntimeError>;
readonly interruptTurn: (turnId?: TurnId) => Effect.Effect<void, CodexSessionRuntimeError>;
readonly startRealtimeVoice: (sdp: string) => Effect.Effect<string, CodexSessionRuntimeError>;
readonly stopRealtimeVoice: Effect.Effect<void, CodexSessionRuntimeError>;
readonly readThread: Effect.Effect<CodexThreadSnapshot, CodexSessionRuntimeError>;
readonly rollbackThread: (
numTurns: number,
Expand All @@ -210,6 +213,15 @@ export interface CodexSessionRuntimeShape {
readonly close: Effect.Effect<void>;
}

export function buildCodexRealtimeStartParams(threadId: string, sdp: string) {
return {
threadId,
outputModality: "audio",
version: "v3",
transport: { type: "webrtc", sdp },
} as const;
}

export type CodexSessionRuntimeError =
| CodexErrors.CodexAppServerError
| CodexSessionRuntimePendingApprovalNotFoundError
Expand Down Expand Up @@ -1130,6 +1142,7 @@ export const makeCodexSessionRuntime = (
const collabChildLiveTurnsRef = yield* Ref.make(new Map<string, string>());
const suppressMemoryConsolidationNotification = makeMemoryConsolidationNotificationFilter();
const closedRef = yield* Ref.make(false);
const pendingRealtimeSdpRef = yield* Ref.make<Deferred.Deferred<string> | null>(null);

// `~` is not shell-expanded when env vars are set via
// `child_process.spawn`; `expandHomePath` lets a configured
Expand Down Expand Up @@ -1712,6 +1725,23 @@ export const makeCodexSessionRuntime = (
),
);

yield* client.handleServerNotification("thread/realtime/sdp", (payload) =>
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
currentSessionProviderThreadId.pipe(
Effect.flatMap((providerThreadId) => {
if (providerThreadId && payload.threadId !== providerThreadId) {
return Effect.void;
}
return Ref.get(pendingRealtimeSdpRef).pipe(
Effect.flatMap((pending) =>
pending === null
? Effect.void
: Deferred.succeed(pending, payload.sdp).pipe(Effect.asVoid),
),
);
}),
),
);

yield* client.handleServerRequest("item/commandExecution/requestApproval", (payload) =>
Effect.gen(function* () {
const requestId = ApprovalRequestId.make(yield* randomUUIDv4("command-approval-request"));
Expand Down Expand Up @@ -2073,6 +2103,12 @@ export const makeCodexSessionRuntime = (
}
yield* settlePendingApprovals("cancel");
yield* settlePendingUserInputs({});
const providerThreadId = currentProviderThreadId(yield* Ref.get(sessionRef));
if (providerThreadId) {
yield* client.raw
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
Outdated
.request("thread/realtime/stop", { threadId: providerThreadId })
.pipe(Effect.ignore);
}
Comment thread
PengLx marked this conversation as resolved.
yield* updateSession(sessionRef, {
status: "closed",
activeTurnId: undefined,
Expand Down Expand Up @@ -2180,6 +2216,42 @@ export const makeCodexSessionRuntime = (
turnId: effectiveTurnId,
});
}),
startRealtimeVoice: (sdp) =>
Effect.gen(function* () {
const providerThreadId = yield* readProviderThreadId;
const pending = yield* Deferred.make<string>();
const installed = yield* Ref.modify(pendingRealtimeSdpRef, (current) =>
current === null ? ([true, pending] as const) : ([false, current] as const),
);
if (!installed) {
return yield* CodexErrors.CodexAppServerRequestError.invalidRequest(
"A Codex realtime voice session is already connecting.",
{ method: "thread/realtime/start" },
);
}
Comment thread
macroscopeapp[bot] marked this conversation as resolved.

return yield* Effect.gen(function* () {
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
Outdated
yield* client.raw.request(
"thread/realtime/start",
buildCodexRealtimeStartParams(providerThreadId, sdp),
);
const answer = yield* Deferred.await(pending).pipe(Effect.timeoutOption("20 seconds"));
if (Option.isNone(answer)) {
yield* client.raw
.request("thread/realtime/stop", { threadId: providerThreadId })
.pipe(Effect.ignore);
return yield* CodexErrors.CodexAppServerRequestError.invalidRequest(
"Codex realtime voice did not return an SDP answer.",
{ method: "thread/realtime/start" },
);
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
Outdated
}
return answer.value;
}).pipe(Effect.ensuring(Ref.set(pendingRealtimeSdpRef, null)));
Comment thread
PengLx marked this conversation as resolved.
Outdated
}),
stopRealtimeVoice: Effect.gen(function* () {
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
const providerThreadId = yield* readProviderThreadId;
yield* client.raw.request("thread/realtime/stop", { threadId: providerThreadId });
}),
readThread: Effect.gen(function* () {
const providerThreadId = yield* readProviderThreadId;
const response = yield* client.request("thread/read", {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@ const fakeCodexAdapter: CodexAdapter.CodexAdapterShape = {
readThread: vi.fn(),
rollbackThread: vi.fn(),
uploadFeedback: vi.fn(),
startRealtimeVoice: vi.fn(),
stopRealtimeVoice: vi.fn(),
stopAll: vi.fn(),
streamEvents: Stream.empty,
};
Expand Down
Loading
Loading