feat(vscodex): add remote Codex collaboration module

This commit is contained in:
fawney
2026-09-01 20:25:35 +08:00
parent 5a69cfe40d
commit 30a75832f8
102 changed files with 38569 additions and 11 deletions
@@ -0,0 +1,46 @@
import { CodexAgentAdapter, CodexAgentAdapterOptions } from "./codexAgentAdapter";
import { RelayClient, RelayClientOptions } from "./relayClient";
import { RelayHost, RelayHostOptions } from "./relayHost";
import { AgentAdapter, Logger, RelayTransport } from "./protocol";
export interface CodexRemoteBridgeOptions {
/** Use a supplied adapter/transport when embedding or testing. */
adapter?: AgentAdapter;
relay?: RelayTransport;
adapterOptions?: CodexAgentAdapterOptions;
relayOptions?: RelayClientOptions;
sessionId?: string;
capabilities?: Iterable<string>;
logger?: Logger;
}
export interface CodexRemoteBridge {
adapter: AgentAdapter;
relay: RelayTransport;
host: RelayHost;
start(): Promise<void>;
stop(): Promise<void>;
}
/** Construct the default outbound VS Code bridge in one call. */
export function createBridge(options: CodexRemoteBridgeOptions): CodexRemoteBridge {
const adapter = options.adapter ?? new CodexAgentAdapter(options.adapterOptions);
const relay = options.relay ?? (() => {
if (!options.relayOptions) throw new Error("relayOptions are required when no relay transport is supplied");
return new RelayClient(options.relayOptions);
})();
const hostOptions: RelayHostOptions = {
adapter,
relay,
...(options.sessionId ? { sessionId: options.sessionId } : {}),
...(options.capabilities ? { capabilities: options.capabilities } : {}),
...(options.logger ? { logger: options.logger } : {}),
};
const host = new RelayHost(hostOptions);
return {
adapter,
relay,
host,
start: () => host.start(),
stop: () => host.stop(),
};
}
@@ -0,0 +1,30 @@
import { CodexAgentAdapter } from "./codexAgentAdapter";
import { RelayHost } from "./relayHost";
import { StdioRelayTransport } from "./relayClient";
/** Standalone bridge: relay frames in stdin, relay frames out on stdout. */
async function main(): Promise<void> {
const logger = {
debug: (message: string, ...args: unknown[]) => console.error(`[debug] ${message}`, ...args),
info: (message: string, ...args: unknown[]) => console.error(`[info] ${message}`, ...args),
warn: (message: string, ...args: unknown[]) => console.error(`[warn] ${message}`, ...args),
error: (message: string, ...args: unknown[]) => console.error(`[error] ${message}`, ...args),
};
const command = process.env.CODEX_COMMAND || "codex";
const args = process.env.CODEX_APP_SERVER_ARGS ? JSON.parse(process.env.CODEX_APP_SERVER_ARGS) as string[] : ["app-server", "--stdio"];
const adapter = new CodexAgentAdapter({ command, args, defaultCwd: process.env.CODEX_WORKSPACE, logger });
const relay = new StdioRelayTransport(process.stdin, process.stdout, logger);
const host = new RelayHost({ adapter, relay, sendHandshake: true, logger });
const shutdown = async (): Promise<void> => {
await host.stop();
process.exit(0);
};
process.once("SIGINT", () => void shutdown());
process.once("SIGTERM", () => void shutdown());
await host.start();
}
void main().catch((error) => {
console.error(error instanceof Error ? error.stack ?? error.message : String(error));
process.exitCode = 1;
});
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,947 @@
/**
* Minimal client for the private Codex desktop/VS Code coordination socket.
*
* This is intentionally separate from the app-server (JSONL/stdio) adapter.
* It attaches to the already running Codex UI through the local IPC router and
* therefore does not spawn another `codex` process. The wire protocol is
* private and versioned by the official extension; keep this module isolated
* so a protocol change can fail without taking down the relay bridge.
*/
import * as crypto from "node:crypto";
import * as net from "node:net";
import * as os from "node:os";
import * as path from "node:path";
import type { JsonObject, JsonValue } from "./protocol";
export const INITIALIZING_CLIENT_ID = "initializing-client";
export const DEFAULT_IPC_REQUEST_TIMEOUT_MS = 5_000;
export const DEFAULT_MAX_IPC_FRAME_BYTES = 256 * 1024 * 1024;
/** Versions shipped by openai.chatgpt 26.820.71523. */
export const CODEX_IPC_METHOD_VERSIONS = Object.freeze({
"thread-stream-state-changed": 11,
"thread-stream-following-changed": 1,
"thread-stream-following-status-requested": 1,
"ipc-connection-reset": 1,
"thread-read-state-changed": 2,
"thread-archived": 2,
"thread-unarchived": 1,
"thread-owner-discovery": 1,
"thread-follower-start-turn": 2,
"thread-follower-load-complete-history": 1,
"thread-follower-compact-thread": 1,
"thread-follower-steer-turn": 1,
"thread-follower-interrupt-turn": 4,
"thread-follower-update-thread-settings": 1,
"thread-follower-edit-last-user-turn": 2,
"thread-follower-command-approval-decision": 1,
"thread-follower-file-approval-decision": 1,
"thread-follower-permissions-request-approval-response": 1,
"thread-follower-submit-user-input": 1,
"thread-follower-submit-mcp-server-elicitation-response": 1,
"thread-follower-set-queued-follow-ups-state": 1,
"thread-queued-followups-changed": 1,
} as const);
export type IpcMethod = keyof typeof CODEX_IPC_METHOD_VERSIONS;
export type IpcRequestId = string | number;
export type IpcPatchPathPart = string | number;
export interface IpcRequest {
type: "request";
requestId: IpcRequestId;
sourceClientId: string;
targetClientId?: string;
version: number;
method: string;
params?: JsonValue;
timeoutMs?: number;
}
export interface IpcResponse {
type: "response";
requestId: IpcRequestId;
resultType: "success" | "error";
method?: string;
handledByClientId?: string;
result?: JsonValue;
error?: string;
}
export interface IpcBroadcast {
type: "broadcast";
method: string;
sourceClientId?: string;
targetClientIds?: string[];
version: number;
params?: JsonValue;
}
export interface IpcClientDiscoveryRequest {
type: "client-discovery-request";
requestId: IpcRequestId;
request: IpcRequest;
}
export interface IpcClientDiscoveryResponse {
type: "client-discovery-response";
requestId: IpcRequestId;
response: { canHandle: boolean };
}
export type IpcMessage =
| IpcRequest
| IpcResponse
| IpcBroadcast
| IpcClientDiscoveryRequest
| IpcClientDiscoveryResponse;
export interface IpcJsonPatch {
op: "add" | "remove" | "replace";
path: IpcPatchPathPart[];
value?: JsonValue;
}
export interface ThreadStreamSnapshot {
type: "snapshot";
revision: number;
conversationState: JsonObject;
}
export interface ThreadStreamPatches {
type: "patches";
baseRevision: number;
revision: number;
patches: IpcJsonPatch[];
}
export type ThreadStreamChange = ThreadStreamSnapshot | ThreadStreamPatches;
export interface ConversationStreamState {
conversationId: string;
hostId: string;
ownerClientId: string;
revision: number;
conversationState: JsonObject;
}
export type ConversationStreamEvent =
| (ConversationStreamState & { kind: "snapshot"; raw: IpcBroadcast })
| (ConversationStreamState & { kind: "patches"; patches: IpcJsonPatch[]; baseRevision: number; raw: IpcBroadcast })
| {
kind: "desync";
conversationId: string;
hostId: string;
ownerClientId: string;
expectedRevision: number;
receivedBaseRevision: number;
receivedRevision: number;
raw: IpcBroadcast;
};
export interface CodexIpcClientOptions {
/** Explicit socket path; otherwise `$CODEX_HOME/ipc/ipc.sock` or `~/.codex`. */
socketPath?: string;
codexHome?: string;
homeDir?: string;
env?: NodeJS.ProcessEnv;
platform?: NodeJS.Platform;
clientType?: string;
requestTimeoutMs?: number;
maxFrameBytes?: number;
strictVersions?: boolean;
/** Reconnect after a socket close and re-send all active following subscriptions. */
autoReconnect?: boolean;
reconnectDelayMs?: number;
/** Optional handler for discovery requests. Default is fail-closed (`false`). */
canHandleRequest?: (request: IpcRequest) => boolean | Promise<boolean>;
}
export interface FollowerTurnStartOptions {
request?: JsonObject;
context?: JsonObject;
clientUserMessageId?: string;
ownerClientId?: string;
timeoutMs?: number;
}
export interface FollowerSteerOptions {
clientUserMessageId?: string;
serviceTier?: string | null;
attachments?: JsonValue[];
additionalContext?: JsonObject | null;
restoreMessage?: JsonValue | null;
ownerClientId?: string;
timeoutMs?: number;
}
export interface FollowerInterruptOptions {
mode?: "user-stop" | "system" | "descendant-cleanup" | string;
expectedTurnId?: string | null;
ownerClientId?: string;
timeoutMs?: number;
}
export interface FollowOptions {
hostId?: string;
targetClientIds?: string[];
}
export interface RequestOptions {
targetClientId?: string;
timeoutMs?: number;
version?: number;
requestId?: IpcRequestId;
}
export interface IpcErrorOptions {
code: string;
response?: IpcResponse;
}
export class CodexIpcError extends Error {
readonly code: string;
readonly response?: IpcResponse;
constructor(message: string, options: IpcErrorOptions) {
super(message);
this.name = "CodexIpcError";
this.code = options.code;
this.response = options.response;
}
}
export function resolveCodexIpcSocketPath(options: {
socketPath?: string;
codexHome?: string;
homeDir?: string;
env?: NodeJS.ProcessEnv;
platform?: NodeJS.Platform;
} = {}): string {
if (options.socketPath?.trim()) return options.socketPath.trim();
const platform = options.platform ?? process.platform;
if (platform === "win32") return "\\\\.\\pipe\\codex-ipc";
const env = options.env ?? process.env;
const homeDir = options.homeDir ?? os.homedir();
const configuredHome = options.codexHome?.trim() || env.CODEX_HOME?.trim() || path.join(homeDir, ".codex");
const codexHome = configuredHome === "~"
? homeDir
: configuredHome.startsWith("~/")
? path.join(homeDir, configuredHome.slice(2))
: configuredHome;
return path.join(codexHome, "ipc", "ipc.sock");
}
/** Encode one private IPC frame: uint32 little-endian byte length + UTF-8 JSON. */
export function encodeIpcFrame(message: IpcMessage, maxFrameBytes = DEFAULT_MAX_IPC_FRAME_BYTES): Buffer {
const json = JSON.stringify(message);
const payload = Buffer.from(json, "utf8");
if (payload.length === 0 || payload.length > maxFrameBytes) {
throw new RangeError(`IPC frame exceeds ${maxFrameBytes} bytes`);
}
const frame = Buffer.allocUnsafe(4 + payload.length);
frame.writeUInt32LE(payload.length, 0);
payload.copy(frame, 4);
return frame;
}
/** Incremental decoder that accepts arbitrary TCP/Unix-socket chunk boundaries. */
export class IpcFrameDecoder {
private buffer = Buffer.alloc(0);
constructor(private readonly maxFrameBytes = DEFAULT_MAX_IPC_FRAME_BYTES) {}
push(chunk: Uint8Array): IpcMessage[] {
if (chunk.length === 0) return [];
this.buffer = this.buffer.length === 0 ? Buffer.from(chunk) : Buffer.concat([this.buffer, chunk]);
const messages: IpcMessage[] = [];
while (this.buffer.length >= 4) {
const payloadLength = this.buffer.readUInt32LE(0);
if (payloadLength === 0 || payloadLength > this.maxFrameBytes) {
throw new CodexIpcError(`Invalid IPC frame length (${payloadLength} bytes)`, { code: "invalid-frame-length" });
}
if (this.buffer.length < payloadLength + 4) break;
const payload = this.buffer.subarray(4, payloadLength + 4).toString("utf8");
this.buffer = this.buffer.subarray(payloadLength + 4);
let decoded: unknown;
try {
decoded = JSON.parse(payload);
} catch (error) {
throw new CodexIpcError(`Invalid IPC JSON: ${error instanceof Error ? error.message : String(error)}`, {
code: "invalid-json",
});
}
if (!isRecord(decoded) || typeof decoded.type !== "string") {
throw new CodexIpcError("IPC frame must be an object with a type", { code: "invalid-message" });
}
messages.push(decoded as unknown as IpcMessage);
}
return messages;
}
reset(): void {
this.buffer = Buffer.alloc(0);
}
}
type Listener<T> = (value: T) => void;
export interface IpcSubscription { dispose(): void; }
function subscribe<T>(set: Set<Listener<T>>, listener: Listener<T>): IpcSubscription {
set.add(listener);
return { dispose: () => set.delete(listener) };
}
function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === "object" && value !== null && !Array.isArray(value);
}
function isJsonObject(value: unknown): value is JsonObject {
return isRecord(value);
}
function requestIdKey(id: IpcRequestId): string {
return `${typeof id}:${String(id)}`;
}
function cloneJson<T extends JsonValue>(value: T): T {
return JSON.parse(JSON.stringify(value)) as T;
}
function versionFor(method: string, params?: JsonValue): number {
// The official client accepts interrupt v3 when expectedTurnId is absent;
// v4 is used when the active-turn precondition is present.
if (method === "thread-follower-interrupt-turn"
&& (!isRecord(params) || params.expectedTurnId === undefined || params.expectedTurnId === null)) return 3;
return CODEX_IPC_METHOD_VERSIONS[method as IpcMethod] ?? 0;
}
function textInput(text: string): JsonObject {
return { type: "text", text, text_elements: [] };
}
function normalizeInput(input: string | JsonValue[]): JsonValue[] {
return typeof input === "string"
? [textInput(input)]
: input.map((entry) => typeof entry === "string" ? textInput(entry) : entry);
}
function hasTarget(frame: IpcBroadcast, clientId: string): boolean {
return frame.targetClientIds == null || frame.targetClientIds.includes(clientId);
}
/** Apply the JSON patch arrays generated by Immer in the official webview. */
export function applyIpcPatches(root: JsonValue, patches: IpcJsonPatch[]): JsonValue {
let result = cloneJson(root);
for (const patch of patches) {
if (!Array.isArray(patch.path)) throw new CodexIpcError("IPC patch path must be an array", { code: "invalid-patch" });
if (patch.path.length === 0) {
if (patch.op === "remove") throw new CodexIpcError("Removing the conversation root is unsupported", { code: "invalid-patch" });
if (patch.value === undefined) throw new CodexIpcError("Patch value is missing", { code: "invalid-patch" });
result = cloneJson(patch.value);
continue;
}
const parentPath = patch.path.slice(0, -1);
const key = patch.path[patch.path.length - 1];
assertSafePatchPart(key);
const parent = getAtPath(result, parentPath);
if (Array.isArray(parent)) {
const index = key === "-" ? parent.length : toArrayIndex(key);
if (patch.op === "add") {
if (patch.value === undefined) throw new CodexIpcError("Patch value is missing", { code: "invalid-patch" });
parent.splice(index, 0, cloneJson(patch.value));
} else if (patch.op === "replace") {
if (patch.value === undefined || index < 0 || index >= parent.length) throw new CodexIpcError("Invalid array replace patch", { code: "invalid-patch" });
parent[index] = cloneJson(patch.value);
} else {
if (index < 0 || index >= parent.length) throw new CodexIpcError("Invalid array remove patch", { code: "invalid-patch" });
parent.splice(index, 1);
}
continue;
}
if (!isRecord(parent) || typeof key !== "string") {
throw new CodexIpcError("IPC patch parent is not an object or array", { code: "invalid-patch" });
}
if (patch.op === "remove") {
delete parent[key];
} else {
if (patch.value === undefined) throw new CodexIpcError("Patch value is missing", { code: "invalid-patch" });
parent[key] = cloneJson(patch.value);
}
}
return result;
}
function getAtPath(root: JsonValue, pathParts: IpcPatchPathPart[]): JsonValue {
let current: JsonValue = root;
for (const part of pathParts) {
if (Array.isArray(current)) {
const index = toArrayIndex(part);
if (index < 0 || index >= current.length) throw new CodexIpcError("IPC patch path is out of bounds", { code: "invalid-patch" });
current = current[index];
} else if (isRecord(current) && typeof part === "string" && Object.prototype.hasOwnProperty.call(current, part)) {
assertSafePatchPart(part);
current = current[part];
} else {
throw new CodexIpcError("IPC patch path does not exist", { code: "invalid-patch" });
}
}
return current;
}
function toArrayIndex(value: IpcPatchPathPart): number {
if (typeof value === "number" && Number.isInteger(value)) return value;
if (typeof value === "string" && /^\d+$/.test(value)) return Number(value);
throw new CodexIpcError(`Invalid array patch index: ${String(value)}`, { code: "invalid-patch" });
}
function assertSafePatchPart(value: IpcPatchPathPart): void {
if (value === "__proto__" || value === "prototype" || value === "constructor") {
throw new CodexIpcError("Unsafe IPC patch path", { code: "invalid-patch" });
}
}
export class CodexIpcClient {
readonly socketPath: string;
private readonly options: Required<Pick<CodexIpcClientOptions, "clientType" | "requestTimeoutMs" | "maxFrameBytes" | "strictVersions" | "autoReconnect" | "reconnectDelayMs">> & CodexIpcClientOptions;
private socket: net.Socket | undefined;
private decoder: IpcFrameDecoder;
private connectPromise: Promise<string> | undefined;
private reconnectTimer: NodeJS.Timeout | undefined;
private disposed = false;
private clientId = INITIALIZING_CLIENT_ID;
private readonly pending = new Map<string, { method: string; resolve: (response: IpcResponse) => void; reject: (error: Error) => void; timer: NodeJS.Timeout }>();
private readonly followed = new Map<string, string>();
private readonly streams = new Map<string, ConversationStreamState>();
private readonly messageListeners = new Set<Listener<IpcMessage>>();
private readonly broadcastListeners = new Set<Listener<IpcBroadcast>>();
private readonly streamListeners = new Set<Listener<ConversationStreamEvent>>();
private readonly errorListeners = new Set<Listener<Error>>();
private readonly closeListeners = new Set<Listener<Error | undefined>>();
private readonly discoveryHandler?: (request: IpcRequest) => boolean | Promise<boolean>;
constructor(options: CodexIpcClientOptions = {}) {
this.options = {
...options,
clientType: options.clientType ?? "codex-remote-collab",
requestTimeoutMs: options.requestTimeoutMs ?? DEFAULT_IPC_REQUEST_TIMEOUT_MS,
maxFrameBytes: options.maxFrameBytes ?? DEFAULT_MAX_IPC_FRAME_BYTES,
strictVersions: options.strictVersions ?? true,
autoReconnect: options.autoReconnect ?? false,
reconnectDelayMs: options.reconnectDelayMs ?? 1_000,
};
this.socketPath = resolveCodexIpcSocketPath(options);
this.decoder = new IpcFrameDecoder(this.options.maxFrameBytes);
this.discoveryHandler = options.canHandleRequest;
}
getClientId(): string { return this.clientId; }
getConversationState(conversationId: string): ConversationStreamState | undefined {
const state = this.streams.get(conversationId);
return state == null ? undefined : { ...state, conversationState: cloneJson(state.conversationState) };
}
getFollowedConversations(): ReadonlyMap<string, string> { return this.followed; }
onMessage(listener: Listener<IpcMessage>): IpcSubscription { return subscribe(this.messageListeners, listener); }
onBroadcast(listener: Listener<IpcBroadcast>): IpcSubscription { return subscribe(this.broadcastListeners, listener); }
onStreamEvent(listener: Listener<ConversationStreamEvent>): IpcSubscription { return subscribe(this.streamListeners, listener); }
onError(listener: Listener<Error>): IpcSubscription { return subscribe(this.errorListeners, listener); }
onClose(listener: Listener<Error | undefined>): IpcSubscription { return subscribe(this.closeListeners, listener); }
async connect(): Promise<string> {
if (this.disposed) throw new CodexIpcError("IPC client is disposed", { code: "disposed" });
if (this.reconnectTimer) {
clearTimeout(this.reconnectTimer);
this.reconnectTimer = undefined;
}
if (this.socket?.writable && this.clientId !== INITIALIZING_CLIENT_ID) return this.clientId;
if (this.connectPromise) return this.connectPromise;
this.connectPromise = new Promise<string>((resolve, reject) => {
const socket = net.createConnection(this.socketPath);
this.socket = socket;
this.decoder.reset();
let settled = false;
const finishError = (error: Error): void => {
if (!settled) {
settled = true;
reject(error);
}
this.emitError(error);
};
socket.setNoDelay?.(true);
socket.on("connect", () => {
const requestId = crypto.randomUUID();
const timer = setTimeout(() => {
this.pending.delete(requestIdKey(requestId));
finishError(new CodexIpcError("IPC initialize timed out", { code: "timeout" }));
socket.destroy();
}, this.options.requestTimeoutMs);
this.pending.set(requestIdKey(requestId), {
method: "initialize",
resolve: (response) => {
clearTimeout(timer);
if (response.resultType !== "success" || !isRecord(response.result) || typeof response.result.clientId !== "string") {
finishError(new CodexIpcError("IPC initialize returned an invalid response", { code: "initialize-failed", response }));
socket.destroy();
return;
}
this.clientId = response.result.clientId;
settled = true;
resolve(this.clientId);
this.resubscribeAfterConnect().catch((error) => this.emitError(asError(error)));
},
reject: (error) => {
clearTimeout(timer);
finishError(error);
socket.destroy();
},
timer,
});
this.write({
type: "request",
requestId,
sourceClientId: INITIALIZING_CLIENT_ID,
version: 0,
method: "initialize",
params: { clientType: this.options.clientType },
});
});
socket.on("data", (chunk) => {
try {
for (const message of this.decoder.push(chunk)) this.handleMessage(message);
} catch (error) {
const normalized = asError(error);
finishError(normalized);
socket.destroy(normalized);
}
});
socket.on("error", (error) => {
if (!settled) finishError(error);
else this.emitError(error);
});
socket.on("close", () => {
this.handleClose();
});
}).finally(() => {
this.connectPromise = undefined;
});
return this.connectPromise;
}
async followConversation(conversationId: string, following = true, options: FollowOptions = {}): Promise<void> {
const hostId = options.hostId ?? "local";
await this.connect();
if (following) this.followed.set(conversationId, hostId);
else {
this.followed.delete(conversationId);
this.streams.delete(conversationId);
}
const params: JsonObject = { conversationId, hostId, following };
const frame: IpcBroadcast = {
type: "broadcast",
method: "thread-stream-following-changed",
sourceClientId: this.clientId,
version: CODEX_IPC_METHOD_VERSIONS["thread-stream-following-changed"],
params,
};
if (options.targetClientIds) frame.targetClientIds = options.targetClientIds;
this.write(frame);
}
async findThreadOwner(conversationId: string, hostId = "local", timeoutMs = this.options.requestTimeoutMs): Promise<string | null> {
try {
const response = await this.request("thread-owner-discovery", { conversationId, hostId }, { timeoutMs });
return response.handledByClientId ?? null;
} catch (error) {
if (error instanceof CodexIpcError
&& (error.code === "no-client-found" || error.code.startsWith("no-client-found:"))) return null;
throw error;
}
}
async request(method: string, params?: JsonValue, options: RequestOptions = {}): Promise<IpcResponse> {
await this.connect();
const requestId = options.requestId ?? crypto.randomUUID();
const timeoutMs = options.timeoutMs ?? this.options.requestTimeoutMs;
const frame: IpcRequest = {
type: "request",
requestId,
sourceClientId: this.clientId,
version: options.version ?? versionFor(method, params),
method,
params,
};
if (options.targetClientId) frame.targetClientId = options.targetClientId;
if (timeoutMs > 0) frame.timeoutMs = timeoutMs;
return new Promise<IpcResponse>((resolve, reject) => {
const key = requestIdKey(requestId);
const timer = setTimeout(() => {
this.pending.delete(key);
reject(new CodexIpcError(`${method} timed out`, { code: "timeout" }));
}, timeoutMs > 0 ? timeoutMs : 2 ** 31 - 1);
this.pending.set(key, { method, resolve, reject, timer });
try {
this.write(frame);
} catch (error) {
clearTimeout(timer);
this.pending.delete(key);
reject(asError(error));
}
}).then((response) => {
if (response.resultType === "error") {
throw new CodexIpcError(response.error ?? `${method} failed`, { code: response.error ?? "ipc-error", response });
}
if (response.method != null && response.method !== method) {
throw new CodexIpcError(`IPC response method mismatch: expected ${method}, got ${response.method}`, {
code: "response-method-mismatch",
response,
});
}
return response;
});
}
async requestFollower(method: string, conversationId: string, params: JsonObject = {}, options: RequestOptions & { ownerClientId?: string } = {}): Promise<IpcResponse> {
const ownerClientId = options.ownerClientId ?? this.streams.get(conversationId)?.ownerClientId;
if (!ownerClientId) throw new CodexIpcError(`No owner is known for conversation ${conversationId}`, { code: "owner-unknown" });
// Do not allow a caller-provided params object to accidentally retarget a
// request after the owner has been selected from the stream snapshot.
const body: JsonObject = { ...params, conversationId };
const { ownerClientId: _owner, ...requestOptions } = options;
return this.request(method, body, { ...requestOptions, targetClientId: ownerClientId });
}
/** Send the exact private `turnStart` envelope expected by the owner. */
async startTurn(conversationId: string, input: string | JsonValue[], options: FollowerTurnStartOptions = {}): Promise<JsonValue | undefined> {
const request: JsonObject = {
...(options.request ?? {}),
threadId: conversationId,
input: options.request?.input ?? normalizeInput(input),
};
const context: JsonObject = { inheritThreadSettings: true, ...(options.context ?? {}) };
if (options.clientUserMessageId) request.clientUserMessageId = options.clientUserMessageId;
const response = await this.requestFollower("thread-follower-start-turn", conversationId, {
turnStart: { request, context },
}, {
ownerClientId: options.ownerClientId,
timeoutMs: options.timeoutMs,
});
return response.result;
}
async steerTurn(conversationId: string, input: string | JsonValue[], options: FollowerSteerOptions = {}): Promise<JsonValue | undefined> {
const params: JsonObject = {
clientUserMessageId: options.clientUserMessageId ?? crypto.randomUUID(),
input: normalizeInput(input),
attachments: options.attachments ?? [],
};
if (options.serviceTier !== undefined) params.serviceTier = options.serviceTier;
if (options.additionalContext !== undefined) params.additionalContext = options.additionalContext;
if (options.restoreMessage !== undefined) params.restoreMessage = options.restoreMessage;
const response = await this.requestFollower("thread-follower-steer-turn", conversationId, params, {
ownerClientId: options.ownerClientId,
timeoutMs: options.timeoutMs,
});
return response.result;
}
/**
* Persist settings for the next turn through the official conversation
* owner. The owner-side follower handler expects the settings nested under
* `threadSettings`; `requestFollower` adds the conversation id to the
* outer envelope, yielding:
* `{ conversationId, threadSettings }`.
*/
async updateThreadSettings(
conversationId: string,
threadSettings: JsonObject,
options: RequestOptions & { ownerClientId?: string } = {},
): Promise<JsonValue | undefined> {
if (!isJsonObject(threadSettings)) {
throw new CodexIpcError("thread settings must be a JSON object", { code: "invalid-thread-settings" });
}
const response = await this.requestFollower(
"thread-follower-update-thread-settings",
conversationId,
{ threadSettings: cloneJson(threadSettings) },
options,
);
return response.result;
}
/** Alias matching the official app-server manager method name. */
async updateThreadSettingsForNextTurn(
conversationId: string,
threadSettings: JsonObject,
options: RequestOptions & { ownerClientId?: string } = {},
): Promise<JsonValue | undefined> {
return this.updateThreadSettings(conversationId, threadSettings, options);
}
async interruptTurn(conversationId: string, options: FollowerInterruptOptions = {}): Promise<JsonValue | undefined> {
const params: JsonObject = { mode: options.mode ?? "user-stop" };
if (options.expectedTurnId !== undefined && options.expectedTurnId !== null) params.expectedTurnId = options.expectedTurnId;
const response = await this.requestFollower("thread-follower-interrupt-turn", conversationId, params, {
ownerClientId: options.ownerClientId,
timeoutMs: options.timeoutMs,
});
return response.result;
}
async loadCompleteHistory(conversationId: string, options: RequestOptions & { ownerClientId?: string } = {}): Promise<JsonValue | undefined> {
const response = await this.requestFollower("thread-follower-load-complete-history", conversationId, {}, options);
return response.result;
}
async respondCommandApproval(conversationId: string, requestId: IpcRequestId, decision: JsonValue, options: RequestOptions & { ownerClientId?: string } = {}): Promise<JsonValue | undefined> {
return this.respondFollower("thread-follower-command-approval-decision", conversationId, { requestId, decision }, options);
}
async respondFileApproval(conversationId: string, requestId: IpcRequestId, decision: JsonValue, options: RequestOptions & { ownerClientId?: string } = {}): Promise<JsonValue | undefined> {
return this.respondFollower("thread-follower-file-approval-decision", conversationId, { requestId, decision }, options);
}
async respondPermissionsApproval(conversationId: string, requestId: IpcRequestId, response: JsonValue, options: RequestOptions & { ownerClientId?: string } = {}): Promise<JsonValue | undefined> {
return this.respondFollower("thread-follower-permissions-request-approval-response", conversationId, { requestId, response }, options);
}
async respondUserInput(conversationId: string, requestId: IpcRequestId, response: JsonValue, options: RequestOptions & { ownerClientId?: string } = {}): Promise<JsonValue | undefined> {
return this.respondFollower("thread-follower-submit-user-input", conversationId, { requestId, response }, options);
}
async respondMcpElicitation(conversationId: string, requestId: IpcRequestId, response: JsonValue, options: RequestOptions & { ownerClientId?: string } = {}): Promise<JsonValue | undefined> {
return this.respondFollower("thread-follower-submit-mcp-server-elicitation-response", conversationId, { requestId, response }, options);
}
private async respondFollower(method: string, conversationId: string, params: JsonObject, options: RequestOptions & { ownerClientId?: string }): Promise<JsonValue | undefined> {
const response = await this.requestFollower(method, conversationId, params, options);
return response.result;
}
async dispose(): Promise<void> {
this.disposed = true;
if (this.reconnectTimer) clearTimeout(this.reconnectTimer);
this.reconnectTimer = undefined;
for (const pending of this.pending.values()) {
clearTimeout(pending.timer);
pending.reject(new CodexIpcError("IPC client disposed", { code: "disposed" }));
}
this.pending.clear();
this.socket?.destroy();
this.socket = undefined;
this.clientId = INITIALIZING_CLIENT_ID;
}
private write(message: IpcMessage): void {
if (!this.socket?.writable) throw new CodexIpcError("IPC socket is not connected", { code: "not-connected" });
this.socket.write(encodeIpcFrame(message, this.options.maxFrameBytes));
}
private handleMessage(message: IpcMessage): void {
for (const listener of this.messageListeners) safeCall(listener, message, (error) => this.emitError(error));
switch (message.type) {
case "response":
this.handleResponse(message);
return;
case "broadcast":
this.handleBroadcast(message);
return;
case "client-discovery-request":
this.handleDiscoveryRequest(message).catch((error) => this.emitError(asError(error)));
return;
case "request":
this.handleUnexpectedRequest(message);
return;
case "client-discovery-response":
// Discovery responses are consumed by the router, not by clients.
return;
}
}
private handleResponse(response: IpcResponse): void {
const key = requestIdKey(response.requestId);
const pending = this.pending.get(key);
if (!pending) return;
this.pending.delete(key);
clearTimeout(pending.timer);
pending.resolve(response);
}
private handleBroadcast(frame: IpcBroadcast): void {
if (!hasTarget(frame, this.clientId)) return;
for (const listener of this.broadcastListeners) safeCall(listener, frame, (error) => this.emitError(error));
if (frame.method === "thread-stream-state-changed") {
this.handleStreamStateBroadcast(frame);
} else if (frame.method === "thread-stream-following-status-requested") {
this.handleFollowingStatusRequested(frame);
}
}
/** Re-announce active subscriptions when an owner reconnects or hands off. */
private handleFollowingStatusRequested(frame: IpcBroadcast): void {
if (!isRecord(frame.params)) return;
if (this.options.strictVersions
&& frame.version !== CODEX_IPC_METHOD_VERSIONS["thread-stream-following-status-requested"]) {
this.emitError(new CodexIpcError(`Unsupported thread following status version ${frame.version}`, { code: "version-mismatch" }));
return;
}
const conversationId = typeof frame.params.conversationId === "string"
? frame.params.conversationId
: undefined;
const hostId = typeof frame.params.hostId === "string" ? frame.params.hostId : "local";
const requester = frame.sourceClientId;
if (!conversationId || !requester || requester === this.clientId) return;
if (this.followed.get(conversationId) !== hostId) return;
void this.followConversation(conversationId, true, {
hostId,
targetClientIds: [requester],
}).catch((error) => this.emitError(asError(error)));
}
private handleStreamStateBroadcast(frame: IpcBroadcast): void {
if (!isRecord(frame.params)) return;
const conversationId = typeof frame.params.conversationId === "string" ? frame.params.conversationId : undefined;
const hostId = typeof frame.params.hostId === "string" ? frame.params.hostId : "local";
const change = frame.params.change;
if (!conversationId || !isRecord(change) || typeof change.type !== "string") return;
if (this.options.strictVersions && frame.version !== CODEX_IPC_METHOD_VERSIONS["thread-stream-state-changed"]) {
this.emitError(new CodexIpcError(`Unsupported thread stream version ${frame.version}`, { code: "version-mismatch" }));
return;
}
const ownerClientId = frame.sourceClientId ?? "";
if (change.type === "snapshot") {
if (typeof change.revision !== "number" || !isJsonObject(change.conversationState)) return;
const state: ConversationStreamState = {
conversationId,
hostId,
ownerClientId,
revision: change.revision,
conversationState: cloneJson(change.conversationState),
};
this.streams.set(conversationId, state);
this.emitStream({ kind: "snapshot", ...state, raw: frame });
return;
}
if (change.type !== "patches" || typeof change.baseRevision !== "number" || typeof change.revision !== "number" || !Array.isArray(change.patches)) return;
const current = this.streams.get(conversationId);
if (!current || current.ownerClientId !== ownerClientId || current.revision !== change.baseRevision) {
const expectedRevision = current?.revision ?? 0;
this.emitStream({
kind: "desync",
conversationId,
hostId,
ownerClientId,
expectedRevision,
receivedBaseRevision: change.baseRevision,
receivedRevision: change.revision,
raw: frame,
});
// Re-sending `following:true` is how the official follower asks the
// owner for a fresh snapshot when a patch base revision is missed.
if (this.followed.has(conversationId)) {
this.followConversation(conversationId, true, { hostId }).catch((error) => this.emitError(asError(error)));
}
return;
}
try {
const patches = change.patches as unknown as IpcJsonPatch[];
const nextConversationState = applyIpcPatches(current.conversationState, patches);
if (!isJsonObject(nextConversationState)) throw new CodexIpcError("Patched conversation state is not an object", { code: "invalid-patch" });
const next: ConversationStreamState = {
...current,
revision: change.revision,
conversationState: nextConversationState,
};
this.streams.set(conversationId, next);
this.emitStream({ kind: "patches", ...next, patches, baseRevision: change.baseRevision, raw: frame });
} catch (error) {
this.emitError(asError(error));
}
}
private async handleDiscoveryRequest(message: IpcClientDiscoveryRequest): Promise<void> {
const request = message.request;
let canHandle = false;
try {
canHandle = this.discoveryHandler ? await this.discoveryHandler(request) : false;
} catch {
canHandle = false;
}
this.write({
type: "client-discovery-response",
requestId: message.requestId,
response: { canHandle },
});
}
private handleUnexpectedRequest(request: IpcRequest): void {
try {
this.write({
type: "response",
requestId: request.requestId,
resultType: "error",
error: "no-handler-for-request",
});
} catch (error) {
this.emitError(asError(error));
}
}
private async resubscribeAfterConnect(): Promise<void> {
const subscriptions = [...this.followed.entries()];
for (const [conversationId, hostId] of subscriptions) {
this.write({
type: "broadcast",
method: "thread-stream-following-changed",
sourceClientId: this.clientId,
version: CODEX_IPC_METHOD_VERSIONS["thread-stream-following-changed"],
params: { conversationId, hostId, following: true },
});
}
}
private handleClose(): void {
const socket = this.socket;
this.socket = undefined;
this.decoder.reset();
const closeError = new CodexIpcError("IPC socket closed", { code: "connection-closed" });
for (const pending of this.pending.values()) {
clearTimeout(pending.timer);
pending.reject(closeError);
}
this.pending.clear();
this.clientId = INITIALIZING_CLIENT_ID;
for (const listener of this.closeListeners) safeCall(listener, closeError, (error) => this.emitError(error));
if (!this.disposed && this.options.autoReconnect && socket) {
this.reconnectTimer = setTimeout(() => {
this.reconnectTimer = undefined;
this.connect().catch((error) => this.emitError(asError(error)));
}, this.options.reconnectDelayMs);
}
}
private emitStream(event: ConversationStreamEvent): void {
for (const listener of this.streamListeners) safeCall(listener, event, (error) => this.emitError(error));
}
private emitError(error: Error): void {
for (const listener of this.errorListeners) safeCall(listener, error, () => undefined);
}
}
function safeCall<T>(listener: Listener<T>, value: T, onError: (error: Error) => void): void {
try {
listener(value);
} catch (error) {
onError(asError(error));
}
}
function asError(error: unknown): Error {
return error instanceof Error ? error : new Error(String(error));
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,152 @@
import { accessSync, constants, Dirent, readdirSync, statSync } from "node:fs";
import { homedir } from "node:os";
import { delimiter, isAbsolute, join, sep } from "node:path";
export interface CodexPathOptions {
/** Environment used for PATH lookup. Defaults to the extension host environment. */
env?: NodeJS.ProcessEnv;
/** Home directory used when looking for bundled installations. */
homeDir?: string;
/** Platform override for deterministic tests. */
platform?: NodeJS.Platform;
}
/**
* Resolve the executable used by the VS Code bridge.
*
* VS Code launched from Finder/Dock often receives a smaller PATH than a shell.
* The default `codex` command therefore gets a few explicit installation
* fallbacks, while a user-supplied command remains authoritative.
*/
export function resolveCodexCommand(configuredCommand = "codex", options: CodexPathOptions = {}): string {
const command = configuredCommand.trim() || "codex";
const env = options.env ?? process.env;
const platform = options.platform ?? process.platform;
const home = options.homeDir ?? homedir();
if (hasPathComponent(command, platform)) {
const resolved = executablePath(command, platform);
if (resolved) return resolved;
throw missingCodexError(command, platform);
}
const fromPath = findOnPath(command, env.PATH, platform, env.PATHEXT);
if (fromPath) return fromPath;
// Only the default command gets installation-specific fallbacks. A custom
// bare command should fail loudly instead of silently running another binary.
if (!isDefaultCommand(command, platform)) throw missingCodexError(command, platform);
for (const candidate of bundledCandidates(home, platform)) {
const resolved = executablePath(candidate, platform);
if (resolved) return resolved;
}
throw missingCodexError(command, platform);
}
export function missingCodexError(command: string, platform: NodeJS.Platform = process.platform): Error {
const examples = platform === "darwin"
? ' Set "codexRemoteCollab.codexCommand" to the full path, for example "/Applications/ChatGPT.app/Contents/Resources/codex".'
: ' Set "codexRemoteCollab.codexCommand" to the full path of the Codex executable.';
return new Error(`Codex executable "${command}" was not found.${examples}`);
}
function isDefaultCommand(command: string, platform: NodeJS.Platform): boolean {
return platform === "win32" ? command.toLowerCase() === "codex" || command.toLowerCase() === "codex.exe" : command === "codex";
}
function hasPathComponent(command: string, platform: NodeJS.Platform): boolean {
return isAbsolute(command) || command.includes(sep) || (platform === "win32" && command.includes("\\"));
}
function executablePath(candidate: string, platform: NodeJS.Platform): string | undefined {
try {
const info = statSync(candidate);
if (!info.isFile()) return undefined;
// X_OK is meaningful on POSIX; Windows still benefits from the file check.
if (platform !== "win32") accessSync(candidate, constants.X_OK);
return candidate;
} catch {
return undefined;
}
}
function findOnPath(command: string, pathValue: string | undefined, platform: NodeJS.Platform, pathextValue?: string): string | undefined {
if (!pathValue) return undefined;
const extensions = platform === "win32" ? windowsExtensions(command, pathextValue) : [""];
for (const directory of pathValue.split(delimiter)) {
if (!directory) continue;
for (const extension of extensions) {
const candidate = join(directory, `${command}${extension}`);
const resolved = executablePath(candidate, platform);
if (resolved) return resolved;
}
}
return undefined;
}
function windowsExtensions(command: string, pathextValue: string | undefined): string[] {
if (/[.][^./\\]+$/.test(command)) return [""];
const extensions = (pathextValue ?? ".COM;.EXE;.BAT;.CMD")
.split(";")
.map((value) => value.trim())
.filter(Boolean);
return ["", ...extensions];
}
function bundledCandidates(home: string, platform: NodeJS.Platform): string[] {
if (platform !== "darwin") return [];
const candidates = [
join(home, "Applications", "ChatGPT.app", "Contents", "Resources", "codex"),
"/Applications/ChatGPT.app/Contents/Resources/codex",
join(home, ".local", "bin", "codex"),
join(home, ".npm-global", "bin", "codex"),
];
for (const extensionsRoot of [
join(home, ".vscode", "extensions"),
join(home, ".vscode-insiders", "extensions"),
]) {
candidates.push(...officialExtensionCandidates(extensionsRoot));
}
return candidates;
}
function officialExtensionCandidates(extensionsRoot: string): string[] {
let entries: Dirent<string>[];
try {
entries = readdirSync(extensionsRoot, { withFileTypes: true, encoding: "utf8" });
} catch {
return [];
}
const matches = entries
.filter((entry) => entry.isDirectory() && entry.name.startsWith("openai.chatgpt-"))
.map((entry) => {
const directory = join(extensionsRoot, entry.name);
let modified = 0;
try {
modified = statSync(directory).mtimeMs;
} catch {
// Keep an unreadable entry at the end of the deterministic sort.
}
return { directory, modified };
})
.sort((left, right) => right.modified - left.modified || right.directory.localeCompare(left.directory));
const candidates: string[] = [];
for (const match of matches) {
let architectures: Dirent<string>[];
try {
architectures = readdirSync(join(match.directory, "bin"), { withFileTypes: true, encoding: "utf8" });
} catch {
continue;
}
for (const architecture of architectures) {
if (architecture.isDirectory()) candidates.push(join(match.directory, "bin", architecture.name, "codex"));
}
}
return candidates;
}
@@ -0,0 +1,131 @@
import { Disposable, RelayFrame, RelayTransport } from "./protocol";
export interface NamedRelayTransport {
id: string;
transport: RelayTransport;
required?: boolean;
}
/**
* Fans host events out to local and cloud relays while presenting one
* transport lifecycle to RelayHost. A temporary cloud outage must not stop
* the local bridge (and vice versa).
*/
export class CompositeRelayTransport implements RelayTransport {
readonly handlesHandshake = true;
private readonly entries: NamedRelayTransport[];
private readonly subscriptions: Disposable[] = [];
private readonly openEntries = new Set<string>();
private readonly messageListeners = new Set<(frame: RelayFrame) => void>();
private readonly openListeners = new Set<() => void>();
private readonly closeListeners = new Set<(error?: Error) => void>();
private started = false;
private sessionId?: string;
constructor(entries: NamedRelayTransport[]) {
if (entries.length === 0) throw new Error("CompositeRelayTransport requires at least one relay");
const ids = new Set<string>();
for (const entry of entries) {
if (!entry.id || ids.has(entry.id)) throw new Error(`duplicate relay id: ${entry.id || "(empty)"}`);
ids.add(entry.id);
}
this.entries = [...entries];
}
setSessionId(sessionId: string): void {
this.sessionId = sessionId;
for (const { transport } of this.entries) {
(transport as RelayTransport & { setSessionId?: (value: string) => void }).setSessionId?.(sessionId);
}
}
async connect(): Promise<void> {
if (this.started) return;
this.started = true;
this.bindTransports();
if (this.sessionId) this.setSessionId(this.sessionId);
const results = await Promise.allSettled(this.entries.map(({ transport }) => transport.connect()));
const failures = results
.map((result, index) => ({ result, entry: this.entries[index] }))
.filter((item): item is { result: PromiseRejectedResult; entry: NamedRelayTransport } => item.result.status === "rejected");
const requiredFailure = failures.find(({ entry }) => entry.required);
const connected = results.length - failures.length;
if (requiredFailure || connected === 0) {
this.started = false;
this.disposeSubscriptions();
for (const { transport } of this.entries) transport.close();
const detail = failures.map(({ entry, result }) => `${entry.id}: ${errorMessage(result.reason)}`).join("; ");
throw new Error(`unable to connect relay${failures.length === 1 ? "" : "s"}: ${detail}`);
}
}
send(frame: RelayFrame): void {
const failures: string[] = [];
for (const { id, transport } of this.entries) {
try {
transport.send(frame);
} catch (error) {
failures.push(`${id}: ${errorMessage(error)}`);
}
}
if (failures.length === this.entries.length) {
throw new Error(`all relay sends failed: ${failures.join("; ")}`);
}
}
onMessage(listener: (frame: RelayFrame) => void): Disposable {
this.messageListeners.add(listener);
return { dispose: () => this.messageListeners.delete(listener) };
}
onOpen(listener: () => void): Disposable {
this.openListeners.add(listener);
return { dispose: () => this.openListeners.delete(listener) };
}
onClose(listener: (error?: Error) => void): Disposable {
this.closeListeners.add(listener);
return { dispose: () => this.closeListeners.delete(listener) };
}
isConnected(id: string): boolean {
return this.openEntries.has(id);
}
close(): void {
this.started = false;
this.openEntries.clear();
this.disposeSubscriptions();
for (const { transport } of this.entries) transport.close();
}
private bindTransports(): void {
for (const { id, transport } of this.entries) {
this.subscriptions.push(transport.onMessage((frame) => {
for (const listener of this.messageListeners) listener(frame);
}));
if (transport.onOpen) this.subscriptions.push(transport.onOpen(() => {
this.openEntries.add(id);
// RelayHost publishes an authoritative snapshot after an authenticated
// reconnect. Surface every member reconnect so a recovered cloud relay
// is hydrated even while the local relay remained online.
for (const listener of this.openListeners) listener();
}));
if (transport.onClose) this.subscriptions.push(transport.onClose((error) => {
const wasOpen = this.openEntries.delete(id);
if (wasOpen && this.openEntries.size === 0) {
for (const listener of this.closeListeners) listener(error);
}
}));
}
}
private disposeSubscriptions(): void {
for (const subscription of this.subscriptions.splice(0)) subscription.dispose();
}
}
function errorMessage(error: unknown): string {
return error instanceof Error ? error.message : String(error);
}
@@ -0,0 +1,574 @@
import { hostname } from "node:os";
import * as vscode from "vscode";
import { CodexAgentAdapter } from "./codexAgentAdapter";
import { resolveCodexCommand } from "./codexPath";
import { CodexIpcAgentAdapter } from "./codexIpcAgentAdapter";
import { CompositeRelayTransport } from "./compositeRelay";
import { LocalRelayController, localRelayTarget } from "./localRelay";
import { AgentAdapter, ControlMode, Disposable, JsonObject, Logger } from "./protocol";
import { RelayClient } from "./relayClient";
import { RelayHost } from "./relayHost";
import { SwitchableAgentAdapter } from "./switchableAgentAdapter";
let activeHost: RelayHost | undefined;
let activeAdapter: AgentAdapter | undefined;
let activeRelay: CompositeRelayTransport | undefined;
let activeAdapterStatusSubscription: Disposable | undefined;
let statusItem: vscode.StatusBarItem | undefined;
let autoStartRetryTimer: NodeJS.Timeout | undefined;
let autoStartRetryMs = 3_000;
let localRelayController: LocalRelayController | undefined;
const t = (message: string, ...args: Array<string | number | boolean>): string => vscode.l10n.t(message, ...args);
export async function activate(context: vscode.ExtensionContext): Promise<void> {
const output = vscode.window.createOutputChannel(t("Codex Remote Collaboration"));
context.subscriptions.push(output);
const logger = {
debug: (message: string, ...args: unknown[]) => output.appendLine(`[debug] ${message} ${formatArgs(args)}`),
info: (message: string, ...args: unknown[]) => output.appendLine(`[info] ${message} ${formatArgs(args)}`),
warn: (message: string, ...args: unknown[]) => output.appendLine(`[warn] ${message} ${formatArgs(args)}`),
error: (message: string, ...args: unknown[]) => output.appendLine(`[error] ${message} ${formatArgs(args)}`),
};
localRelayController = new LocalRelayController({ extensionPath: context.extensionPath, logger });
statusItem = vscode.window.createStatusBarItem(vscode.StatusBarAlignment.Right, 100);
statusItem.command = "codexRemoteCollab.openWeb";
statusItem.text = "$(plug) Codex Remote";
statusItem.tooltip = t("Connecting to the local Codex collaboration service");
statusItem.show();
context.subscriptions.push(statusItem);
const start = async (automatic = false): Promise<void> => {
if (!automatic && autoStartRetryTimer) {
clearTimeout(autoStartRetryTimer);
autoStartRetryTimer = undefined;
}
if (activeHost) {
if (!automatic) vscode.window.showInformationMessage(t("The Codex remote bridge is already running."));
return;
}
const configuration = vscode.workspace.getConfiguration("codexRemoteCollab");
const relayConfiguration = resolveRelayConfiguration(configuration);
const localRelayUrl = relayConfiguration.localUrl;
if (!localRelayUrl) {
vscode.window.showWarningMessage(t("Set codexRemoteCollab.localRelayUrl before starting the bridge."));
return;
}
const localTarget = localRelayTarget(localRelayUrl);
if (!localTarget) {
vscode.window.showErrorMessage(t("codexRemoteCollab.localRelayUrl must be a loopback ws:// address."));
return;
}
if (localTarget && configuration.get<boolean>("autoStartLocalRelay", true)) {
setStatus("$(sync~spin) Codex Remote", t("Starting {0}", localTarget.webUrl), "codexRemoteCollab.openWeb");
try {
await localRelayController?.ensureRunning(localRelayUrl);
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
setStatus("$(error) Codex Remote", t("Unable to start the local collaboration service: {0}", message), "codexRemoteCollab.openWeb");
if (automatic) {
logger.warn(`Automatic local relay start failed; retrying in ${autoStartRetryMs}ms`, error);
scheduleAutoStartRetry(start);
} else {
vscode.window.showErrorMessage(t("Unable to start the local Codex collaboration service: {0}", message));
}
return;
}
}
const legacyToken = await context.secrets.get("codexRemoteCollab.relayToken");
const localToken = relayConfiguration.legacyRemote ? undefined : legacyToken;
if (!localToken) logger.info("Using the loopback-only unauthenticated local relay");
const initialControlMode = resolveInitialControlMode(configuration);
let currentControlMode: ControlMode = initialControlMode;
const createAdapter = (controlMode: ControlMode): AgentAdapter => {
if (controlMode === "sync") {
const configuredThreadId = configuration.get<string>("threadId", "").trim();
const socketPath = configuration.get<string>("ipcSocketPath", "").trim();
logger.info(`Synchronous mode enabled; following the VS Code Codex panel${configuredThreadId ? ` (initial conversation ${configuredThreadId})` : ""}`);
return new CodexIpcAgentAdapter({
threadId: configuredThreadId || undefined,
socketPath: socketPath || undefined,
hostId: configuration.get<string>("hostId", "local"),
autoDiscoverThread: configuration.get<boolean>("autoDiscoverThread", true),
// Synchronous mode has one navigation owner: the official panel.
followVscodeSession: true,
preferredCwds: workspaceRoots(),
strictVersions: configuration.get<boolean>("ipcStrictVersions", true),
logger,
approvalTimeoutMs: configuration.get<number>("approvalTimeoutMs", 300_000),
openNewSession: () => openOfficialNewSession(logger),
});
}
const configuredCommand = configuration.get<string>("codexCommand", "codex");
const command = resolveCodexCommand(configuredCommand);
const args = configuration.get<string[]>("codexArgs", ["app-server", "--stdio"]);
const defaultCwd = configuration.get<string>("defaultCwd", "") || firstWorkspaceRoot();
logger.info(`Asynchronous mode enabled; using independent Codex executable: ${command}`);
return new CodexAgentAdapter({
command,
args,
defaultCwd: defaultCwd || undefined,
logger,
approvalTimeoutMs: configuration.get<number>("approvalTimeoutMs", 300_000),
});
};
const adapter = new SwitchableAgentAdapter({
initialMode: initialControlMode,
createAdapter,
logger,
onModeChanged: async (nextMode) => {
currentControlMode = nextMode;
await configuration.update("controlMode", nextMode, vscode.ConfigurationTarget.Global);
setControlModeStatus(nextMode, nextMode === "async" || Boolean((await adapter.snapshot()).threadId));
},
});
const localRelay = new RelayClient({
url: localRelayUrl,
...(localToken ? { accessToken: localToken } : {}),
reconnect: configuration.get<boolean>("relayReconnect", true),
logger,
});
const relayEntries = [{ id: "local", transport: localRelay, required: true }];
const cloudRelayUrl = relayConfiguration.cloudUrl;
const cloudToken = await context.secrets.get("codexRemoteCollab.cloudRelayToken")
?? (relayConfiguration.legacyRemote ? legacyToken : undefined);
if (cloudRelayUrl && cloudToken) {
relayEntries.push({
id: "aether-cloud",
transport: new RelayClient({
url: cloudRelayUrl,
accessToken: cloudToken,
reconnect: configuration.get<boolean>("relayReconnect", true),
logger,
}),
required: false,
});
logger.info(`Aether cloud relay enabled: ${cloudRelayUrl}`);
} else if (cloudRelayUrl) {
logger.warn("Aether cloud relay URL is configured without a device credential; cloud sync is disabled until pairing is completed");
}
const relay = new CompositeRelayTransport(relayEntries);
const capabilities = ["read_output", "send_task_input", "cancel_task", "approve_low_risk"];
if (configuration.get<boolean>("allowHighRiskApprovals", false)) capabilities.push("approve_high_risk");
const host = new RelayHost({ adapter, relay, logger, capabilities });
let controlReady = initialControlMode === "async";
activeAdapterStatusSubscription?.dispose();
const adapterStatusSubscription = adapter.onEvent((event) => {
if (activeAdapter !== adapter) return;
if (event.type === "control.mode.changed") {
const changedMode = event.payload.controlMode;
if (changedMode === "sync" || changedMode === "async") currentControlMode = changedMode;
}
if (event.type !== "session.snapshot") return;
const metadata = event.payload.metadata;
if (metadata !== null && typeof metadata === "object" && !Array.isArray(metadata)) {
const snapshotMode = (metadata as JsonObject).controlMode;
if (snapshotMode === "sync" || snapshotMode === "async") currentControlMode = snapshotMode;
}
const waiting = event.payload.state === "waiting_for_host"
|| (metadata !== null && typeof metadata === "object" && !Array.isArray(metadata)
&& (metadata as JsonObject).waitingForSession === true);
const threadId = event.threadId
?? (typeof event.payload.threadId === "string" ? event.payload.threadId : undefined);
controlReady = currentControlMode === "async" || (Boolean(threadId) && !waiting);
setControlModeStatus(currentControlMode, controlReady);
});
activeAdapterStatusSubscription = adapterStatusSubscription;
activeAdapter = adapter;
activeRelay = relay;
activeHost = host;
if (configuration.get<boolean>("autoStartLocalRelay", true)) {
localRelay.onClose(() => {
if (activeHost !== host) return;
setStatus("$(sync~spin) Codex Remote", t("Restoring the local collaboration service"), "codexRemoteCollab.openWeb");
void localRelayController?.ensureRunning(localRelayUrl).catch((error) => {
logger.warn("Unable to recover bundled local relay", error);
setStatus("$(error) Codex Remote", t("Unable to restore the local collaboration service"), "codexRemoteCollab.openWeb");
});
});
localRelay.onOpen(() => {
if (activeHost === host) {
setControlModeStatus(currentControlMode, controlReady);
}
});
}
try {
await host.start();
autoStartRetryMs = 3_000;
const snapshot = await adapter.snapshot();
const snapshotMode = snapshot.metadata?.controlMode;
if (snapshotMode === "sync" || snapshotMode === "async") currentControlMode = snapshotMode;
controlReady = currentControlMode === "async" || Boolean(snapshot.threadId);
setControlModeStatus(currentControlMode, controlReady);
if (!automatic && (currentControlMode === "async" || controlReady)) {
vscode.window.showInformationMessage(currentControlMode === "sync"
? t("The Codex remote bridge attached to the existing VS Code Codex conversation.")
: t("The independent Codex remote mode connected."));
}
} catch (error) {
activeHost = undefined;
activeAdapter = undefined;
activeRelay = undefined;
if (activeAdapterStatusSubscription === adapterStatusSubscription) {
activeAdapterStatusSubscription.dispose();
activeAdapterStatusSubscription = undefined;
}
await host.stop().catch(() => undefined);
const message = error instanceof Error ? error.message : String(error);
if (initialControlMode === "sync" && isAttachSessionUnavailable(message)) {
setStatus("$(sync~spin) Codex Remote", t("Waiting for a Codex conversation to open in VS Code. It will connect automatically."), "codexRemoteCollab.openWeb");
logger.info(`No attachable VS Code Codex session is available; retrying in ${autoStartRetryMs}ms`);
scheduleAutoStartRetry(start);
return;
}
setStatus("$(error) Codex Remote", initialControlMode === "sync" ? t("The Codex conversation is not connected") : t("The independent Codex mode is not connected"), "codexRemoteCollab.openWeb");
if (automatic) {
logger.warn(`Automatic bridge start failed; retrying in ${autoStartRetryMs}ms`, error);
scheduleAutoStartRetry(start);
} else {
const detail = localTarget && /ECONNREFUSED|connect refused/i.test(message)
? t("The local collaboration service at {0} is temporarily unavailable. The extension will keep retrying.", localTarget.webUrl)
: t("Unable to start the Codex remote bridge: {0}", message);
vscode.window.showErrorMessage(detail);
}
}
};
const stop = async (): Promise<void> => {
if (autoStartRetryTimer) {
clearTimeout(autoStartRetryTimer);
autoStartRetryTimer = undefined;
}
autoStartRetryMs = 3_000;
const host = activeHost;
activeHost = undefined;
activeAdapter = undefined;
activeRelay = undefined;
activeAdapterStatusSubscription?.dispose();
activeAdapterStatusSubscription = undefined;
if (host) await host.stop();
setStatus("$(plug) Codex Remote", t("Bridge paused. Click to open the web control and resume automatically."), "codexRemoteCollab.openWeb");
};
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.openWeb", async () => {
const localRelayUrl = resolveRelayConfiguration(vscode.workspace.getConfiguration("codexRemoteCollab")).localUrl;
const webUrl = localRelayController?.getWebUrl(localRelayUrl);
if (!webUrl) {
vscode.window.showErrorMessage(t("The local collaboration URL is invalid. Check codexRemoteCollab.localRelayUrl."));
return;
}
if (!activeHost) await start(false);
if (activeHost) await vscode.env.openExternal(vscode.Uri.parse(webUrl));
}));
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.start", start));
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.stop", stop));
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.setThreadId", async () => {
const configuration = vscode.workspace.getConfiguration("codexRemoteCollab");
const current = configuration.get<string>("threadId", "");
const value = await vscode.window.showInputBox({
prompt: t("Existing Codex conversation ID (leave blank for auto-discovery)"),
value: current,
ignoreFocusOut: true,
});
if (value === undefined) return;
await configuration.update("threadId", value.trim(), vscode.ConfigurationTarget.Global);
vscode.window.showInformationMessage(value.trim()
? t("Codex Remote will attach to {0} after the next bridge start.", value.trim())
: t("Codex Remote will auto-discover the latest VS Code Codex conversation after the next bridge start."));
}));
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.setRelayToken", async () => {
const token = await vscode.window.showInputBox({ prompt: t("Relay access token (leave blank for the local relay)"), password: true, ignoreFocusOut: true });
if (token === undefined) return;
await context.secrets.store("codexRemoteCollab.relayToken", token);
vscode.window.showInformationMessage(t("Relay token stored in VS Code SecretStorage."));
}));
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.configureCloud", async () => {
const configuration = vscode.workspace.getConfiguration("codexRemoteCollab");
const currentUrl = configuration.get<string>("cloudRelayUrl", "");
const url = await vscode.window.showInputBox({
prompt: t("Aether cloud relay WebSocket URL"),
value: currentUrl,
placeHolder: "wss://aether.example.com/api/vscodex/ws",
ignoreFocusOut: true,
validateInput: validateCloudRelayUrl,
});
if (url === undefined) return;
if (!url.trim()) {
await configuration.update("cloudRelayUrl", "", vscode.ConfigurationTarget.Global);
await context.secrets.delete("codexRemoteCollab.cloudRelayToken");
vscode.window.showInformationMessage(t("Aether cloud connection removed. Local control remains enabled."));
return;
}
const token = await vscode.window.showInputBox({
prompt: t("Device credential from the Aether pairing flow"),
password: true,
ignoreFocusOut: true,
});
if (token === undefined) return;
if (!token.trim()) {
vscode.window.showWarningMessage(t("A non-empty Aether device credential is required."));
return;
}
await configuration.update("cloudRelayUrl", url.trim(), vscode.ConfigurationTarget.Global);
await context.secrets.store("codexRemoteCollab.cloudRelayToken", token.trim());
vscode.window.showInformationMessage(t("Aether cloud connection saved. Restart the Codex Remote bridge to connect; local control remains available."));
}));
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.pairCloud", async () => {
const configuration = vscode.workspace.getConfiguration("codexRemoteCollab");
const currentBaseUrl = configuration.get<string>("aetherUrl", "");
const baseUrl = await vscode.window.showInputBox({
prompt: t("Aether server URL"),
value: currentBaseUrl,
placeHolder: "https://aether.example.com",
ignoreFocusOut: true,
validateInput: validateAetherBaseUrl,
});
if (baseUrl === undefined || !baseUrl.trim()) return;
const code = await vscode.window.showInputBox({
prompt: t("One-time pairing code shown in Aether"),
placeHolder: "ABCD-EFGH",
ignoreFocusOut: true,
validateInput: (value) => normalizePairingCode(value).length === 8 ? undefined : t("Enter the 8-character pairing code."),
});
if (code === undefined || !code.trim()) return;
try {
const normalizedBaseUrl = baseUrl.trim().replace(/\/+$/, "");
const response = await fetch(`${normalizedBaseUrl}/api/vscodex/pair`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ code: normalizePairingCode(code), name: hostname() || "VS Code" }),
});
const raw = await response.text();
let result: unknown;
try {
result = JSON.parse(raw);
} catch {
result = null;
}
if (!response.ok) {
const detail = isJsonRecord(result) && typeof result.error === "string" ? result.error : `HTTP ${response.status}`;
throw new Error(detail);
}
if (!isJsonRecord(result) || typeof result.device_token !== "string" || typeof result.ws_url !== "string") {
throw new Error(t("Aether returned an invalid pairing response."));
}
const wsError = validateCloudRelayUrl(result.ws_url);
if (wsError) throw new Error(wsError);
await configuration.update("aetherUrl", normalizedBaseUrl, vscode.ConfigurationTarget.Global);
await configuration.update("cloudRelayUrl", result.ws_url, vscode.ConfigurationTarget.Global);
await context.secrets.store("codexRemoteCollab.cloudRelayToken", result.device_token);
if (activeHost) await stop();
await start(false);
if (!activeHost) return;
if (activeRelay?.isConnected("aether-cloud")) {
vscode.window.showInformationMessage(t("Aether pairing completed. Local and cloud control are both active."));
} else {
vscode.window.showWarningMessage(t("Aether pairing was saved, but the cloud connection is currently unavailable. Local control remains active and the cloud connection will retry."));
}
} catch (error) {
vscode.window.showErrorMessage(t("Unable to pair with Aether: {0}", error instanceof Error ? error.message : String(error)));
}
}));
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.sendInput", async () => {
if (!activeAdapter) {
vscode.window.showWarningMessage(t("Start the Codex remote bridge first."));
return;
}
const text = await vscode.window.showInputBox({ prompt: t("Send input to the active Codex turn"), ignoreFocusOut: true });
if (text === undefined || !text.trim()) return;
try {
await activeAdapter.sendInput(text);
} catch (error) {
vscode.window.showErrorMessage(t("Unable to send Codex input: {0}", error instanceof Error ? error.message : String(error)));
}
}));
context.subscriptions.push(vscode.commands.registerCommand("codexRemoteCollab.snapshot", async () => {
if (!activeAdapter) return vscode.window.showWarningMessage(t("Start the Codex remote bridge first."));
const snapshot = await activeAdapter.snapshot();
output.appendLine(JSON.stringify(snapshot));
output.show(true);
}));
if (vscode.workspace.getConfiguration("codexRemoteCollab").get<boolean>("autoStart", true)) await start(true);
}
export async function deactivate(): Promise<void> {
if (autoStartRetryTimer) {
clearTimeout(autoStartRetryTimer);
autoStartRetryTimer = undefined;
}
const host = activeHost;
activeHost = undefined;
activeAdapter = undefined;
activeAdapterStatusSubscription?.dispose();
activeAdapterStatusSubscription = undefined;
if (host) await host.stop();
const localRelay = localRelayController;
localRelayController = undefined;
if (localRelay) await localRelay.stop();
}
function firstWorkspaceRoot(): string | undefined {
return vscode.workspace.workspaceFolders?.[0]?.uri.fsPath;
}
function workspaceRoots(): string[] {
return (vscode.workspace.workspaceFolders ?? []).map((folder) => folder.uri.fsPath);
}
function setStatus(text: string, tooltip: string, command?: string): void {
if (!statusItem) return;
statusItem.text = text;
statusItem.tooltip = tooltip;
statusItem.command = command;
}
function setAttachStatus(ready: boolean): void {
setStatus(
ready ? "$(check) Codex Remote" : "$(sync~spin) Codex Remote",
ready ? t("Attached to the existing Codex conversation. Click to open the web control.") : t("Waiting for a Codex conversation to open in VS Code. It will connect automatically."),
"codexRemoteCollab.openWeb",
);
}
function setControlModeStatus(mode: ControlMode, ready: boolean): void {
if (mode === "sync") {
setAttachStatus(ready);
return;
}
setStatus(
ready ? "$(check) Codex Remote" : "$(sync~spin) Codex Remote",
ready
? t("Independent Codex mode is connected. Click to open the web control.")
: t("Starting the independent Codex mode."),
"codexRemoteCollab.openWeb",
);
}
function scheduleAutoStartRetry(start: (automatic?: boolean) => Promise<void>): void {
if (autoStartRetryTimer) return;
const delay = autoStartRetryMs;
autoStartRetryMs = Math.min(autoStartRetryMs * 2, 30_000);
autoStartRetryTimer = setTimeout(() => {
autoStartRetryTimer = undefined;
void start(true);
}, delay);
}
function formatArgs(args: unknown[]): string {
return args.length ? args.map((arg) => (typeof arg === "string" ? arg : JSON.stringify(arg))).join(" ") : "";
}
function isAttachSessionUnavailable(message: string): boolean {
return message.includes("没有找到已打开的 VS Code Codex 会话")
|| /找不到会话\s+.+\s+的 VS Code Codex owner/.test(message);
}
function validateCloudRelayUrl(value: string): string | undefined {
if (!value.trim()) return undefined;
try {
const url = new URL(value.trim());
if (url.protocol !== "wss:" && url.protocol !== "ws:") return t("Use a ws:// or wss:// URL.");
if (url.protocol === "ws:" && !isLoopbackHostname(url.hostname)) {
return t("Remote Aether connections must use wss://.");
}
return undefined;
} catch {
return t("Enter a valid WebSocket URL.");
}
}
function resolveRelayConfiguration(configuration: vscode.WorkspaceConfiguration): {
localUrl: string;
cloudUrl: string;
legacyRemote: boolean;
} {
const defaultLocalUrl = "ws://127.0.0.1:8787/v1/connect";
const explicitLocal = inspectedValue<string>(configuration.inspect<string>("localRelayUrl"));
const explicitCloud = inspectedValue<string>(configuration.inspect<string>("cloudRelayUrl"));
const explicitLegacy = inspectedValue<string>(configuration.inspect<string>("relayUrl"));
const legacyUrl = explicitLegacy?.trim() || "";
const legacyRemote = Boolean(legacyUrl && !localRelayTarget(legacyUrl));
const localUrl = (explicitLocal?.trim()
|| (!legacyRemote ? legacyUrl : "")
|| configuration.get<string>("localRelayUrl", defaultLocalUrl).trim()
|| defaultLocalUrl);
const cloudUrl = explicitCloud?.trim()
|| (legacyRemote ? legacyUrl : "")
|| configuration.get<string>("cloudRelayUrl", "").trim();
return { localUrl, cloudUrl, legacyRemote };
}
function inspectedValue<T>(inspection: ReturnType<vscode.WorkspaceConfiguration["inspect"]> | undefined): T | undefined {
if (!inspection) return undefined;
const values = inspection as {
globalLanguageValue?: T;
workspaceFolderLanguageValue?: T;
workspaceLanguageValue?: T;
workspaceFolderValue?: T;
workspaceValue?: T;
globalValue?: T;
};
return values.workspaceFolderLanguageValue
?? values.workspaceLanguageValue
?? values.globalLanguageValue
?? values.workspaceFolderValue
?? values.workspaceValue
?? values.globalValue;
}
function resolveInitialControlMode(configuration: vscode.WorkspaceConfiguration): ControlMode {
const configured = inspectedValue<ControlMode>(configuration.inspect<ControlMode>("controlMode"));
if (configured === "sync" || configured === "async") return configured;
const legacyMode = inspectedValue<"attach" | "spawn">(configuration.inspect<"attach" | "spawn">("mode"));
return legacyMode === "spawn" ? "async" : "sync";
}
function validateAetherBaseUrl(value: string): string | undefined {
if (!value.trim()) return t("Enter the Aether server URL.");
try {
const url = new URL(value.trim());
if (url.username || url.password || url.search || url.hash) return t("Use the Aether origin without credentials, a query, or a fragment.");
if (url.protocol === "https:") return undefined;
if (url.protocol === "http:" && isLoopbackHostname(url.hostname)) return undefined;
return t("Remote Aether servers must use https://.");
} catch {
return t("Enter a valid URL.");
}
}
function normalizePairingCode(value: string): string {
return value.toUpperCase().replace(/[^A-Z2-9]/g, "");
}
function isLoopbackHostname(value: string): boolean {
const hostname = value.replace(/^\[|\]$/g, "").toLowerCase();
return hostname === "127.0.0.1" || hostname === "localhost" || hostname === "::1";
}
function isJsonRecord(value: unknown): value is Record<string, unknown> {
return typeof value === "object" && value !== null && !Array.isArray(value);
}
/**
* Reuse the official extension's command registry for the header's new-chat
* action. This keeps the remote UI attached to the same VS Code Codex
* installation and avoids launching a second app-server process.
*/
async function openOfficialNewSession(logger: Logger): Promise<JsonObject> {
const commands = await vscode.commands.getCommands(true);
const command = commands.includes("chatgpt.newCodexPanel")
? "chatgpt.newCodexPanel"
: commands.includes("chatgpt.newChat")
? "chatgpt.newChat"
: undefined;
if (!command) {
throw new Error(t("The official Codex extension new-conversation command was not found. Make sure the VS Code Codex extension is enabled."));
}
await vscode.commands.executeCommand(command);
logger.info?.("Opened a new official Codex conversation with " + command);
return { opened: true, command };
}
@@ -0,0 +1,11 @@
export * from "./protocol";
export * from "./jsonlRpc";
export * from "./codexAgentAdapter";
export * from "./relayClient";
export * from "./compositeRelay";
export * from "./relayHost";
export * from "./bridge";
export * from "./codexPath";
export * from "./codexIpc";
export * from "./codexIpcAgentAdapter";
export * from "./switchableAgentAdapter";
@@ -0,0 +1,268 @@
import { ChildProcessWithoutNullStreams, spawn } from "node:child_process";
import { createInterface, Interface as ReadLineInterface } from "node:readline";
import {
asJsonValue,
Disposable,
isRecord,
JsonRpcId,
JsonRpcNotification,
JsonRpcRequest,
JsonValue,
Logger,
jsonRpcIdKey,
} from "./protocol";
export interface JsonlRpcClientOptions {
command?: string;
args?: string[];
cwd?: string;
env?: NodeJS.ProcessEnv;
logger?: Logger;
/** Optional request timeout. Zero disables it, which is useful for long turns. */
requestTimeoutMs?: number;
}
interface PendingRequest {
method: string;
resolve: (value: JsonValue) => void;
reject: (reason: Error) => void;
timer?: NodeJS.Timeout;
}
export class JsonRpcRemoteError extends Error {
constructor(
message: string,
readonly code: number,
readonly data?: JsonValue,
) {
super(message);
this.name = "JsonRpcRemoteError";
}
}
/** Minimal newline-delimited JSON-RPC client used by `codex app-server --stdio`. */
export class JsonlRpcClient {
private readonly options: Required<Pick<JsonlRpcClientOptions, "command" | "args" | "requestTimeoutMs">> &
Omit<JsonlRpcClientOptions, "command" | "args" | "requestTimeoutMs">;
private child?: ChildProcessWithoutNullStreams;
private stdoutLines?: ReadLineInterface;
private nextId = 1;
private readonly pending = new Map<string, PendingRequest>();
private readonly notificationListeners = new Set<(message: JsonRpcNotification) => void>();
private readonly requestListeners = new Set<(message: JsonRpcRequest) => void>();
private readonly exitListeners = new Set<(error?: Error) => void>();
constructor(options: JsonlRpcClientOptions = {}) {
this.options = {
command: options.command ?? "codex",
args: options.args ?? ["app-server", "--stdio"],
requestTimeoutMs: options.requestTimeoutMs ?? 0,
cwd: options.cwd,
env: options.env,
logger: options.logger,
};
}
get running(): boolean {
return Boolean(this.child && this.child.exitCode === null && !this.child.killed);
}
async start(): Promise<void> {
if (this.running) return;
const child = spawn(this.options.command, this.options.args, {
cwd: this.options.cwd,
env: { ...process.env, ...(this.options.env ?? {}) },
stdio: ["pipe", "pipe", "pipe"],
windowsHide: true,
});
this.child = child;
this.stdoutLines = createInterface({ input: child.stdout, crlfDelay: Infinity });
this.stdoutLines.on("line", (line) => this.handleLine(line));
child.stderr.on("data", (chunk: Buffer) => {
const text = redactDiagnostic(chunk.toString("utf8").trim());
if (text) this.options.logger?.debug?.(`[app-server stderr] ${text}`);
});
child.once("exit", (code, signal) => {
const expected = child.killed;
const error = expected
? undefined
: new Error(`codex app-server exited (code=${String(code)}, signal=${String(signal)})`);
this.handleExit(error);
});
await new Promise<void>((resolve, reject) => {
const onSpawn = (): void => {
child.off("error", onError);
resolve();
};
const onError = (error: Error): void => {
child.off("spawn", onSpawn);
const spawnError = error as NodeJS.ErrnoException;
if (spawnError.code === "ENOENT") {
reject(new Error(`Codex executable "${this.options.command}" was not found. Set codexRemoteCollab.codexCommand to its full path.`));
return;
}
reject(error);
};
child.once("spawn", onSpawn);
child.once("error", onError);
});
}
request(method: string, params?: JsonValue): Promise<JsonValue> {
if (!this.running) return Promise.reject(new Error("app-server is not running"));
const id = this.nextId++;
return new Promise<JsonValue>((resolve, reject) => {
const pending: PendingRequest = { method, resolve, reject };
if (this.options.requestTimeoutMs > 0) {
pending.timer = setTimeout(() => {
this.pending.delete(jsonRpcIdKey(id));
reject(new Error(`app-server request timed out: ${method}`));
}, this.options.requestTimeoutMs);
}
this.pending.set(jsonRpcIdKey(id), pending);
try {
this.write({ id, method, ...(params === undefined ? {} : { params }) });
} catch (error) {
this.pending.delete(jsonRpcIdKey(id));
if (pending.timer) clearTimeout(pending.timer);
reject(error instanceof Error ? error : new Error(String(error)));
}
});
}
notify(method: string, params?: JsonValue): void {
this.write({ method, ...(params === undefined ? {} : { params }) });
}
respond(id: JsonRpcId, result: JsonValue): void {
this.write({ id, result });
}
respondError(id: JsonRpcId, code: number, message: string, data?: JsonValue): void {
this.write({ id, error: { code, message, ...(data === undefined ? {} : { data }) } });
}
onNotification(listener: (message: JsonRpcNotification) => void): Disposable {
this.notificationListeners.add(listener);
return { dispose: () => this.notificationListeners.delete(listener) };
}
onServerRequest(listener: (message: JsonRpcRequest) => void): Disposable {
this.requestListeners.add(listener);
return { dispose: () => this.requestListeners.delete(listener) };
}
onExit(listener: (error?: Error) => void): Disposable {
this.exitListeners.add(listener);
return { dispose: () => this.exitListeners.delete(listener) };
}
close(): void {
const child = this.child;
this.child = undefined;
this.stdoutLines?.close();
this.stdoutLines = undefined;
if (child && child.exitCode === null && !child.killed) child.kill();
this.rejectAll(new Error("app-server client closed"));
}
private write(message: unknown): void {
const child = this.child;
if (!child || child.exitCode !== null || child.killed || !child.stdin.writable) {
throw new Error("app-server is not running");
}
child.stdin.write(`${JSON.stringify(message)}\n`, "utf8");
}
private handleLine(line: string): void {
const trimmed = line.trim();
if (!trimmed) return;
let message: unknown;
try {
message = JSON.parse(trimmed);
} catch (error) {
this.options.logger?.warn?.("Ignoring malformed app-server JSON", error, trimmed.slice(0, 500));
return;
}
if (!isRecord(message)) return;
const hasId = typeof message.id === "string" || typeof message.id === "number";
const hasMethod = typeof message.method === "string";
if (hasId && (Object.hasOwn(message, "result") || Object.hasOwn(message, "error")) && !hasMethod) {
this.handleResponse(message as Record<string, unknown> & { id: JsonRpcId });
return;
}
if (hasMethod && hasId) {
const request: JsonRpcRequest = {
id: message.id as JsonRpcId,
method: message.method as string,
...(message.params === undefined ? {} : { params: asJsonValue(message.params) }),
};
for (const listener of this.requestListeners) listener(request);
return;
}
if (hasMethod) {
const notification: JsonRpcNotification = {
method: message.method as string,
...(message.params === undefined ? {} : { params: asJsonValue(message.params) }),
};
for (const listener of this.notificationListeners) listener(notification);
return;
}
this.options.logger?.warn?.("Ignoring unknown app-server message", message);
}
private handleResponse(message: Record<string, unknown> & { id: JsonRpcId }): void {
const pending = this.pending.get(jsonRpcIdKey(message.id));
if (!pending) {
this.options.logger?.warn?.(`Received response for unknown app-server request ${String(message.id)}`);
return;
}
this.pending.delete(jsonRpcIdKey(message.id));
if (pending.timer) clearTimeout(pending.timer);
if (isRecord(message.error)) {
pending.reject(
new JsonRpcRemoteError(
typeof message.error.message === "string" ? message.error.message : `Request failed: ${pending.method}`,
typeof message.error.code === "number" ? message.error.code : -32000,
message.error.data === undefined ? undefined : asJsonValue(message.error.data),
),
);
return;
}
pending.resolve(message.result === undefined ? null : asJsonValue(message.result));
}
private handleExit(error?: Error): void {
this.child = undefined;
this.stdoutLines?.close();
this.stdoutLines = undefined;
this.rejectAll(error ?? new Error("app-server exited"));
for (const listener of this.exitListeners) listener(error);
}
private rejectAll(error: Error): void {
for (const request of this.pending.values()) {
if (request.timer) clearTimeout(request.timer);
request.reject(error);
}
this.pending.clear();
}
}
function redactDiagnostic(text: string): string {
return text
.replace(/Bearer\s+[A-Za-z0-9._~+\-/]+=*/gi, "Bearer [REDACTED]")
.replace(/\b(?:sk-[A-Za-z0-9_-]{12,}|gh[pousr]_[A-Za-z0-9_]{12,})\b/g, "[REDACTED]")
.replace(/((?:token|secret|password|api[_-]?key)\s*[:=]\s*)[^\s,;]+/gi, "$1[REDACTED]");
}
@@ -0,0 +1,193 @@
import * as http from "node:http";
import * as path from "node:path";
import { Logger } from "./protocol";
interface BundledRelay {
start(): Promise<{ host: string; port: number }>;
stop(): Promise<void>;
}
interface BundledRelayModule {
CodexRelay: new (options: Record<string, unknown>) => BundledRelay;
}
export interface LocalRelayTarget {
host: string;
port: number;
healthUrl: string;
webUrl: string;
}
export interface LocalRelayControllerOptions {
extensionPath: string;
logger?: Logger;
probeTimeoutMs?: number;
relayModulePath?: string;
loadRelayModule?: (modulePath: string) => BundledRelayModule;
probeRelayHealth?: (url: string, timeoutMs?: number) => Promise<boolean>;
}
/**
* Owns the loopback relay bundled with the companion extension. Remote and
* TLS relay URLs deliberately stay outside this controller.
*/
export class LocalRelayController {
private readonly options: LocalRelayControllerOptions;
private relay?: BundledRelay;
private target?: LocalRelayTarget;
private starting?: Promise<boolean>;
private generation = 0;
constructor(options: LocalRelayControllerOptions) {
this.options = options;
}
async ensureRunning(relayUrl: string): Promise<boolean> {
const target = localRelayTarget(relayUrl);
if (!target) return false;
if ((this.relay || this.starting) && this.target?.healthUrl !== target.healthUrl) await this.stop();
const generation = this.generation;
this.target = target;
const available = await this.probeHealth(target.healthUrl);
// `stop()` may run while the health request is in flight. Do not let that
// completed probe resurrect a relay owned by a deactivated extension.
if (generation !== this.generation) return false;
if (available) return false;
if (this.relay) {
await this.relay.stop().catch(() => undefined);
this.relay = undefined;
}
if (this.starting) return this.starting;
this.starting = this.startBundledRelay(target).finally(() => {
this.starting = undefined;
});
return this.starting;
}
getWebUrl(relayUrl: string): string | undefined {
return localRelayTarget(relayUrl)?.webUrl;
}
async stop(): Promise<void> {
this.generation += 1;
const starting = this.starting;
if (starting) await starting.catch(() => undefined);
const relay = this.relay;
this.relay = undefined;
this.target = undefined;
if (relay) await relay.stop();
}
private async startBundledRelay(target: LocalRelayTarget): Promise<boolean> {
const modulePath = this.options.relayModulePath
?? path.join(this.options.extensionPath, "dist", "local-relay", "server.js");
let relay: BundledRelay;
try {
const load = this.options.loadRelayModule ?? ((value: string) => require(value) as BundledRelayModule);
const module = load(modulePath);
if (typeof module?.CodexRelay !== "function") throw new Error("bundled relay module is invalid");
relay = new module.CodexRelay({
host: target.host,
port: target.port,
mode: "host",
spawnCodex: false,
authRequired: false,
});
await relay.start();
} catch (error) {
// Another VS Code window can win the listen race after our health
// probe. Treat that as success only when the expected relay responds.
if (await this.probeHealth(target.healthUrl)) {
this.options.logger?.info?.(`Using existing local relay at ${target.webUrl}`);
return false;
}
throw error;
}
this.relay = relay;
this.options.logger?.info?.(`Started bundled local relay at ${target.webUrl}`);
return true;
}
private probeHealth(url: string): Promise<boolean> {
const probe = this.options.probeRelayHealth ?? relayHealthAvailable;
return probe(url, this.options.probeTimeoutMs);
}
}
export function localRelayTarget(relayUrl: string): LocalRelayTarget | undefined {
let url: URL;
try {
url = new URL(relayUrl);
} catch {
return undefined;
}
if (url.protocol !== "ws:" || !isLoopbackHostname(url.hostname)) return undefined;
const port = Number(url.port || 80);
if (!Number.isInteger(port) || port < 1 || port > 65_535) return undefined;
const hostname = normalizeLoopbackHostname(url.hostname);
const authorityHost = hostname.includes(":") ? `[${hostname}]` : hostname;
return {
host: hostname,
port,
healthUrl: `http://${authorityHost}:${port}/api/health`,
webUrl: `http://${authorityHost}:${port}/`,
};
}
function isLoopbackHostname(hostname: string): boolean {
const normalized = hostname.toLowerCase().replace(/^\[|\]$/g, "");
return normalized === "localhost" || normalized === "127.0.0.1" || normalized === "::1";
}
function normalizeLoopbackHostname(hostname: string): string {
const normalized = hostname.toLowerCase().replace(/^\[|\]$/g, "");
return normalized === "localhost" ? "127.0.0.1" : normalized;
}
export function relayHealthAvailable(url: string, timeoutMs = 700): Promise<boolean> {
return new Promise((resolve) => {
let settled = false;
let timer: NodeJS.Timeout | undefined;
const finish = (available: boolean): void => {
if (settled) return;
settled = true;
if (timer) clearTimeout(timer);
resolve(available);
};
const request = http.get(url, (response) => {
if (response.statusCode !== 200) {
response.resume();
finish(false);
return;
}
let body = "";
response.setEncoding("utf8");
response.on("data", (chunk) => {
if (body.length <= 16_384) body += chunk;
});
response.on("end", () => {
try {
const payload = JSON.parse(body) as { ok?: unknown };
finish(payload.ok === true);
} catch {
finish(false);
}
});
response.on("aborted", () => finish(false));
response.on("error", () => finish(false));
response.on("close", () => {
if (!response.complete) finish(false);
});
});
request.setTimeout(timeoutMs, () => {
request.destroy();
finish(false);
});
request.on("error", () => finish(false));
timer = setTimeout(() => {
request.destroy();
finish(false);
}, timeoutMs);
});
}
@@ -0,0 +1,496 @@
/**
* Wire types shared by the relay host and the Codex app-server adapter.
*
* The relay intentionally treats `payload` as JSON. Keeping this boundary
* unopinionated lets the bridge continue working when app-server adds a new
* notification or request before this extension is updated.
*/
export type JsonPrimitive = string | number | boolean | null;
export type JsonValue = JsonPrimitive | JsonValue[] | { [key: string]: JsonValue };
export type JsonObject = { [key: string]: JsonValue };
export type JsonRpcId = string | number;
/** Preserve the JSON-RPC id type when using it as a map key. */
export function jsonRpcIdKey(id: JsonRpcId): string {
return `${typeof id}:${String(id)}`;
}
export function isJsonRpcId(value: unknown): value is JsonRpcId {
return typeof value === "string" || typeof value === "number";
}
export type ApprovalDecisionKind = "allow" | "deny" | "cancel";
const LEGACY_APPROVAL_METHODS = new Set(["applyPatchApproval", "execCommandApproval"]);
const V2_APPROVAL_METHODS = new Set([
"item/commandExecution/requestApproval",
"item/fileChange/requestApproval",
]);
/**
* Classify both current and legacy app-server approval decisions without
* rewriting the wire value. Unknown tagged objects intentionally return
* `undefined` so callers can fail closed instead of accidentally approving a
* newly introduced response shape.
*/
export function approvalDecisionKind(value: unknown): ApprovalDecisionKind | undefined {
if (typeof value === "string") {
if (new Set([
"allow",
"accept",
"acceptForSession",
"approved",
"approved_for_session",
"approved_mcp_policy_amendment",
]).has(value)) return "allow";
if (new Set(["deny", "decline", "denied", "timed_out"]).has(value)) return "deny";
if (new Set(["cancel", "abort"]).has(value)) return "cancel";
return undefined;
}
if (!isRecord(value)) return undefined;
const keys = Object.keys(value);
if (keys.length !== 1) return undefined;
const key = keys[0];
const nested = value[key];
if (isExecpolicyAmendmentTag(key, nested) || isNetworkPolicyAmendmentTag(key, nested)) return "allow";
if (key === "denied" && isRecord(nested) && typeof nested.rejection === "string") return "deny";
return undefined;
}
/**
* Classify a decision against the response schema for one app-server method.
* The generic classifier above is intentionally useful for relay envelopes;
* this method-aware variant prevents a v2 tagged object from being sent to a
* legacy callback (or vice versa), while retaining compatibility aliases that
* the relay may use in its outer `decision` field.
*/
export function approvalDecisionKindForMethod(
value: unknown,
method?: string,
): ApprovalDecisionKind | undefined {
const generic = approvalDecisionKind(value);
if (!generic || !method) return generic;
if (LEGACY_APPROVAL_METHODS.has(method)) {
if (typeof value === "string") {
return new Set([
"approved",
"approved_for_session",
"approved_mcp_policy_amendment",
"timed_out",
"abort",
]).has(value) ? generic : undefined;
}
if (!isRecord(value)) return undefined;
const key = Object.keys(value)[0];
return key === "approved_execpolicy_amendment"
|| key === "network_policy_amendment"
|| key === "denied" ? generic : undefined;
}
if (V2_APPROVAL_METHODS.has(method)) {
if (typeof value === "string") {
return new Set(["accept", "acceptForSession", "decline", "cancel"]).has(value)
? generic
: undefined;
}
if (!isRecord(value)) return undefined;
const key = Object.keys(value)[0];
if (method === "item/fileChange/requestApproval") return undefined;
return key === "acceptWithExecpolicyAmendment" || key === "applyNetworkPolicyAmendment"
? generic
: undefined;
}
if (method === "mcpServer/elicitation/request") {
return typeof value === "string" && new Set(["accept", "decline", "cancel"]).has(value)
? generic
: undefined;
}
return generic;
}
function isExecpolicyAmendmentTag(key: string, nested: unknown): boolean {
if (!isRecord(nested)) return false;
if (key === "acceptWithExecpolicyAmendment") {
return isStringArray(nested.execpolicy_amendment);
}
if (key === "approved_execpolicy_amendment") {
return isStringArray(nested.proposed_execpolicy_amendment);
}
return false;
}
function isNetworkPolicyAmendmentTag(key: string, nested: unknown): boolean {
if (!isRecord(nested)) return false;
if (key === "applyNetworkPolicyAmendment") {
return isNetworkPolicyAmendment(nested.network_policy_amendment);
}
if (key === "network_policy_amendment") {
return isNetworkPolicyAmendment(nested.network_policy_amendment);
}
return false;
}
function isStringArray(value: unknown): value is string[] {
return Array.isArray(value) && value.every((item) => typeof item === "string");
}
function isNetworkPolicyAmendment(value: unknown): boolean {
return isRecord(value)
&& typeof value.host === "string"
&& (value.action === "allow" || value.action === "deny");
}
/** Whether a response explicitly carries a decision/action field. */
export function hasApprovalDecisionField(value: unknown): value is Record<string, unknown> {
return isRecord(value) && (Object.prototype.hasOwnProperty.call(value, "decision")
|| Object.prototype.hasOwnProperty.call(value, "action"));
}
export interface Disposable {
dispose(): void;
}
export interface JsonRpcRequest {
id: JsonRpcId;
method: string;
params?: JsonValue;
}
export interface JsonRpcNotification {
method: string;
params?: JsonValue;
}
export interface JsonRpcResponse {
id: JsonRpcId;
result?: JsonValue;
error?: {
code: number;
message: string;
data?: JsonValue;
};
}
export type JsonRpcMessage = JsonRpcRequest | JsonRpcNotification | JsonRpcResponse;
export type RelayRole = "owner" | "operator" | "approver" | "viewer" | string;
export interface RelayActor {
id?: string;
role?: RelayRole;
}
/** A versioned relay event frame. `seq` is normally assigned by the relay. */
export interface RelayEventFrame {
v: 1;
kind: "event";
type: string;
id: string;
sessionId: string;
seq?: number;
ts: string;
actor?: RelayActor;
payload: JsonObject;
/** Optional typed execution projection attached by a VS Code host. */
status?: AgentStatusSnapshot;
}
export interface RelayCommandFrame {
v?: 1;
kind?: "command";
type: string;
/** Compact relay compatibility form: `{ type: "command", method, params }`. */
method?: string;
params?: JsonObject;
commandId?: string;
id?: string;
sessionId?: string;
actor?: RelayActor;
payload?: JsonObject;
/** Some clients put the command body under `command`. */
command?: {
type?: string;
commandId?: string;
payload?: JsonObject;
[key: string]: JsonValue | undefined;
};
}
export interface RelayHelloFrame {
v: 1;
kind: "hello";
clientType: "host" | "web" | string;
protocol?: number;
accessToken?: string;
token?: string;
lastSeq?: number;
sessionId?: string;
payload?: JsonObject;
}
export interface RelayAckFrame {
v: 1;
kind: "ack";
sessionId: string;
seq: number;
}
export interface RelayErrorFrame {
v: 1;
kind: "error";
code: string;
message: string;
retryable?: boolean;
commandId?: string;
}
export type RelayFrame =
| RelayEventFrame
| RelayCommandFrame
| RelayHelloFrame
| RelayAckFrame
| RelayErrorFrame
| (JsonObject & { kind?: string; v?: number });
/**
* Live execution information projected from the official Codex conversation
* state. The private IPC protocol can add new turn statuses/flags, so the
* string fields intentionally remain open-ended for forward compatibility.
*/
export interface AgentStatusSnapshot {
/** Coarse UI activity, for example `thinking`, `editing`, or `running`. */
activity: string;
/** Raw/normalized turn status (`inProgress`, `completed`, ...). */
turnStatus: string;
/** Runtime flags such as `waitingOnApproval` or `waitingOnUserInput`. */
activeFlags: string[];
startedAtMs?: number | null;
durationMs?: number | null;
/**
* Time spent doing work in the official UI. This deliberately differs
* from `durationMs`: Codex starts the worked-for clock at the first work
* item and stops it when the final assistant response starts.
*/
workedDurationMs?: number | null;
/** Elapsed wall-clock time for an active turn. */
elapsedMs?: number | null;
firstTurnWorkItemStartedAtMs?: number | null;
finalAssistantStartedAtMs?: number | null;
error?: JsonValue;
}
/** Official background-agent lifecycle values emitted by Codex v2 items. */
export type CollabAgentStatus =
| "pendingInit"
| "running"
| "interrupted"
| "completed"
| "errored"
| "shutdown"
| "notFound"
| string;
export type CollabAgentTool =
| "spawnAgent"
| "sendInput"
| "resumeAgent"
| "wait"
| "closeAgent"
| string;
export type CollabAgentToolCallStatus = "inProgress" | "completed" | "failed" | string;
export type SubAgentActivityKind = "started" | "interacted" | "interrupted" | "completed" | string;
/** Last known state for one receiver in a collabAgentToolCall item. */
export interface CollabAgentStateSnapshot {
status: CollabAgentStatus;
message?: string | null;
}
/**
* Browser-safe projection of a background Codex subagent. The official
* webview currently uses the four coarse statuses below; the string union is
* deliberately open so a newer app-server status does not break the relay.
*/
export interface SubagentSnapshot {
threadId: string;
displayName: string | null;
prompt: string | null;
/** Alias used by the subagent side panel for the same prompt text. */
objective?: string | null;
status: "waiting" | "working" | "done" | "failed" | string;
statusMessage: string | null;
startedAtMs?: number | null;
completedAtMs?: number | null;
canInteract?: boolean;
model?: string | null;
agentPath?: string | null;
parentThreadId?: string | null;
}
export interface AgentEvent {
/** Normalized relay event name, for example `output.chunk`. */
type: string;
threadId?: string;
turnId?: string;
requestId?: JsonRpcId;
payload: JsonObject;
/** Original app-server notification/request, when available. */
raw?: JsonValue;
/** Optional typed projection of live Codex turn/runtime status. */
status?: AgentStatusSnapshot;
}
export interface PendingApproval {
requestId: JsonRpcId;
method: string;
threadId?: string;
turnId?: string;
itemId?: string;
action: string;
risk: "low" | "medium" | "high" | "unknown";
summary: string;
/** SHA-256 of canonicalized, unredacted app-server request params. */
commandHash?: string;
createdAt: number;
expiresAt?: number;
payload: JsonObject;
}
export interface SessionSnapshot {
threadId: string | null;
turnId: string | null;
state: string;
pendingApprovals: PendingApproval[];
pendingRequests?: Array<{
requestId: JsonRpcId;
method: string;
params?: JsonValue;
commandHash?: string;
risk?: string;
summary?: string;
createdAt?: number;
expiresAt?: number;
}>;
outputTail: string;
/** Optional role-aware projection used by the browser renderer. */
messages?: JsonValue[];
/** Background/inline subagents reconstructed from official collab items. */
subagents?: SubagentSnapshot[];
/** Live execution projection; retained alongside the legacy `state` field. */
status?: AgentStatusSnapshot;
/** Convenience aliases for clients that do not consume `status` yet. */
activity?: string;
turnStatus?: string;
activeFlags?: string[];
startedAtMs?: number | null;
durationMs?: number | null;
workedDurationMs?: number | null;
elapsedMs?: number | null;
metadata?: JsonObject;
}
/** A live VS Code Codex conversation that the attach bridge has verified. */
export interface SessionListEntry {
threadId: string;
title: string;
updatedAtMs: number | null;
cwd?: string | null;
active: boolean;
/** True for attach-mode results; retained for wire compatibility. */
available: boolean;
}
export interface SessionListResult {
sessions: SessionListEntry[];
activeThreadId: string | null;
}
/** Which owner controls conversation navigation for the remote surface. */
export type ControlMode = "sync" | "async";
export interface AgentAdapter {
start(): Promise<void>;
/** Switch between following VS Code and independently owned conversations. */
setControlMode?(params: JsonObject): Promise<JsonValue>;
/** Return the currently committed control mode without taking a snapshot. */
getControlMode?(): ControlMode;
/** Start a new app-server thread. */
startThread?(params?: JsonObject): Promise<JsonValue>;
/** Ask the official VS Code Codex extension to open a fresh conversation. */
newSession?(params?: JsonObject): Promise<JsonValue>;
/** Start a turn; `threadId` may be supplied in params or use the active thread. */
startTurn?(params: JsonObject): Promise<JsonValue>;
/** Steer the active turn. */
steerTurn?(params: JsonObject): Promise<JsonValue>;
/** Persist model/effort and other owner-managed settings on the thread. */
updateThreadSettings?(params: JsonObject): Promise<JsonValue>;
/** List verified, attachable local conversations without starting another Codex process. */
listSessions?(params?: JsonObject): Promise<JsonValue>;
/** Attach the follower to another already-open conversation. */
selectSession?(params: JsonObject): Promise<JsonValue>;
/** Interrupt a turn. */
interruptTurn?(params: JsonObject): Promise<JsonValue>;
/** Convenience MVP aliases. */
sendInput(text: string, params?: JsonObject): Promise<JsonValue>;
cancel(taskId?: string, params?: JsonObject): Promise<JsonValue>;
respondApproval(
requestId: JsonRpcId,
decision: "allow" | "deny" | "cancel",
reason?: string,
response?: JsonValue,
): Promise<JsonValue>;
/** Resolve all pending approvals/inputs with a deny response. */
denyPending?(reason?: string): Promise<void>;
snapshot(): Promise<SessionSnapshot>;
onEvent(listener: (event: AgentEvent) => void): Disposable;
dispose(): Promise<void>;
}
export interface RelayTransport {
connect(): Promise<void>;
send(frame: RelayFrame): void;
onMessage(listener: (frame: RelayFrame) => void): Disposable;
onOpen?(listener: () => void): Disposable;
onClose?(listener: (error?: Error) => void): Disposable;
close(): void;
}
export interface Logger {
debug?(message: string, ...args: unknown[]): void;
info?(message: string, ...args: unknown[]): void;
warn?(message: string, ...args: unknown[]): void;
error?(message: string, ...args: unknown[]): void;
}
export const noopDisposable = (): Disposable => ({ dispose: () => undefined });
export function isRecord(value: unknown): value is Record<string, unknown> {
return typeof value === "object" && value !== null && !Array.isArray(value);
}
export function asJsonObject(value: unknown): JsonObject {
return isRecord(value) ? (value as JsonObject) : {};
}
export function asJsonValue(value: unknown): JsonValue {
if (value === undefined) return null;
if (value === null || typeof value === "string" || typeof value === "number" || typeof value === "boolean") {
return value;
}
if (Array.isArray(value)) {
return value.map(asJsonValue);
}
if (isRecord(value)) {
const output: JsonObject = {};
for (const [key, item] of Object.entries(value)) {
if (item !== undefined) output[key] = asJsonValue(item);
}
return output;
}
return String(value);
}
@@ -0,0 +1,413 @@
import { createInterface, Interface as ReadLineInterface } from "node:readline";
import WebSocket from "ws";
import {
Disposable,
isRecord,
JsonObject,
Logger,
RelayFrame,
RelayHelloFrame,
RelayTransport,
} from "./protocol";
export interface RelayClientOptions {
url: string;
accessToken?: string;
sessionId?: string;
lastSeq?: number;
reconnect?: boolean;
reconnectInitialMs?: number;
reconnectMaxMs?: number;
maxFrameBytes?: number;
maxQueuedBytes?: number;
logger?: Logger;
/** Injectable constructor for tests or a browser-compatible WebSocket. */
webSocket?: new (url: string) => unknown;
}
type SocketLike = {
readyState?: number;
send(data: string): void;
close(): void;
on?(event: string, listener: (...args: any[]) => void): void;
addEventListener?(event: string, listener: (...args: any[]) => void): void;
};
const OPEN = 1;
// A structured Codex history snapshot is routinely larger than 256 KiB even
// though its plain-text tail is capped. Keep a bounded limit, but leave enough
// room for the message/tool projection of a long attached conversation.
export const DEFAULT_MAX_RELAY_FRAME_BYTES = 16 * 1024 * 1024;
export const DEFAULT_MAX_RELAY_QUEUE_BYTES = DEFAULT_MAX_RELAY_FRAME_BYTES + 2 * 1024 * 1024;
interface QueuedFrame {
serialized: string;
bytes: number;
projectionKey?: string;
}
const QUEUED_PROJECTION_TYPES = new Set(["session.snapshot", "output.snapshot", "output.chunk"]);
function queuedProjectionKey(frame: RelayFrame): string | undefined {
if (!isRecord(frame) || frame.kind !== "event" || typeof frame.type !== "string"
|| !QUEUED_PROJECTION_TYPES.has(frame.type)) return undefined;
const sessionId = typeof frame.sessionId === "string" ? frame.sessionId : "default";
return `${sessionId}:transcript`;
}
/** WebSocket relay transport with bounded reconnect and frame validation. */
export class RelayClient implements RelayTransport {
readonly handlesHandshake = true;
private readonly options: Required<
Pick<RelayClientOptions, "reconnect" | "reconnectInitialMs" | "reconnectMaxMs" | "maxFrameBytes" | "maxQueuedBytes">
> &
Omit<RelayClientOptions, "reconnect" | "reconnectInitialMs" | "reconnectMaxMs" | "maxFrameBytes" | "maxQueuedBytes">;
private socket?: SocketLike;
// A WebSocket can report OPEN while its relay authentication handshake is
// still in flight. Keep this separate from `socket` so events emitted by
// the adapter during reconnect are queued until the relay sends auth.ok.
private authenticatedSocket?: SocketLike;
private connecting?: Promise<void>;
// Incremented whenever a connection attempt is replaced or explicitly
// closed. Late events from an older WebSocket must not mutate newer state.
private connectionGeneration = 0;
private reconnectTimer?: NodeJS.Timeout;
private stopped = false;
private retryMs: number;
private readonly queue: QueuedFrame[] = [];
private queueBytes = 0;
private readonly listeners = new Set<(frame: RelayFrame) => void>();
private readonly openListeners = new Set<() => void>();
private readonly closeListeners = new Set<(error?: Error) => void>();
constructor(options: RelayClientOptions) {
const maxFrameBytes = options.maxFrameBytes ?? DEFAULT_MAX_RELAY_FRAME_BYTES;
const defaultMaxQueuedBytes = Math.max(
maxFrameBytes,
Math.min(DEFAULT_MAX_RELAY_QUEUE_BYTES, maxFrameBytes * 2),
);
this.options = {
...options,
reconnect: options.reconnect ?? true,
reconnectInitialMs: options.reconnectInitialMs ?? 500,
reconnectMaxMs: options.reconnectMaxMs ?? 10_000,
maxFrameBytes,
maxQueuedBytes: Math.max(1, Math.floor(options.maxQueuedBytes ?? defaultMaxQueuedBytes)),
};
this.retryMs = this.options.reconnectInitialMs;
}
/** Let RelayHost assign its stable session id before the first hello. */
setSessionId(sessionId: string): void {
this.options.sessionId = sessionId;
}
async connect(): Promise<void> {
this.stopped = false;
if (this.socket?.readyState === OPEN && this.authenticatedSocket === this.socket) return;
if (this.connecting) return this.connecting;
const generation = ++this.connectionGeneration;
let connectionPromise: Promise<void>;
connectionPromise = new Promise<void>((resolve, reject) => {
let settled = false;
let authenticated = false;
const SocketCtor = this.options.webSocket ?? WebSocket;
let socket: SocketLike;
try {
socket = new SocketCtor(this.options.url) as SocketLike;
} catch (error) {
reject(error instanceof Error ? error : new Error(String(error)));
return;
}
this.socket = socket;
this.authenticatedSocket = undefined;
const isCurrent = (): boolean => this.connectionGeneration === generation && this.socket === socket;
const onOpen = (): void => {
if (!isCurrent() || settled || authenticated) return;
try {
// The TCP/WebSocket open event is only a transport milestone. Do
// not release queued commands until the relay has authenticated us.
this.sendHello(socket);
} catch (error) {
if (!settled) {
settled = true;
reject(error instanceof Error ? error : new Error(String(error)));
}
}
};
const onMessage = (raw: unknown): void => {
if (!isCurrent()) return;
const data = extractMessageData(raw);
if (Buffer.byteLength(data, "utf8") > this.options.maxFrameBytes) {
this.options.logger?.warn?.("Ignoring oversized relay frame");
return;
}
let frame: unknown;
try {
frame = JSON.parse(data);
} catch {
this.options.logger?.warn?.("Ignoring malformed relay JSON");
return;
}
if (!isRecord(frame)) return;
if (frame.type === "auth.ok" && !authenticated && !settled) {
authenticated = true;
settled = true;
this.authenticatedSocket = socket;
this.retryMs = this.options.reconnectInitialMs;
try {
this.flush(socket);
} catch (error) {
this.options.logger?.warn?.("Unable to flush relay queue after authentication", error);
}
for (const listener of this.openListeners) listener();
resolve();
} else if (frame.type === "error" && !authenticated && !settled) {
settled = true;
reject(new Error(typeof frame.message === "string" ? frame.message : "relay authentication failed"));
}
if (frame.kind === "event" && typeof frame.seq === "number") {
this.options.lastSeq = Math.max(this.options.lastSeq ?? 0, frame.seq);
}
for (const listener of this.listeners) listener(frame as RelayFrame);
};
const onError = (raw: unknown): void => {
if (!isCurrent()) return;
const error = raw instanceof Error ? raw : new Error("relay websocket error");
this.options.logger?.warn?.(error.message);
if (!settled) {
settled = true;
reject(error);
}
};
const onClose = (): void => {
const current = isCurrent();
if (current) {
this.socket = undefined;
if (this.authenticatedSocket === socket) this.authenticatedSocket = undefined;
}
const error = new Error("relay websocket closed");
// A stale socket may still need to settle the promise returned to its
// caller, but it must never notify the active host or schedule a
// second reconnect loop.
if (!current) {
if (!settled) {
settled = true;
reject(error);
}
return;
}
for (const listener of this.closeListeners) listener(error);
if (!settled) {
settled = true;
reject(error);
}
if (!this.stopped && this.options.reconnect) this.scheduleReconnect();
};
bindSocket(socket, onOpen, onMessage, onError, onClose);
// A small number of test/browser WebSocket implementations can already
// be OPEN by the time listeners are attached.
if (socket.readyState === OPEN) queueMicrotask(onOpen);
}).finally(() => {
if (this.connectionGeneration === generation && this.connecting === connectionPromise) {
this.connecting = undefined;
}
});
this.connecting = connectionPromise;
return connectionPromise;
}
send(frame: RelayFrame): void {
const serialized = JSON.stringify(frame);
const bytes = Buffer.byteLength(serialized, "utf8");
if (bytes > this.options.maxFrameBytes) {
throw new Error(`relay frame exceeds ${this.options.maxFrameBytes} bytes`);
}
if (this.socket?.readyState === OPEN && this.authenticatedSocket === this.socket) {
this.socket.send(serialized);
return;
}
this.enqueue({ serialized, bytes, projectionKey: queuedProjectionKey(frame) });
}
onMessage(listener: (frame: RelayFrame) => void): Disposable {
this.listeners.add(listener);
return { dispose: () => this.listeners.delete(listener) };
}
onOpen(listener: () => void): Disposable {
this.openListeners.add(listener);
return { dispose: () => this.openListeners.delete(listener) };
}
onClose(listener: (error?: Error) => void): Disposable {
this.closeListeners.add(listener);
return { dispose: () => this.closeListeners.delete(listener) };
}
close(): void {
this.stopped = true;
this.connectionGeneration += 1;
if (this.reconnectTimer) clearTimeout(this.reconnectTimer);
this.reconnectTimer = undefined;
this.connecting = undefined;
const socket = this.socket;
this.socket = undefined;
this.authenticatedSocket = undefined;
if (socket && socket.readyState !== 3) socket.close();
this.queue.length = 0;
this.queueBytes = 0;
}
private sendHello(socket: SocketLike): void {
const hello: RelayHelloFrame = {
v: 1,
kind: "hello",
clientType: "host",
protocol: 1,
...(this.options.sessionId ? { sessionId: this.options.sessionId } : {}),
...(this.options.lastSeq !== undefined ? { lastSeq: this.options.lastSeq } : {}),
};
socket.send(JSON.stringify(hello));
if (this.options.accessToken) {
// Keep authentication separate from hello so a relay can challenge the
// host before accepting a bearer token (and so hello remains cacheable).
socket.send(JSON.stringify({ v: 1, kind: "auth", accessToken: this.options.accessToken }));
}
}
private flush(socket: SocketLike): void {
if (socket.readyState !== OPEN || this.authenticatedSocket !== socket || this.socket !== socket) return;
while (this.queue.length > 0) {
const entry = this.queue.shift() as QueuedFrame;
this.queueBytes = Math.max(0, this.queueBytes - entry.bytes);
socket.send(entry.serialized);
}
}
private enqueue(entry: QueuedFrame): void {
// Transcript events are reconstructible: RelayHost publishes a fresh full
// session snapshot after every authenticated reconnect. Keep only the
// newest projection per session while preserving approval/command events.
if (entry.projectionKey) {
for (let index = this.queue.length - 1; index >= 0; index -= 1) {
if (this.queue[index].projectionKey === entry.projectionKey) this.removeQueuedFrame(index);
}
}
if (entry.bytes > this.options.maxQueuedBytes) {
this.options.logger?.warn?.("Dropping relay frame that exceeds the reconnect queue byte limit");
return;
}
while (this.queue.length >= 100 || this.queueBytes + entry.bytes > this.options.maxQueuedBytes) {
const projectionIndex = this.queue.findIndex((queued) => Boolean(queued.projectionKey));
if (projectionIndex >= 0) {
this.removeQueuedFrame(projectionIndex);
continue;
}
// Never evict an approval/command solely to retain a transcript delta;
// the authoritative snapshot emitted after auth restores that state.
if (entry.projectionKey) {
this.options.logger?.debug?.("Dropping supersedable relay projection while reconnect queue is full");
return;
}
this.removeQueuedFrame(0);
}
this.queue.push(entry);
this.queueBytes += entry.bytes;
}
private removeQueuedFrame(index: number): void {
const [removed] = this.queue.splice(index, 1);
if (removed) this.queueBytes = Math.max(0, this.queueBytes - removed.bytes);
}
private scheduleReconnect(): void {
if (this.reconnectTimer || this.stopped) return;
const delay = this.retryMs;
this.retryMs = Math.min(this.options.reconnectMaxMs, Math.max(this.retryMs * 2, this.options.reconnectInitialMs));
this.reconnectTimer = setTimeout(() => {
this.reconnectTimer = undefined;
void this.connect().catch((error) => this.options.logger?.debug?.("relay reconnect failed", error));
}, delay);
}
}
/**
* Line-oriented transport for local development and CI. Pipe it to a relay
* process with `node dist/cli.js`; each line is one JSON relay frame.
*/
export class StdioRelayTransport implements RelayTransport {
private readonly listeners = new Set<(frame: RelayFrame) => void>();
private readonly lineReader: ReadLineInterface;
private closed = false;
constructor(
private readonly input: NodeJS.ReadableStream = process.stdin,
private readonly output: NodeJS.WritableStream = process.stdout,
private readonly logger?: Logger,
) {
this.lineReader = createInterface({ input, crlfDelay: Infinity });
this.lineReader.on("line", (line) => {
if (!line.trim()) return;
try {
const frame = JSON.parse(line);
if (isRecord(frame)) for (const listener of this.listeners) listener(frame as RelayFrame);
} catch (error) {
this.logger?.warn?.("Ignoring malformed relay stdin frame", error);
}
});
}
async connect(): Promise<void> {
this.closed = false;
}
send(frame: RelayFrame): void {
if (this.closed) throw new Error("stdio relay transport is closed");
this.output.write(`${JSON.stringify(frame)}\n`);
}
onMessage(listener: (frame: RelayFrame) => void): Disposable {
this.listeners.add(listener);
return { dispose: () => this.listeners.delete(listener) };
}
close(): void {
this.closed = true;
this.lineReader.close();
}
}
function bindSocket(
socket: SocketLike,
onOpen: () => void,
onMessage: (data: unknown) => void,
onError: (error: unknown) => void,
onClose: () => void,
): void {
if (typeof socket.on === "function") {
socket.on("open", onOpen);
socket.on("message", onMessage);
socket.on("error", onError);
socket.on("close", onClose);
} else if (typeof socket.addEventListener === "function") {
socket.addEventListener("open", onOpen);
socket.addEventListener("message", onMessage);
socket.addEventListener("error", onError);
socket.addEventListener("close", onClose);
} else {
onError(new Error("WebSocket implementation has no event API"));
}
}
function extractMessageData(raw: unknown): string {
if (typeof raw === "string") return raw;
if (Buffer.isBuffer(raw)) return raw.toString("utf8");
if (isRecord(raw) && "data" in raw) return extractMessageData(raw.data);
return String(raw ?? "");
}
@@ -0,0 +1,650 @@
import { randomUUID } from "node:crypto";
import {
AgentAdapter,
AgentEvent,
approvalDecisionKind,
approvalDecisionKindForMethod,
asJsonObject,
asJsonValue,
Disposable,
hasApprovalDecisionField,
isRecord,
JsonObject,
JsonRpcId,
Logger,
JsonValue,
RelayActor,
RelayCommandFrame,
RelayEventFrame,
RelayFrame,
RelayTransport,
} from "./protocol";
export interface RelayHostOptions {
adapter: AgentAdapter;
relay: RelayTransport;
sessionId?: string;
actor?: RelayActor;
/** Capabilities enforced locally even when relay authorization is bypassed. */
capabilities?: Iterable<string>;
logger?: Logger;
/** Emit a handshake on transports that do not implement one themselves. */
sendHandshake?: boolean;
}
/**
* Maps relay commands to the app-server AgentAdapter and publishes normalized
* adapter events. This is the policy boundary for the VS Code host.
*/
export class RelayHost {
private readonly adapter: AgentAdapter;
private readonly relay: RelayTransport;
private readonly options: RelayHostOptions;
private readonly subscriptions: Disposable[] = [];
private readonly commandResults = new Map<string, RelayEventFrame>();
private readonly inFlightCommands = new Set<string>();
private readonly capabilities: Set<string>;
private eventSeq = 0;
private sessionId: string;
private started = false;
private adapterReady = false;
constructor(options: RelayHostOptions);
constructor(adapter: AgentAdapter, relay: RelayTransport, options?: Omit<RelayHostOptions, "adapter" | "relay">);
constructor(
optionsOrAdapter: RelayHostOptions | AgentAdapter,
relayArg?: RelayTransport,
legacyOptions: Omit<RelayHostOptions, "adapter" | "relay"> = {},
) {
if (isAgentAdapter(optionsOrAdapter)) {
this.adapter = optionsOrAdapter;
if (!relayArg) throw new Error("RelayHost requires a relay transport");
this.relay = relayArg;
this.options = { ...legacyOptions, adapter: this.adapter, relay: this.relay };
} else {
this.options = optionsOrAdapter;
this.adapter = optionsOrAdapter.adapter;
this.relay = optionsOrAdapter.relay;
}
this.capabilities = new Set(this.options.capabilities ?? [
"read_output",
"send_task_input",
"cancel_task",
"approve_low_risk",
]);
this.sessionId = this.options.sessionId ?? `sess_${randomUUID()}`;
}
get id(): string {
return this.sessionId;
}
async start(): Promise<void> {
if (this.started) return;
this.started = true;
this.subscriptions.push(this.adapter.onEvent((event) => {
if (event.type === "connection.opened") this.adapterReady = true;
if (event.type === "connection.closed") this.adapterReady = false;
this.publishAgentEvent(event);
}));
this.subscriptions.push(this.relay.onMessage((frame) => {
void this.handleFrame(frame).catch((error) => {
this.options.logger?.warn?.("Invalid relay frame", error);
if (isRecord(frame) && typeof frame.commandId === "string") {
this.sendCommandResult(frame.commandId, false, undefined, error instanceof Error ? error.message : String(error), typeof frame.method === "string" ? frame.method : typeof frame.type === "string" ? frame.type : undefined);
}
});
}));
if (this.relay.onClose) this.subscriptions.push(this.relay.onClose((error) => {
// A relay disconnect must not leave an app-server request waiting for a
// browser that can no longer answer. The adapter's local deny path is
// deliberately fail-closed. Do not publish `connection.closed` here:
// that event describes the app-server process, while this callback only
// describes the outbound transport and is followed by connection.opened
// on a successful reconnect.
void this.adapter.denyPending?.("relay disconnected");
this.options.logger?.debug?.("Relay transport closed", error?.message ?? "");
}));
if (this.relay.onOpen) this.subscriptions.push(this.relay.onOpen(() => {
// RelayClient fires onOpen only after auth.ok. On reconnect the adapter
// is already initialized, so the synthetic event restores relay state;
// during initial startup the adapter event below is authoritative.
if (this.adapterReady) {
this.publishConnectionEvent("connection.opened");
void this.publishSnapshot();
}
}));
const configurableRelay = this.relay as RelayTransport & { setSessionId?: (sessionId: string) => void };
configurableRelay.setSessionId?.(this.sessionId);
try {
await this.relay.connect();
if (this.options.sendHandshake !== false && !transportHandlesHandshake(this.relay)) {
this.safeSend({ v: 1, kind: "hello", clientType: "host", protocol: 1, sessionId: this.sessionId });
}
// Start app-server only after the relay handshake is queued/sent. This
// keeps standalone stdout frames protocol-ordered and prevents an early
// notification from racing the host hello.
await this.adapter.start();
if (!this.adapterReady) {
this.adapterReady = true;
this.publishConnectionEvent("connection.opened");
}
await this.publishSnapshot();
} catch (error) {
this.started = false;
this.adapterReady = false;
for (const subscription of this.subscriptions.splice(0)) subscription.dispose();
this.relay.close();
await this.adapter.dispose().catch(() => undefined);
throw error;
}
}
async stop(): Promise<void> {
if (!this.started) return;
this.started = false;
this.adapterReady = false;
this.inFlightCommands.clear();
for (const subscription of this.subscriptions.splice(0)) subscription.dispose();
this.relay.close();
await this.adapter.dispose();
}
/** Public for unit tests and local stdin bridges. */
async handleFrame(frame: RelayFrame): Promise<void> {
if (!isRecord(frame)) return;
if (frame.kind === "command" || isCommandLike(frame)) {
await this.handleCommand(frame as unknown as RelayCommandFrame);
return;
}
if (frame.kind === "event" && typeof frame.seq === "number") {
this.safeSend({ v: 1, kind: "ack", sessionId: frame.sessionId, seq: frame.seq });
}
}
private async handleCommand(frame: RelayCommandFrame): Promise<void> {
const command = normalizeCommand(frame);
const commandId = command.commandId;
if (commandId) {
const previous = this.commandResults.get(commandId);
if (previous) {
this.safeSend(previous);
return;
}
if (this.inFlightCommands.has(commandId)) {
this.safeSend({
v: 1,
kind: "event",
// Do not call this `command.accepted`: the relay treats that event
// as the terminal result for its pending command. A retry while the
// original operation is running is only an informational event.
type: "command.pending",
id: `evt_${randomUUID()}`,
sessionId: this.sessionId,
seq: ++this.eventSeq,
ts: new Date().toISOString(),
actor: this.options.actor ?? { id: "host", role: "host" },
payload: { commandId, duplicate: true, pending: true },
});
return;
}
this.inFlightCommands.add(commandId);
}
const role = frame.actor?.role ?? "operator";
const denied = authorize(command.type, role, this.capabilities);
if (denied) {
// A viewer may not force a pending approval to deny (that would turn a
// read-only role into a denial-of-service primitive). Authorized roles
// can still be rejected by local capability/policy checks, in which
// case denying the app-server request is the safe terminal action.
if (role === "owner" || role === "operator" || role === "approver" || role === "host") {
await this.denyApprovalIfNeeded(command.type, command.payload, denied);
}
this.sendCommandResult(commandId, false, undefined, denied, command.type);
if (commandId) this.inFlightCommands.delete(commandId);
return;
}
try {
const result = await this.executeCommand(command.type, command.payload);
this.sendCommandResult(commandId, true, result, undefined, command.type);
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
this.options.logger?.warn?.(`Relay command ${command.type} failed`, error);
await this.denyApprovalIfNeeded(command.type, command.payload, message);
this.sendCommandResult(commandId, false, undefined, message, command.type);
} finally {
if (commandId) this.inFlightCommands.delete(commandId);
}
}
private async denyApprovalIfNeeded(type: string, payload: JsonObject, reason: string): Promise<void> {
const command = canonicalCommandType(type);
if (command !== "approval.respond" && command !== "input.respond" && command !== "server.request.respond") return;
const requestId = payload.requestId;
if (requestId === undefined || (typeof requestId !== "string" && typeof requestId !== "number")) return;
try {
await this.adapter.respondApproval(requestId, "deny", reason);
} catch {
// The request may already have expired or been resolved. Keep the
// original command rejection as the observable result.
}
}
private async executeCommand(type: string, payload: JsonObject): Promise<unknown> {
switch (canonicalCommandType(type)) {
case "control.mode.get": {
const mode = this.adapter.getControlMode?.();
if (mode) return { mode };
const snapshot = await this.adapter.snapshot();
const controlMode = snapshot.metadata?.controlMode;
if (controlMode !== "sync" && controlMode !== "async") {
throw new Error("adapter does not expose a control mode");
}
return { mode: controlMode };
}
case "control.mode.set":
if (!this.adapter.setControlMode) throw new Error("adapter does not support control mode switching");
return this.adapter.setControlMode(payload);
case "thread.start":
if (!this.adapter.startThread) throw new Error("adapter does not support thread/start");
return this.adapter.startThread(payload);
case "session.new":
if (!this.adapter.newSession) throw new Error("adapter does not support session/new");
return this.adapter.newSession(payload);
case "thread.settings.update":
if (!this.adapter.updateThreadSettings) throw new Error("adapter does not support thread/settings/update");
return this.adapter.updateThreadSettings(payload);
case "session.list":
if (!this.adapter.listSessions) throw new Error("adapter does not support session/list");
return this.adapter.listSessions(payload);
case "session.select":
if (!this.adapter.selectSession) throw new Error("adapter does not support session/select");
return this.adapter.selectSession(payload);
case "turn.start":
if (this.adapter.startTurn) return this.adapter.startTurn(payload);
return this.adapter.sendInput(extractCommandText(payload), payload);
case "turn.steer":
if (this.adapter.steerTurn) return this.adapter.steerTurn(payload);
return this.adapter.sendInput(extractCommandText(payload), payload);
case "turn.interrupt":
if (this.adapter.interruptTurn) return this.adapter.interruptTurn(payload);
return this.adapter.cancel(typeof payload.turnId === "string" ? payload.turnId : undefined, payload);
case "task.input": {
const text = typeof payload.text === "string" ? payload.text : typeof payload.message === "string" ? payload.message : undefined;
if (!text) throw new Error("task.input requires payload.text");
return this.adapter.sendInput(text, payload);
}
case "task.cancel": {
const taskId = typeof payload.taskId === "string" ? payload.taskId : typeof payload.turnId === "string" ? payload.turnId : undefined;
return this.adapter.cancel(taskId, payload);
}
case "approval.respond": {
const requestId = payload.requestId;
if (typeof requestId !== "string" && typeof requestId !== "number") throw new Error("approval.respond requires requestId");
const requestedValue = payload.decision;
const decision = approvalDecisionKind(requestedValue);
if (!decision) throw new Error("decision must be a recognized allow, deny, or cancel value");
const snapshot = await this.adapter.snapshot();
// JSON-RPC distinguishes numeric and string ids. Keep the lookup
// type-safe so id `1` cannot accidentally authorize response `"1"`.
const approval = snapshot.pendingApprovals.find((item) => item.requestId === requestId);
const method = approval?.method ?? (typeof payload.method === "string" ? payload.method : undefined);
const response = payload.response ?? implicitApprovalResponse(requestedValue, decision, method, payload.scope);
if (decision === "allow") {
if (approval && typeof payload.commandHash === "string" && payload.commandHash !== approval.commandHash) {
throw new Error("approval commandHash does not match the pending request");
}
if (approval?.risk === "high" && !this.capabilities.has("approve_high_risk") && !this.capabilities.has("*")) {
throw new Error("host policy requires approve_high_risk for this approval");
}
}
validateApprovalResponse(decision, response, method);
return this.adapter.respondApproval(
requestId,
decision,
typeof payload.reason === "string" ? payload.reason : undefined,
response,
);
}
case "input.respond":
case "server.request.respond": {
const requestId = payload.requestId;
if (typeof requestId !== "string" && typeof requestId !== "number") throw new Error(`${type} requires requestId`);
const response = payload.response ?? (payload.answers !== undefined ? payload.answers : undefined);
// Tool-input requests do not carry an allow/deny field in their wire
// response, while MCP elicitation uses `action`. Prefer an explicit
// response action when present; otherwise honor the relay decision and
// fail closed when a denial has no custom response.
if (response !== undefined && !isRecord(response)) {
throw new Error("input response must be a JSON object");
}
const responseDecision = isRecord(response)
? explicitResponseDecision(response)
: undefined;
const requestedDecision = payload.decision === undefined
? undefined
: approvalDecisionKind(payload.decision);
if (payload.decision !== undefined && !requestedDecision) {
throw new Error("decision must be allow, deny, or cancel");
}
if (responseDecision && requestedDecision
&& responseDecision !== requestedDecision
// RelayHost uses `decision: "allow"` as a generic envelope for
// MCP/input responses; the nested action remains authoritative in
// that one compatibility case.
&& requestedDecision !== "allow") {
throw new Error(`input response implies ${responseDecision}, but decision is ${requestedDecision}`);
}
// For MCP, `response.action` is the actual app-server decision and is
// authoritative even if a relay uses `decision: "allow"` as a generic
// input-response envelope. With no custom response, an explicit relay
// decision (or the fail-closed deny default) controls the result.
const decision = responseDecision ?? requestedDecision ?? (response === undefined ? "deny" : "allow");
const responseForAdapter = requestedDecision && requestedDecision !== "allow" && !responseDecision
? undefined
: response;
return this.adapter.respondApproval(
requestId,
decision,
typeof payload.reason === "string" ? payload.reason : undefined,
responseForAdapter,
);
}
case "session.snapshot":
case "snapshot":
return this.adapter.snapshot();
case "ping":
return { pong: true, ts: new Date().toISOString() };
default:
throw new Error(`unsupported relay command: ${type}`);
}
}
private sendCommandResult(commandId: string | undefined, ok: boolean, result?: unknown, error?: string, method?: string): void {
const frame: RelayEventFrame = {
v: 1,
kind: "event",
type: ok ? "command.accepted" : "command.rejected",
id: `evt_${randomUUID()}`,
sessionId: this.sessionId,
seq: ++this.eventSeq,
ts: new Date().toISOString(),
actor: this.options.actor ?? { id: "host", role: "host" },
payload: {
...(commandId ? { commandId } : {}),
...(method ? { method } : {}),
ok,
...(ok ? { result: asJsonValue(result) } : { error: error ?? "command rejected" }),
},
};
if (commandId) {
this.commandResults.set(commandId, frame);
if (this.commandResults.size > 1000) this.commandResults.delete(this.commandResults.keys().next().value as string);
}
this.safeSend(frame);
}
private publishAgentEvent(event: AgentEvent): void {
if (event.threadId && this.sessionId.startsWith("sess_")) {
// Keep a stable relay session id while exposing the app-server thread id
// in the payload; a relay session may contain more than one thread.
}
const payload: JsonObject = {
...event.payload,
...(event.status ? { executionStatus: asJsonValue(event.status) } : {}),
...(event.threadId ? { threadId: event.threadId } : {}),
...(event.turnId ? { turnId: event.turnId } : {}),
...(event.requestId !== undefined ? { requestId: asJsonValue(event.requestId) } : {}),
...(event.raw !== undefined ? { raw: event.raw } : {}),
};
const frame: RelayEventFrame = {
v: 1,
kind: "event",
type: event.type,
id: `evt_${randomUUID()}`,
sessionId: this.sessionId,
seq: ++this.eventSeq,
ts: new Date().toISOString(),
actor: this.options.actor ?? { id: "host", role: "host" },
payload,
...(event.status ? { status: { ...event.status, activeFlags: [...event.status.activeFlags] } } : {}),
};
try {
this.safeSend(frame);
} catch (error) {
this.options.logger?.warn?.("Unable to publish relay event", error);
}
}
private safeSend(frame: RelayFrame): void {
try {
this.relay.send(frame);
} catch (error) {
this.options.logger?.warn?.("Unable to send relay frame", error);
}
}
private publishConnectionEvent(type: string, error?: Error): void {
if (!this.started && type === "connection.closed") return;
this.publishAgentEvent({ type, payload: error ? { message: error.message } : {} });
}
private async publishSnapshot(): Promise<void> {
try {
const snapshot = await this.adapter.snapshot();
this.publishAgentEvent({
type: "session.snapshot",
threadId: snapshot.threadId ?? undefined,
turnId: snapshot.turnId ?? undefined,
payload: asJsonObject(snapshot),
status: snapshot.status,
});
} catch (error) {
this.options.logger?.warn?.("Unable to publish adapter snapshot", error);
}
}
}
interface NormalizedCommand {
type: string;
commandId?: string;
payload: JsonObject;
}
function normalizeCommand(frame: RelayCommandFrame): NormalizedCommand {
const nested = isRecord(frame.command) ? frame.command : undefined;
const type = typeof nested?.type === "string"
? nested.type
: typeof frame.method === "string"
? frame.method
: frame.type === "command"
? ""
: frame.type;
const commandId = typeof frame.commandId === "string"
? frame.commandId
: typeof nested?.commandId === "string"
? nested.commandId
: typeof frame.id === "string"
? frame.id
: undefined;
if (!type) throw new Error("relay command has no type");
if (isRecord(nested?.payload)) return { type, commandId, payload: asJsonObject(nested.payload) };
if (isRecord(frame.payload)) return { type, commandId, payload: asJsonObject(frame.payload) };
if (isRecord(frame.params)) return { type, commandId, payload: asJsonObject(frame.params) };
const payload: JsonObject = {};
for (const [key, value] of Object.entries(frame)) {
if (["v", "kind", "type", "method", "params", "commandId", "id", "sessionId", "actor", "command"].includes(key)) continue;
if (value !== undefined) payload[key] = asJsonValue(value);
}
return { type, commandId, payload };
}
function canonicalCommandType(type: string): string {
const normalized = type.trim().replace(/\//g, ".").replace(/\s+/g, ".").toLowerCase();
if (normalized === "control.mode.get" || normalized === "controlmode.get" || normalized === "controlmodeget" || normalized === "mode.get" || normalized === "modeget") return "control.mode.get";
if (normalized === "control.mode.set" || normalized === "controlmode.set" || normalized === "controlmodeset" || normalized === "mode.set" || normalized === "modeset") return "control.mode.set";
if (normalized === "thread.start" || normalized === "threadstart") return "thread.start";
if (normalized === "session.new" || normalized === "sessionnew" || normalized === "thread.new" || normalized === "threadnew") return "session.new";
if (normalized === "thread.settings.update" || normalized === "threadsettings.update" || normalized === "threadsettingsupdate") return "thread.settings.update";
if (normalized === "session.list" || normalized === "thread.list" || normalized === "sessionlist" || normalized === "threadlist") return "session.list";
if (normalized === "session.select" || normalized === "session.switch" || normalized === "thread.select" || normalized === "thread.attach" || normalized === "sessionswitch" || normalized === "threadselect") return "session.select";
if (normalized === "turn.start" || normalized === "turnstart") return "turn.start";
if (normalized === "turn.steer" || normalized === "turnsteer") return "turn.steer";
if (normalized === "turn.interrupt" || normalized === "turninterrupt") return "turn.interrupt";
if (normalized === "approval.respond" || normalized === "approvalrespond") return "approval.respond";
if (normalized === "task.input" || normalized === "taskinput") return "task.input";
if (normalized === "task.cancel" || normalized === "taskcancel") return "task.cancel";
if (normalized === "input.respond" || normalized === "inputrespond") return "input.respond";
if (normalized === "server.request.respond" || normalized === "serverrequest.respond") return "server.request.respond";
if (normalized === "session.snapshot" || normalized === "snapshot") return "session.snapshot";
return normalized;
}
function authorize(type: string, role: string, capabilities: Set<string>): string | undefined {
const command = canonicalCommandType(type);
const readOnly = command === "session.snapshot" || command === "snapshot" || command === "session.list" || command === "control.mode.get" || command === "ping";
if (readOnly) return undefined;
if (role === "viewer") return "viewer role cannot issue control commands";
if (role !== "owner" && role !== "operator" && role !== "approver" && role !== "host") return `role ${role} is not authorized`;
if ((command === "approval.respond" || command === "input.respond" || command === "server.request.respond") && role !== "owner" && role !== "operator" && role !== "approver" && role !== "host") {
return "role is not authorized to resolve approvals";
}
const required = command === "approval.respond" ? "approve_low_risk" : command === "task.cancel" || command === "turn.interrupt" ? "cancel_task" : command === "task.input" || command.startsWith("turn.") || command === "thread.start" || command === "thread.settings.update" || command === "session.select" || command === "session.new" || command === "control.mode.set" ? "send_task_input" : undefined;
if (required && !capabilities.has(required) && !capabilities.has("*") && role !== "owner" && role !== "host") return `missing capability: ${required}`;
return undefined;
}
function isCommandLike(frame: Record<string, unknown>): boolean {
if (frame.kind === "command") return true;
if (frame.type === "command" && typeof frame.method === "string") return true;
if (frame.kind !== undefined) return false;
if (typeof frame.method === "string") return true;
if (typeof frame.commandId !== "string") return false;
return KNOWN_COMMAND_TYPES.has(String(frame.type).trim().replace(/\//g, ".").toLowerCase());
}
const KNOWN_COMMAND_TYPES = new Set([
"control.mode.get",
"control.mode.set",
"thread.start",
"session.new",
"thread.settings.update",
"session.list",
"session.select",
"turn.start",
"turn.steer",
"turn.interrupt",
"approval.respond",
"task.input",
"task.cancel",
"input.respond",
"server.request.respond",
"session.snapshot",
"snapshot",
"ping",
]);
function isAgentAdapter(value: unknown): value is AgentAdapter {
return isRecord(value) && typeof value.start === "function" && typeof value.onEvent === "function" && typeof value.sendInput === "function" && typeof value.cancel === "function" && typeof value.respondApproval === "function" && typeof value.snapshot === "function";
}
function extractCommandText(payload: JsonObject): string {
if (typeof payload.text === "string") return payload.text;
if (typeof payload.message === "string") return payload.message;
if (typeof payload.prompt === "string") return payload.prompt;
if (Array.isArray(payload.input)) {
const first = payload.input[0];
if (isRecord(first) && typeof first.text === "string") return first.text;
}
throw new Error("turn command requires text or input");
}
function validateApprovalResponse(
decision: "allow" | "deny" | "cancel",
response: JsonValue | undefined,
method?: string,
): void {
if (response === undefined) return;
if (!isRecord(response)) throw new Error("approval response must be a JSON object");
const hasDecision = Object.prototype.hasOwnProperty.call(response, "decision");
const hasAction = Object.prototype.hasOwnProperty.call(response, "action");
if (hasDecision || hasAction) {
const decisionKind = hasDecision ? approvalDecisionKindForMethod(response.decision, method) : undefined;
const actionKind = hasAction ? approvalDecisionKindForMethod(response.action, method) : undefined;
if (hasDecision && !decisionKind) throw new Error("unsupported approval response decision");
if (hasAction && !actionKind) throw new Error("unsupported approval response action");
if (decisionKind && actionKind && decisionKind !== actionKind) {
throw new Error("approval response decision and action conflict");
}
const implied = decisionKind ?? actionKind;
if (implied && implied !== decision) {
throw new Error(`approval response implies ${implied}, but decision is ${decision}`);
}
return;
}
// Permissions approvals intentionally carry a profile rather than a
// decision field. Keep the profile shape narrow; malformed/unknown objects
// must not be interpreted as an approval.
if (method === "item/permissions/requestApproval"
&& isRecord(response.permissions)
&& (response.scope === "turn" || response.scope === "session")
&& (response.strictAutoReview === undefined || typeof response.strictAutoReview === "boolean")
&& Object.keys(response).every((key) => key === "permissions" || key === "scope" || key === "strictAutoReview")) {
return;
}
throw new Error("approval response has no recognized decision or permission profile");
}
/**
* Convert a relay's compact outer decision into a wire response only when it
* carries a non-canonical app-server value. Canonical `allow`/`deny`/`cancel`
* remain undefined so the adapter can choose the method-specific default.
*/
function implicitApprovalResponse(
requestedValue: JsonValue,
decision: "allow" | "deny" | "cancel",
method?: string,
scope?: JsonValue,
): JsonValue | undefined {
// Permission approvals have a profile response, not a decision wrapper.
// Let the adapter construct the requested turn-scoped profile by default;
// callers that need session scope must provide the full profile explicitly.
if (method === "item/permissions/requestApproval") return undefined;
if (requestedValue === "allow" || requestedValue === "deny" || requestedValue === "cancel") {
if (requestedValue === "allow" && scope === "session") {
if (method === "applyPatchApproval" || method === "execCommandApproval") return { decision: "approved_for_session" };
if (method === "item/commandExecution/requestApproval" || method === "item/fileChange/requestApproval") return { decision: "acceptForSession" };
}
return undefined;
}
// The generic classifier has already rejected unknown/conflicting values.
// Preserve recognized legacy/v2 tags exactly under the app-server wrapper.
if (decision === "allow" || decision === "deny" || decision === "cancel") {
return { decision: requestedValue };
}
return undefined;
}
function explicitResponseDecision(response: Record<string, unknown>): "allow" | "deny" | "cancel" | undefined {
if (!hasApprovalDecisionField(response)) return undefined;
const hasDecision = Object.prototype.hasOwnProperty.call(response, "decision");
const hasAction = Object.prototype.hasOwnProperty.call(response, "action");
const decision = hasDecision ? approvalDecisionKind(response.decision) : undefined;
const action = hasAction ? approvalDecisionKind(response.action) : undefined;
if (hasDecision && !decision) throw new Error("unsupported input response decision");
if (hasAction && !action) throw new Error("unsupported input response action");
if (decision && action && decision !== action) throw new Error("input response decision and action conflict");
return decision ?? action;
}
function transportHandlesHandshake(transport: RelayTransport): boolean {
return Boolean((transport as RelayTransport & { handlesHandshake?: boolean }).handlesHandshake);
}
@@ -0,0 +1,425 @@
import {
AgentAdapter,
AgentEvent,
asJsonObject,
ControlMode,
Disposable,
JsonObject,
JsonRpcId,
JsonValue,
Logger,
SessionSnapshot,
} from "./protocol";
export type AgentAdapterFactory = (mode: ControlMode) => AgentAdapter | Promise<AgentAdapter>;
export interface SwitchableAgentAdapterOptions {
initialMode: ControlMode;
createAdapter: AgentAdapterFactory;
/** Persist the committed mode. Persistence errors do not roll back a live adapter. */
onModeChanged?: (mode: ControlMode, previousMode: ControlMode) => void | Promise<void>;
logger?: Logger;
}
interface AdapterBinding {
adapter: AgentAdapter;
mode: ControlMode;
generation: number;
committed: boolean;
bufferedEvents: AgentEvent[];
subscription: Disposable;
}
interface ModeCapabilities extends JsonObject {
followsVscodeRoute: boolean;
sessionList: boolean;
sessionSelect: boolean;
sessionCreate: boolean;
threadSettings: boolean;
}
/**
* Keeps RelayHost bound to one stable AgentAdapter while atomically replacing
* the implementation behind it when the control owner changes.
*/
export class SwitchableAgentAdapter implements AgentAdapter {
private readonly options: SwitchableAgentAdapterOptions;
private readonly listeners = new Set<(event: AgentEvent) => void>();
private binding: AdapterBinding | null = null;
private controlMode: ControlMode;
private modeEpoch = 0;
private startPromise: Promise<void> | null = null;
private switchPromise: Promise<JsonValue> | null = null;
private started = false;
private disposed = false;
constructor(options: SwitchableAgentAdapterOptions) {
this.options = options;
this.controlMode = validateControlMode(options.initialMode);
}
async start(): Promise<void> {
if (this.started) return;
if (this.disposed) throw new Error("switchable adapter has been disposed");
if (this.startPromise) return this.startPromise;
const operation = this.startInitialAdapter();
this.startPromise = operation;
try {
await operation;
} finally {
if (this.startPromise === operation) this.startPromise = null;
}
}
getControlMode(): ControlMode {
return this.controlMode;
}
async setControlMode(params: JsonObject): Promise<JsonValue> {
const nextMode = controlModeFromParams(params);
this.ensureStarted();
if (this.switchPromise) throw new Error("a control mode switch is already in progress");
if (nextMode === this.controlMode) {
return {
changed: false,
controlMode: this.controlMode,
previousControlMode: this.controlMode,
modeEpoch: this.modeEpoch,
};
}
const operation = this.performModeSwitch(nextMode);
this.switchPromise = operation;
try {
return await operation;
} finally {
if (this.switchPromise === operation) this.switchPromise = null;
}
}
async startThread(params: JsonObject = {}): Promise<JsonValue> {
this.assertIndependentNavigation("thread/start");
const adapter = this.activeAdapterForMutation();
if (!adapter.startThread) throw unsupported("thread/start", this.controlMode);
return adapter.startThread(params);
}
async newSession(params: JsonObject = {}): Promise<JsonValue> {
this.assertIndependentNavigation("session/new");
const adapter = this.activeAdapterForMutation();
if (adapter.newSession) return adapter.newSession(params);
if (adapter.startThread) return adapter.startThread(params);
throw unsupported("session/new", this.controlMode);
}
async startTurn(params: JsonObject): Promise<JsonValue> {
const adapter = this.activeAdapterForMutation();
if (!adapter.startTurn) throw unsupported("turn/start", this.controlMode);
return adapter.startTurn(params);
}
async steerTurn(params: JsonObject): Promise<JsonValue> {
const adapter = this.activeAdapterForMutation();
if (!adapter.steerTurn) throw unsupported("turn/steer", this.controlMode);
return adapter.steerTurn(params);
}
async updateThreadSettings(params: JsonObject): Promise<JsonValue> {
const adapter = this.activeAdapterForMutation();
if (!adapter.updateThreadSettings) throw unsupported("thread/settings/update", this.controlMode);
return adapter.updateThreadSettings(params);
}
async listSessions(params: JsonObject = {}): Promise<JsonValue> {
this.assertIndependentNavigation("session/list");
const adapter = this.activeAdapter();
if (!adapter.listSessions) throw unsupported("session/list", this.controlMode);
return adapter.listSessions(params);
}
async selectSession(params: JsonObject): Promise<JsonValue> {
this.assertIndependentNavigation("session/select");
const adapter = this.activeAdapterForMutation();
if (!adapter.selectSession) throw unsupported("session/select", this.controlMode);
return adapter.selectSession(params);
}
async interruptTurn(params: JsonObject): Promise<JsonValue> {
const adapter = this.activeAdapter();
if (!adapter.interruptTurn) throw unsupported("turn/interrupt", this.controlMode);
return adapter.interruptTurn(params);
}
async sendInput(text: string, params: JsonObject = {}): Promise<JsonValue> {
return this.activeAdapterForMutation().sendInput(text, params);
}
async cancel(taskId?: string, params: JsonObject = {}): Promise<JsonValue> {
return this.activeAdapter().cancel(taskId, params);
}
async respondApproval(
requestId: JsonRpcId,
decision: "allow" | "deny" | "cancel",
reason?: string,
response?: JsonValue,
): Promise<JsonValue> {
return this.activeAdapter().respondApproval(requestId, decision, reason, response);
}
async denyPending(reason?: string): Promise<void> {
await this.activeAdapter().denyPending?.(reason);
}
async snapshot(): Promise<SessionSnapshot> {
const binding = this.activeBinding();
const snapshot = await binding.adapter.snapshot();
return this.decorateSnapshot(snapshot, binding);
}
onEvent(listener: (event: AgentEvent) => void): Disposable {
this.listeners.add(listener);
return { dispose: () => this.listeners.delete(listener) };
}
async dispose(): Promise<void> {
if (this.disposed) return;
this.disposed = true;
const starting = this.startPromise;
const switching = this.switchPromise;
await starting?.catch(() => undefined);
await switching?.catch(() => undefined);
const binding = this.binding;
this.binding = null;
this.started = false;
if (!binding) return;
binding.committed = false;
binding.subscription.dispose();
await binding.adapter.dispose();
}
private async startInitialAdapter(): Promise<void> {
const binding = await this.createBinding(this.controlMode, this.modeEpoch);
try {
await binding.adapter.start();
if (this.disposed) throw new Error("switchable adapter was disposed while starting");
binding.committed = true;
this.binding = binding;
this.started = true;
this.flushBufferedEvents(binding);
} catch (error) {
await this.releaseBinding(binding);
throw error;
}
}
private async performModeSwitch(nextMode: ControlMode): Promise<JsonValue> {
const previousBinding = this.activeBinding();
const previousMode = this.controlMode;
this.assertSnapshotIdle(await previousBinding.adapter.snapshot());
const nextEpoch = this.modeEpoch + 1;
const candidate = await this.createBinding(nextMode, nextEpoch);
if (candidate.adapter === previousBinding.adapter) {
candidate.subscription.dispose();
throw new Error("adapter factory must return a distinct adapter when switching control modes");
}
try {
await candidate.adapter.start();
if (this.disposed) throw new Error("switchable adapter was disposed while switching modes");
// VS Code can start a turn independently while the candidate boots.
// Recheck immediately before the synchronous commit point.
this.assertSnapshotIdle(await previousBinding.adapter.snapshot());
const candidateSnapshot = await candidate.adapter.snapshot();
// The candidate snapshot is an await point, so make the old adapter's
// liveness check the final operation before committing synchronously.
this.assertSnapshotIdle(await previousBinding.adapter.snapshot());
candidate.committed = true;
this.binding = candidate;
this.controlMode = nextMode;
this.modeEpoch = nextEpoch;
previousBinding.committed = false;
const result: JsonObject = {
changed: true,
controlMode: nextMode,
previousControlMode: previousMode,
modeEpoch: nextEpoch,
};
this.emit({ type: "control.mode.changed", payload: result });
this.flushBufferedEvents(candidate);
const snapshot = this.decorateSnapshot(candidateSnapshot, candidate);
this.emit({
type: "session.snapshot",
threadId: snapshot.threadId ?? undefined,
turnId: snapshot.turnId ?? undefined,
payload: asJsonObject(snapshot),
status: snapshot.status,
});
previousBinding.subscription.dispose();
await previousBinding.adapter.dispose().catch((error) => {
this.options.logger?.warn?.("Unable to dispose the previous control mode adapter", error);
});
await Promise.resolve(this.options.onModeChanged?.(nextMode, previousMode)).catch((error) => {
this.options.logger?.warn?.("Unable to persist the committed control mode", error);
});
return result;
} catch (error) {
if (this.binding !== candidate) await this.releaseBinding(candidate);
throw error;
}
}
private async createBinding(mode: ControlMode, generation: number): Promise<AdapterBinding> {
const adapter = await this.options.createAdapter(mode);
if (!adapter) throw new Error(`adapter factory returned no adapter for ${mode} mode`);
const binding: AdapterBinding = {
adapter,
mode,
generation,
committed: false,
bufferedEvents: [],
subscription: { dispose: () => undefined },
};
binding.subscription = adapter.onEvent((event) => this.receiveAdapterEvent(binding, event));
return binding;
}
private receiveAdapterEvent(binding: AdapterBinding, event: AgentEvent): void {
if (!binding.committed) {
binding.bufferedEvents.push(event);
return;
}
if (this.binding !== binding || binding.generation !== this.modeEpoch) return;
this.emit(this.decorateEvent(event, binding));
}
private flushBufferedEvents(binding: AdapterBinding): void {
const events = binding.bufferedEvents.splice(0);
for (const event of events) {
if (this.binding !== binding || binding.generation !== this.modeEpoch) return;
this.emit(this.decorateEvent(event, binding));
}
}
private emit(event: AgentEvent): void {
for (const listener of this.listeners) {
try {
listener(event);
} catch (error) {
this.options.logger?.warn?.("Switchable adapter event listener failed", error);
}
}
}
private activeBinding(): AdapterBinding {
this.ensureStarted();
if (!this.binding) throw new Error("switchable adapter has no active adapter");
return this.binding;
}
private activeAdapter(): AgentAdapter {
return this.activeBinding().adapter;
}
private activeAdapterForMutation(): AgentAdapter {
if (this.switchPromise) throw new Error("control mode is switching; retry after it completes");
return this.activeAdapter();
}
private ensureStarted(): void {
if (this.disposed) throw new Error("switchable adapter has been disposed");
if (!this.started || !this.binding) throw new Error("switchable adapter is not started");
}
private assertIndependentNavigation(operation: string): void {
if (this.controlMode === "sync") {
throw new Error(`${operation} is unavailable in sync mode; conversation navigation follows VS Code`);
}
}
private assertSnapshotIdle(snapshot: SessionSnapshot): void {
const pendingApprovalCount = snapshot.pendingApprovals.length;
const pendingRequestCount = snapshot.pendingRequests?.length ?? 0;
const state = normalizeStatus(snapshot.state);
const turnStatus = normalizeStatus(snapshot.status?.turnStatus ?? snapshot.turnStatus ?? "");
const activeFlags = snapshot.status?.activeFlags ?? snapshot.activeFlags ?? [];
const hasActiveState = ACTIVE_STATUSES.has(state) || ACTIVE_STATUSES.has(turnStatus) || activeFlags.length > 0;
if (snapshot.turnId || hasActiveState || pendingApprovalCount > 0 || pendingRequestCount > 0) {
throw new Error("cannot switch control mode while a turn or request is active");
}
}
private decorateSnapshot(snapshot: SessionSnapshot, binding: AdapterBinding): SessionSnapshot {
return {
...snapshot,
metadata: {
...(snapshot.metadata ?? {}),
mode: binding.mode,
controlMode: binding.mode,
modeEpoch: binding.generation,
capabilities: this.capabilities(binding),
},
};
}
private decorateEvent(event: AgentEvent, binding: AdapterBinding): AgentEvent {
if (event.type !== "session.snapshot") return event;
return {
...event,
payload: {
...event.payload,
metadata: {
...asJsonObject(event.payload.metadata),
mode: binding.mode,
controlMode: binding.mode,
modeEpoch: binding.generation,
capabilities: this.capabilities(binding),
},
},
};
}
private capabilities(binding: AdapterBinding): ModeCapabilities {
const independent = binding.mode === "async";
return {
followsVscodeRoute: !independent,
sessionList: independent && typeof binding.adapter.listSessions === "function",
sessionSelect: independent && typeof binding.adapter.selectSession === "function",
sessionCreate: independent && (typeof binding.adapter.newSession === "function"
|| typeof binding.adapter.startThread === "function"),
threadSettings: typeof binding.adapter.updateThreadSettings === "function",
};
}
private async releaseBinding(binding: AdapterBinding): Promise<void> {
binding.committed = false;
binding.subscription.dispose();
await binding.adapter.dispose().catch(() => undefined);
}
}
function controlModeFromParams(params: JsonObject): ControlMode {
return validateControlMode(params.mode ?? params.controlMode);
}
function validateControlMode(value: unknown): ControlMode {
if (value === "sync" || value === "async") return value;
throw new Error("control mode must be sync or async");
}
function unsupported(operation: string, mode: ControlMode): Error {
return new Error(`${operation} is not supported by the ${mode} adapter`);
}
const ACTIVE_STATUSES = new Set(["active", "inprogress", "running", "starting", "thinking", "editing", "working"]);
function normalizeStatus(value: string): string {
return value.trim().replace(/[\s_-]+/g, "").toLowerCase();
}