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
33 changes: 33 additions & 0 deletions apps/server/src/codexAppServerManager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -628,6 +628,11 @@ describe("isRecoverableThreadResumeError", () => {
expect(
isRecoverableThreadResumeError(new Error("thread/resume failed: thread not found")),
).toBe(true);
expect(
isRecoverableThreadResumeError(
new Error("thread/resume failed: no rollout found for thread id thread_1"),
),
).toBe(true);
});

it("ignores non-resume errors", () => {
Expand Down Expand Up @@ -1775,6 +1780,34 @@ describe("CodexAppServerManager discovery", () => {
});

describe("thread checkpoint control", () => {
it("starts WebRTC realtime with feature enablement and the compatible protocol", async () => {
const { manager, context, requireSession, sendRequest } = createThreadControlHarness();
sendRequest.mockResolvedValue({});

await manager.startRealtime({
threadId: asThreadId("thread_1"),
sdp: "v=0\r\ns=-\r\n",
voice: "cedar",
});

expect(requireSession).toHaveBeenCalledWith("thread_1");
expect(sendRequest).toHaveBeenNthCalledWith(1, context, "experimentalFeature/enablement/set", {
enablement: {
realtime_conversation: true,
},
});
expect(sendRequest).toHaveBeenNthCalledWith(2, context, "thread/realtime/start", {
threadId: "thread_1",
outputModality: "audio",
version: "v3",
transport: {
type: "webrtc",
sdp: "v=0\r\ns=-\r\n",
},
voice: "cedar",
});
});

it("lists resumable desktop threads with normalized metadata", async () => {
const manager = new CodexAppServerManager();
const context = { discovery: true };
Expand Down
84 changes: 80 additions & 4 deletions apps/server/src/codexAppServerManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import {
ProviderItemId,
type ProviderListModelsResult,
type ProviderListPluginsResult,
type ProviderListRealtimeVoicesResult,
type ProviderMentionReference,
type ProviderPluginAppSummary,
type ProviderPluginDescriptor,
Expand All @@ -32,6 +33,7 @@ import {
type ProviderEvent,
type ProviderSession,
type ProviderSessionStartInput,
type ProviderStartRealtimeInput,
type ProviderTurnStartResult,
RuntimeMode,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -113,6 +115,7 @@ interface CodexSessionContext {
nextRequestId: number;
stopping: boolean;
discovery?: boolean;
realtimeFeatureEnabled?: boolean;
}

interface CodexSkillListInput {
Expand Down Expand Up @@ -188,6 +191,11 @@ export interface CodexAppServerSendTurnInput {
readonly effort?: string;
}

export type CodexAppServerStartRealtimeInput = Pick<
ProviderStartRealtimeInput,
"threadId" | "sdp" | "voice"
>;

type CodexAppServerReviewTarget = ProviderStartReviewInput["target"];

export interface CodexAppServerStartSessionInput {
Expand Down Expand Up @@ -257,6 +265,7 @@ const BENIGN_PROCESS_OUTPUT_REGEXES = [/^(?:\^C)?Token usage:/i];
const RECOVERABLE_THREAD_RESUME_ERROR_SNIPPETS = [
"not found",
"missing thread",
"no rollout found",
"no such thread",
"unknown thread",
"does not exist",
Expand Down Expand Up @@ -449,10 +458,14 @@ function spawnCodexAppServer(input: {
readonly cwd: string;
readonly env: NodeJS.ProcessEnv;
}): ChildProcessWithoutNullStreams {
const prepared = prepareWindowsSafeProcess(input.binaryPath, ["app-server"], {
cwd: input.cwd,
env: input.env,
});
const prepared = prepareWindowsSafeProcess(
input.binaryPath,
["app-server", "--enable", "realtime_conversation"],
{
cwd: input.cwd,
env: input.env,
},
);
return spawn(prepared.command, prepared.args, {
cwd: input.cwd,
env: input.env,
Expand Down Expand Up @@ -722,6 +735,9 @@ export class CodexAppServerManager extends EventEmitter<CodexAppServerManagerEve
await this.sendRequest(context, "initialize", buildCodexInitializeParams());

this.writeMessage(context, { method: "initialized" });
// Codex snapshots feature enablement into each Thread when it is opened.
// Enabling realtime at the later startRealtime call is too late.
await this.enableRealtimeConversation(context);
await this.registerChitauriSkillsRoot(context);
try {
const modelListResponse = await this.sendRequest(context, "model/list", {});
Expand Down Expand Up @@ -1000,6 +1016,41 @@ export class CodexAppServerManager extends EventEmitter<CodexAppServerManagerEve
};
}

async startRealtime(input: CodexAppServerStartRealtimeInput): Promise<void> {
const context = this.requireSession(input.threadId);
const providerThreadId = this.requireProviderThreadId(context);
await this.enableRealtimeConversation(context);
await this.sendRequest(context, "thread/realtime/start", {
threadId: providerThreadId,
outputModality: "audio",
// V2 cannot use WebRTC, while legacy V1 currently routes some accounts
// to the retired quicksilver endpoint. V3 is the supported WebRTC path.
version: "v3",
transport: {
type: "webrtc",
sdp: input.sdp,
},
...(input.voice !== undefined ? { voice: input.voice } : {}),
});
}

async stopRealtime(threadId: ThreadId): Promise<void> {
const context = this.requireSession(threadId);
await this.sendRequest(context, "thread/realtime/stop", {
threadId: this.requireProviderThreadId(context),
});
}

async listRealtimeVoices(threadId: ThreadId): Promise<ProviderListRealtimeVoicesResult> {
const context = this.requireSession(threadId);
await this.enableRealtimeConversation(context);
return this.sendRequest<ProviderListRealtimeVoicesResult>(
context,
"thread/realtime/listVoices",
{},
);
}

async steerTurn(input: CodexAppServerSendTurnInput): Promise<ProviderTurnStartResult> {
const context = this.requireSession(input.threadId);
context.collabReceiverTurns.clear();
Expand Down Expand Up @@ -1352,6 +1403,7 @@ export class CodexAppServerManager extends EventEmitter<CodexAppServerManagerEve

await this.sendRequest(context, "initialize", buildCodexInitializeParams());
this.writeMessage(context, { method: "initialized" });
await this.enableRealtimeConversation(context);
await this.registerChitauriSkillsRoot(context);
try {
const accountReadResponse = await this.sendRequest(context, "account/read", {});
Expand Down Expand Up @@ -2392,6 +2444,30 @@ export class CodexAppServerManager extends EventEmitter<CodexAppServerManagerEve
return result as TResponse;
}

private requireProviderThreadId(context: CodexSessionContext): string {
const providerThreadId = readResumeThreadId({
threadId: context.session.threadId,
runtimeMode: context.session.runtimeMode,
resumeCursor: context.session.resumeCursor,
});
if (!providerThreadId) {
throw new Error("Session is missing provider resume thread id.");
}
return providerThreadId;
}

private async enableRealtimeConversation(context: CodexSessionContext): Promise<void> {
if (context.realtimeFeatureEnabled) {
return;
}
await this.sendRequest(context, "experimentalFeature/enablement/set", {
enablement: {
realtime_conversation: true,
},
});
context.realtimeFeatureEnabled = true;
}

private writeMessage(context: CodexSessionContext, message: unknown): void {
const encoded = JSON.stringify(message);
if (!context.child.stdin.writable) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -103,7 +103,11 @@ function createProviderServiceHarness(
getCapabilities: () => Effect.succeed({ sessionModelSwitch: "in-session" }),
rollbackConversation,
compactThread: () => unsupported(),
startRealtime: () => unsupported(),
stopRealtime: () => unsupported(),
listRealtimeVoices: () => unsupported(),
streamEvents: Stream.fromPubSub(runtimeEventPubSub),
streamRealtimeEvents: Stream.empty,
};

const emit = (event: LegacyProviderRuntimeEvent): void => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ describe("OrchestrationReactor", () => {
),
Layer.provideMerge(
Layer.succeed(ProviderCommandReactor, {
ensureSession: (threadId) => Effect.succeed(threadId),
start: Effect.sync(() => {
started.push("provider-command-reactor");
}),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -326,7 +326,11 @@ describe("ProviderCommandReactor", () => {
}),
rollbackConversation,
compactThread: () => unsupported(),
startRealtime: () => unsupported(),
stopRealtime: () => unsupported(),
listRealtimeVoices: () => unsupported(),
streamEvents: Stream.fromPubSub(runtimeEventPubSub),
streamRealtimeEvents: Stream.empty,
};

const orchestrationLayer = OrchestrationEngineLive.pipe(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2193,6 +2193,9 @@ const make = Effect.gen(function* () {

const worker = yield* makeDrainableWorker(processDomainEventSafely);

const ensureSession: ProviderCommandReactorShape["ensureSession"] = (threadId) =>
ensureSessionForThread(threadId, new Date().toISOString());

const start: ProviderCommandReactorShape["start"] = Effect.all([
Stream.runForEach(orchestrationEngine.streamDomainEvents, (event) => {
if (
Expand Down Expand Up @@ -2221,6 +2224,7 @@ const make = Effect.gen(function* () {
]).pipe(Effect.asVoid);

return {
ensureSession,
start,
drain: worker.drain,
} satisfies ProviderCommandReactorShape;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,11 @@ function createProviderServiceHarness() {
}),
rollbackConversation: () => unsupported(),
compactThread: () => unsupported(),
startRealtime: () => unsupported(),
stopRealtime: () => unsupported(),
listRealtimeVoices: () => unsupported(),
streamEvents: Stream.fromPubSub(runtimeEventPubSub),
streamRealtimeEvents: Stream.empty,
};

const setSession = (session: ProviderSession): void => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,21 @@
*
* @module ProviderCommandReactor
*/
import type { ThreadId } from "@t3tools/contracts";
import { ServiceMap } from "effect";
import type { Effect, Scope } from "effect";

/**
* ProviderCommandReactorShape - Service API for provider command reactors.
*/
export interface ProviderCommandReactorShape {
/**
* Establish or recover the canonical provider session for a projected Thread.
* Non-turn provider features must use this instead of inventing a partial
* session startup path.
*/
readonly ensureSession: (threadId: ThreadId) => Effect.Effect<ThreadId, unknown>;

/**
* Start reacting to provider-intent orchestration domain events.
*
Expand Down
53 changes: 53 additions & 0 deletions apps/server/src/orchestration/workerToolsMcp.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -157,6 +157,59 @@ describe("Worker MCP tools", () => {
});
expect(snapshot.threads).toHaveLength(2);

const delegated = (await call(threadA, "threads_create", {
title: "Claude implementation",
provider: "claudeAgent",
model: "claude-sonnet-4-5",
prompt: "Implement the isolated change and report the result.",
})) as {
id: string;
provider: string;
parentThreadId: string;
dispatched: boolean;
};
expect(delegated).toMatchObject({
provider: "claudeAgent",
parentThreadId: threadA,
dispatched: true,
});

const threads = (await call(threadA, "threads_list", {})) as Array<{
id: string;
provider: string;
}>;
expect(threads).toEqual(
expect.arrayContaining([
expect.objectContaining({ id: delegated.id, provider: "claudeAgent" }),
]),
);

const delegatedRead = (await call(threadA, "threads_read", {
thread_id: delegated.id,
})) as {
messages: Array<{ role: string; text: string }>;
};
expect(delegatedRead.messages).toEqual([
expect.objectContaining({
role: "user",
text: "Implement the isolated change and report the result.",
}),
]);

await call(threadA, "threads_send", {
thread_id: delegated.id,
prompt: "Also run the focused tests.",
});
const delegatedReadAfterSend = (await call(threadA, "threads_read", {
thread_id: delegated.id,
message_limit: 1,
})) as {
messages: Array<{ role: string; text: string }>;
};
expect(delegatedReadAfterSend.messages).toEqual([
expect.objectContaining({ role: "user", text: "Also run the focused tests." }),
]);

await system.runtime.dispose();
});
});
Loading
Loading