Skip to content
Open
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
32 changes: 32 additions & 0 deletions src/CodexAcpClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -379,6 +379,38 @@ export class CodexAcpClient {
};
}

async forkSession(
request: acp.ForkSessionRequest,
onSubscribed: (sessionId: string) => void
): Promise<SessionMetadata> {
const additionalDirectories = readAdditionalDirectories(request.cwd, request.additionalDirectories, request._meta);
await this.refreshSkills(request.cwd, additionalDirectories);

const response = await this.codexClient.threadFork({
config: await this.createSessionConfig(request.cwd, additionalDirectories, request.mcpServers ?? []),
cwd: request.cwd,
ephemeral: false,
modelProvider: await this.getResumeModelProvider(),
threadId: request.sessionId,
});
if (response.thread.id === request.sessionId) {
throw new Error("Codex thread/fork did not return a child session id");
}
// Codex has subscribed to the child now, so the caller must be able to clean it up if later work fails.
onSubscribed(response.thread.id);
const codexModels = await this.fetchAvailableModels();
const currentModelId = this.createModelId(codexModels, response.model, response.reasoningEffort).toString();
return {
sessionId: response.thread.id,
currentModelId: currentModelId,
models: codexModels,
collaborationMode: this.getCollaborationMode(response.thread.id),
modelProvider: response.modelProvider,
currentServiceTier: response.serviceTier as ServiceTier ?? null,
additionalDirectories,
};
}

async newSession(request: acp.NewSessionRequest): Promise<SessionMetadata> {
const additionalDirectories = readAdditionalDirectories(request.cwd, request.additionalDirectories, request._meta);
await this.refreshSkills(request.cwd, additionalDirectories);
Expand Down
132 changes: 116 additions & 16 deletions src/CodexAcpServer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -124,6 +124,17 @@ interface ActiveAuthState {
authConfigured: boolean;
}

interface InstallSessionStateOptions {
cwd: string;
sessionMetadata: SessionMetadata;
authState: ActiveAuthState;
authProvider: string | null;
requestedMcpServers: Array<acp.McpServer>;
mcpServerStartupVersion: number | null;
recoverMcpServers: boolean;
sessionTitleSource: SessionState["sessionTitleSource"];
}

interface PendingMcpStartupSession {
requestedServers: Set<string>;
afterVersion: number;
Expand Down Expand Up @@ -232,6 +243,7 @@ export class CodexAcpServer {
},
sessionCapabilities: {
resume: { },
fork: { },
list: { },
close: { },
delete: { },
Expand Down Expand Up @@ -308,8 +320,12 @@ export class CodexAcpServer {
}

async getOrCreateSession(request: acp.NewSessionRequest | acp.ResumeSessionRequest): Promise<[SessionId, LegacySessionModelState, SessionModeState]> {
return await this.withSessionOpenErrorHandling(() => this.tryCreateSession(request));
}

private async withSessionOpenErrorHandling<T>(openSession: () => Promise<T>): Promise<T> {
try {
return await this.tryCreateSession(request);
return await openSession();
} catch (e) {
const error = e instanceof Error ? e : new Error(String(e));
await this.handleError(error);
Expand Down Expand Up @@ -439,9 +455,41 @@ export class CodexAcpServer {
resumeSubscribed = false;
await this.closeStaleSessionOpen(sessionId, sessionGeneration);
}
const sessionMcpServers = this.resolveSessionMcpServers(requestedMcpServers, "sessionId" in request);
const sessionState = this.installSessionState({
cwd: request.cwd,
sessionMetadata,
authState,
authProvider,
requestedMcpServers,
mcpServerStartupVersion,
recoverMcpServers: "sessionId" in request,
sessionTitleSource: "sessionId" in request ? "unknown" : "unset",
});
resumeSubscribed = false;

this.publishAvailableCommandsAsync(sessionState);
if ("sessionId" in request) {
this.publishCurrentGoalAsync(sessionState, sessionGeneration);
}
const sessionModelState: LegacySessionModelState = this.createModelState(models, currentModelId);
const sessionModeState: SessionModeState = sessionState.agentMode.toSessionModeState();

return [sessionId, sessionModelState, sessionModeState];
}

private installSessionState(options: InstallSessionStateOptions): SessionState {
const {
cwd,
sessionMetadata,
authState,
authProvider,
requestedMcpServers,
mcpServerStartupVersion,
recoverMcpServers,
sessionTitleSource,
} = options;
const {sessionId, currentModelId, models} = sessionMetadata;
const currentModel = this.findCurrentModel(models, currentModelId);
const currentModelSupportsFast = modelSupportsFast(currentModel);
const sessionState: SessionState = {
sessionId: sessionId,
currentModelId: currentModelId,
Expand All @@ -458,18 +506,17 @@ export class CodexAcpServer {
account: authState.account,
authConfigured: authState.authConfigured,
authProvider: authProvider,
cwd: request.cwd,
cwd: cwd,
additionalDirectories: sessionMetadata.additionalDirectories,
fastModeEnabled: sessionMetadata.currentServiceTier === "fast",
currentModelSupportsFast: currentModelSupportsFast,
sessionMcpServers: sessionMcpServers,
currentModelSupportsFast: modelSupportsFast(currentModel),
sessionMcpServers: this.resolveSessionMcpServers(requestedMcpServers, recoverMcpServers),
terminalOutputMode: this.terminalOutputMode,
goalRevision: 0,
sessionTitle: null,
sessionTitleSource: "sessionId" in request ? "unknown" : "unset",
sessionTitleSource: sessionTitleSource,
};
this.sessions.set(sessionId, sessionState);
resumeSubscribed = false;

if (requestedMcpServers.length > 0 && mcpServerStartupVersion !== null) {
this.pendingMcpStartupSessions.set(sessionId, {
Expand All @@ -479,14 +526,7 @@ export class CodexAcpServer {
this.publishMcpStartupStatusAsync(sessionId);
}

this.publishAvailableCommandsAsync(sessionState);
if ("sessionId" in request) {
this.publishCurrentGoalAsync(sessionState, sessionGeneration);
}
const sessionModelState: LegacySessionModelState = this.createModelState(models, currentModelId);
const sessionModeState: SessionModeState = sessionState.agentMode.toSessionModeState();

return [sessionId, sessionModelState, sessionModeState];
return sessionState;
}

private async getAuthStateForProvider(authProvider: string | null): Promise<ActiveAuthState> {
Expand Down Expand Up @@ -560,6 +600,66 @@ export class CodexAcpServer {
};
}

async unstable_forkSession(params: acp.ForkSessionRequest): Promise<acp.ForkSessionResponse> {
logger.log("Forking session...", {sessionId: params.sessionId});
return await this.withSessionOpenErrorHandling(() => this.tryForkSession(params));
}

private async tryForkSession(params: acp.ForkSessionRequest): Promise<acp.ForkSessionResponse> {
await this.checkAuthorization();
const requestedMcpServers = params.mcpServers ?? [];
const mcpServerStartupVersion = requestedMcpServers.length > 0
? this.codexAcpClient.getMcpServerStartupVersion()
: null;
let subscribedSessionId: string | null = null;
let sessionGeneration: number | null = null;
let sessionMetadata: SessionMetadata;
let sessionState: SessionState;
try {
sessionMetadata = await this.runWithProcessCheck(() =>
this.codexAcpClient.forkSession(params, (sessionId) => {
subscribedSessionId = sessionId;
sessionGeneration = this.beginSessionOpen(sessionId);
})
);
const {sessionId} = sessionMetadata;
if (sessionGeneration === null) {
throw new Error("Codex session/fork did not report its child subscription");
}
const authProvider = sessionMetadata.modelProvider ?? this.codexAcpClient.getModelProvider();
const authState = await this.getAuthStateForProvider(authProvider);
if (!this.sessionOpenCanInstall(sessionId, sessionGeneration)) {
subscribedSessionId = null;
await this.closeStaleSessionOpen(sessionId, sessionGeneration);
}
sessionState = this.installSessionState({
cwd: params.cwd,
sessionMetadata,
authState,
authProvider,
requestedMcpServers,
mcpServerStartupVersion,
recoverMcpServers: false,
sessionTitleSource: "unknown",
});
subscribedSessionId = null;
} catch (err) {
if (subscribedSessionId !== null && sessionGeneration !== null) {
await this.cleanupStaleSessionOpen(subscribedSessionId, sessionGeneration);
}
throw err;
}

this.publishAvailableCommandsAsync(sessionState);
this.publishCurrentGoalAsync(sessionState, sessionGeneration);
logger.log("Session forked", {parentSessionId: params.sessionId, sessionId: sessionMetadata.sessionId});
return {
sessionId: sessionMetadata.sessionId,
modes: sessionState.agentMode.toSessionModeState(),
...this.createSessionConfigOptionsResponse(sessionState),
};
}

async listSessions(params: acp.ListSessionsRequest): Promise<acp.ListSessionsResponse> {
logger.log("Listing sessions...", {cwd: params.cwd, cursor: params.cursor});
await this.checkAuthorization();
Expand Down
6 changes: 6 additions & 0 deletions src/CodexAppServerClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@ import type {
ThreadGoalGetResponse,
ThreadGoalSetParams,
ThreadGoalSetResponse,
ThreadForkParams,
ThreadForkResponse,
ThreadLoadedListParams,
ThreadLoadedListResponse,
ThreadListParams,
Expand Down Expand Up @@ -528,6 +530,10 @@ export class CodexAppServerClient {
return await this.sendRequest({ method: "thread/resume", params: params });
}

async threadFork(params: ThreadForkParams): Promise<ThreadForkResponse> {
return await this.sendRequest({ method: "thread/fork", params: params });
}

getThreadSettings(threadId: string): ThreadSettings | undefined {
return this.threadSettings.get(threadId);
}
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
{
"eventType": "request",
"method": "thread/fork",
"params": {
"config": {
"projects": {
"/workspace": {
"trust_level": "trusted"
}
}
},
"cwd": "/workspace",
"ephemeral": false,
"modelProvider": "openai",
"threadId": "parent-session"
}
}
{
"eventType": "response",
"placeholder": "thread/fork"
}
{
"eventType": "request",
"method": "thread/unsubscribe",
"params": {
"threadId": "child-session"
}
}
{
"eventType": "response",
"status": "unsubscribed"
}
29 changes: 29 additions & 0 deletions src/__tests__/CodexACPAgent/data/session-fork.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
{
"eventType": "request",
"method": "thread/fork",
"params": {
"config": {
"projects": {
"/workspace": {
"trust_level": "trusted"
},
"/shared": {
"trust_level": "trusted"
}
},
"sandbox_workspace_write": {
"writable_roots": [
"/shared"
]
}
},
"cwd": "/workspace",
"ephemeral": false,
"modelProvider": "openai",
"threadId": "parent-session"
}
}
{
"eventType": "response",
"placeholder": "thread/fork"
}
1 change: 1 addition & 0 deletions src/__tests__/CodexACPAgent/initialize.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ describe('CodexACPAgent - initialize', () => {
},
sessionCapabilities: {
resume: {},
fork: {},
list: {},
close: {},
delete: {},
Expand Down
Loading