diff --git a/src/providers/codex.ts b/src/providers/codex.ts index 6e4a27ce..12962243 100644 --- a/src/providers/codex.ts +++ b/src/providers/codex.ts @@ -264,8 +264,12 @@ function shouldSkipCodexSandboxConfig( } function shouldDisableCodexWebsockets(): boolean { - const envRaw = Deno.env.get(CODEX_DISABLE_WEBSOCKETS_ENV); - return Boolean(envRaw && parseTruthy(envRaw)); + // Newer Codex releases reserve built-in provider IDs, so overriding + // `model_providers.openai.supports_websockets` now prevents app-server + // startup. Keep the env var harmless while the old transport workaround ages + // out. + Deno.env.get(CODEX_DISABLE_WEBSOCKETS_ENV); + return false; } function tomlString(value: string): string { diff --git a/src/providers/codex_app_server.test.ts b/src/providers/codex_app_server.test.ts index 46da3290..f300bef5 100644 --- a/src/providers/codex_app_server.test.ts +++ b/src/providers/codex_app_server.test.ts @@ -35,7 +35,7 @@ Deno.test("codex app-server refresh host failures are returned as RPC errors", a } }); -Deno.test("codex config can disable OpenAI websocket responses", () => { +Deno.test("codex websocket disable env does not override built-in OpenAI provider", () => { const previous = Deno.env.get("GAMBIT_CODEX_DISABLE_WEBSOCKETS"); Deno.env.set("GAMBIT_CODEX_DISABLE_WEBSOCKETS", "1"); @@ -44,7 +44,7 @@ Deno.test("codex config can disable OpenAI websocket responses", () => { assertEquals( args.includes("model_providers.openai.supports_websockets=false"), - true, + false, ); } finally { if (previous === undefined) { diff --git a/src/runtime_host_service.test.ts b/src/runtime_host_service.test.ts new file mode 100644 index 00000000..164207b0 --- /dev/null +++ b/src/runtime_host_service.test.ts @@ -0,0 +1,171 @@ +import { assertEquals } from "@std/assert"; +import { assertThrows } from "@std/assert/throws"; +import { join } from "@std/path"; +import { + callRuntimeHostServiceRaw, + CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD, + DRAFT_COWORKER_TASK_HOST_SERVICE_METHOD, + validateRuntimeHostServiceMethodAndParams, +} from "./runtime_host_service.ts"; + +async function readJsonLine(conn: Deno.Conn): Promise> { + const reader = conn.readable.getReader(); + const decoder = new TextDecoder(); + let text = ""; + try { + while (true) { + const { value, done } = await reader.read(); + if (done) break; + text += decoder.decode(value, { stream: true }); + const newlineIndex = text.indexOf("\n"); + if (newlineIndex >= 0) { + return JSON.parse(text.slice(0, newlineIndex)); + } + } + } finally { + reader.releaseLock(); + } + throw new Error("expected JSON line"); +} + +Deno.test("runtime host-service raw calls return non-Codex host results unchanged", async () => { + const root = await Deno.makeTempDir({ dir: "/tmp", prefix: "rhs-" }); + const socketPath = join(root, "host-services.sock"); + const listener = Deno.listen({ path: socketPath, transport: "unix" }); + const server = (async () => { + const conn = await listener.accept(); + try { + const request = await readJsonLine(conn); + const writer = conn.writable.getWriter(); + await writer.write( + new TextEncoder().encode( + `${ + JSON.stringify({ + error: null, + id: request.id, + result: { + payload: { taskId: "workspace-delegation-smoke" }, + status: 200, + }, + }) + }\n`, + ), + ); + writer.releaseLock(); + } finally { + conn.close(); + listener.close(); + } + })(); + + try { + const result = await callRuntimeHostServiceRaw({ + method: DRAFT_COWORKER_TASK_HOST_SERVICE_METHOD, + params: { + targetCoworker: "assistant-to-chief-of-staff", + title: "Workspace delegation smoke", + purpose: "Verify host-owned task drafting.", + request: "Report the current working directory.", + }, + socketPath, + token: "test-token", + }); + + assertEquals(result, { + payload: { taskId: "workspace-delegation-smoke" }, + status: 200, + }); + } finally { + await server; + await Deno.remove(root, { recursive: true }).catch(() => undefined); + } +}); + +Deno.test("runtime host-service raw calls support TCP endpoints", async () => { + const listener = Deno.listen({ hostname: "127.0.0.1", port: 0 }); + const address = listener.addr; + if (address.transport !== "tcp") { + throw new Error("expected TCP listener"); + } + const server = (async () => { + const conn = await listener.accept(); + try { + const request = await readJsonLine(conn); + const writer = conn.writable.getWriter(); + await writer.write( + new TextEncoder().encode( + `${ + JSON.stringify({ + error: null, + id: request.id, + result: { + payload: { taskId: "tcp-workspace-delegation-smoke" }, + status: 200, + }, + }) + }\n`, + ), + ); + writer.releaseLock(); + } finally { + conn.close(); + listener.close(); + } + })(); + + try { + const result = await callRuntimeHostServiceRaw({ + method: DRAFT_COWORKER_TASK_HOST_SERVICE_METHOD, + params: { + targetCoworker: "assistant-to-chief-of-staff", + title: "TCP workspace delegation smoke", + purpose: "Verify host-owned task drafting over TCP.", + request: "Report the current working directory.", + }, + socketPath: `tcp://127.0.0.1:${address.port}`, + token: "test-token", + }); + + assertEquals(result, { + payload: { taskId: "tcp-workspace-delegation-smoke" }, + status: 200, + }); + } finally { + await server; + } +}); + +Deno.test("runtime host-service validates create writeback preview params", () => { + assertEquals( + validateRuntimeHostServiceMethodAndParams({ + method: CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD, + params: { + changedPaths: ["notes/smoke.md"], + summary: "Preview the runtime note.", + workspaceRoot: + "/runtime/cache/chief-session-workspaces/session-123/merged/coworkers/agents/assistant-to-chief-of-staff", + }, + }), + { + method: CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD, + params: { + changedPaths: ["notes/smoke.md"], + summary: "Preview the runtime note.", + workspaceRoot: + "/runtime/cache/chief-session-workspaces/session-123/merged/coworkers/agents/assistant-to-chief-of-staff", + }, + }, + ); + + assertThrows( + () => + validateRuntimeHostServiceMethodAndParams({ + method: CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD, + params: { + summary: "Preview the runtime note.", + }, + }), + Error, + "workspaceRoot", + ); +}); diff --git a/src/runtime_host_service.ts b/src/runtime_host_service.ts index 8682151e..35467939 100644 --- a/src/runtime_host_service.ts +++ b/src/runtime_host_service.ts @@ -5,6 +5,12 @@ export const RUNTIME_HOST_SERVICE_TOKEN_ENV = export const CODEX_REFRESH_HOST_SERVICE_METHOD = "providerAuth.codex.refreshChatgptTokens"; +export const DRAFT_COWORKER_TASK_HOST_SERVICE_METHOD = + "workloop.tasks.draftCoworkerTask"; +export const QUEUE_COWORKER_TASK_HOST_SERVICE_METHOD = + "workloop.tasks.queueCoworkerTask"; +export const CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD = + "workloop.writebacks.createPreview"; export type RuntimeHostServiceFailureReason = | "host_auth_missing" @@ -25,6 +31,38 @@ export type CodexRefreshHostServiceResult = { type: "chatgptAuthTokens"; }; +export type DraftCoworkerTaskHostServiceParams = { + targetCoworker: string; + taskId?: string | null; + title: string; + purpose: string; + request: string; + acceptanceCriteria?: Array; +}; + +export type QueueCoworkerTaskHostServiceParams = { + targetCoworker: string; + taskId: string; +}; + +export type CreateWritebackPreviewHostServiceParams = { + summary: string; + workspaceRoot: string; + changedPaths?: Array; +}; + +export type RuntimeHostServiceMethod = + | typeof CODEX_REFRESH_HOST_SERVICE_METHOD + | typeof DRAFT_COWORKER_TASK_HOST_SERVICE_METHOD + | typeof QUEUE_COWORKER_TASK_HOST_SERVICE_METHOD + | typeof CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD; + +export type RuntimeHostServiceParams = + | CodexRefreshHostServiceParams + | DraftCoworkerTaskHostServiceParams + | QueueCoworkerTaskHostServiceParams + | CreateWritebackPreviewHostServiceParams; + export type RuntimeHostServiceRequest = { id: string; method: string; @@ -108,20 +146,126 @@ export function validateCodexRefreshHostServiceResult( }; } +function normalizeStringArray(value: unknown, label: string): Array { + if (value == null) return []; + if (!Array.isArray(value)) { + throw new Error(`${label} must be an array when provided.`); + } + return value.map((entry, index) => { + if (typeof entry !== "string" || entry.trim().length === 0) { + throw new Error(`${label}[${index}] must be a non-empty string.`); + } + return entry.trim(); + }); +} + +export function validateDraftCoworkerTaskHostServiceParams( + value: unknown, +): DraftCoworkerTaskHostServiceParams { + if (!isRecord(value)) { + throw new Error("host service params must be a JSON object."); + } + const taskId = value.taskId == null + ? null + : normalizeOptionalString(value.taskId); + if (value.taskId != null && taskId == null) { + throw new Error( + "workloop.tasks.draftCoworkerTask taskId must be a non-empty string when provided.", + ); + } + return { + targetCoworker: normalizeRequiredString( + value.targetCoworker, + "targetCoworker", + ), + taskId, + title: normalizeRequiredString(value.title, "title"), + purpose: normalizeRequiredString(value.purpose, "purpose"), + request: normalizeRequiredString(value.request, "request"), + acceptanceCriteria: normalizeStringArray( + value.acceptanceCriteria, + "acceptanceCriteria", + ), + }; +} + +export function validateQueueCoworkerTaskHostServiceParams( + value: unknown, +): QueueCoworkerTaskHostServiceParams { + if (!isRecord(value)) { + throw new Error("host service params must be a JSON object."); + } + return { + targetCoworker: normalizeRequiredString( + value.targetCoworker, + "targetCoworker", + ), + taskId: normalizeRequiredString(value.taskId, "taskId"), + }; +} + +export function validateCreateWritebackPreviewHostServiceParams( + value: unknown, +): CreateWritebackPreviewHostServiceParams { + if (!isRecord(value)) { + throw new Error("host service params must be a JSON object."); + } + return { + summary: normalizeRequiredString(value.summary, "summary"), + workspaceRoot: normalizeRequiredString( + value.workspaceRoot, + "workspaceRoot", + ), + changedPaths: normalizeStringArray(value.changedPaths, "changedPaths"), + }; +} + export function validateRuntimeHostServiceMethodAndParams(input: { method: string; params: unknown; -}): { - method: typeof CODEX_REFRESH_HOST_SERVICE_METHOD; - params: CodexRefreshHostServiceParams; -} { - if (input.method !== CODEX_REFRESH_HOST_SERVICE_METHOD) { - throw new Error(`unsupported runtime host service method: ${input.method}`); +}): + | { + method: typeof CODEX_REFRESH_HOST_SERVICE_METHOD; + params: CodexRefreshHostServiceParams; + } + | { + method: typeof DRAFT_COWORKER_TASK_HOST_SERVICE_METHOD; + params: DraftCoworkerTaskHostServiceParams; + } + | { + method: typeof QUEUE_COWORKER_TASK_HOST_SERVICE_METHOD; + params: QueueCoworkerTaskHostServiceParams; + } + | { + method: typeof CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD; + params: CreateWritebackPreviewHostServiceParams; + } { + switch (input.method) { + case CODEX_REFRESH_HOST_SERVICE_METHOD: + return { + method: CODEX_REFRESH_HOST_SERVICE_METHOD, + params: validateCodexRefreshHostServiceParams(input.params), + }; + case DRAFT_COWORKER_TASK_HOST_SERVICE_METHOD: + return { + method: DRAFT_COWORKER_TASK_HOST_SERVICE_METHOD, + params: validateDraftCoworkerTaskHostServiceParams(input.params), + }; + case QUEUE_COWORKER_TASK_HOST_SERVICE_METHOD: + return { + method: QUEUE_COWORKER_TASK_HOST_SERVICE_METHOD, + params: validateQueueCoworkerTaskHostServiceParams(input.params), + }; + case CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD: + return { + method: CREATE_WRITEBACK_PREVIEW_HOST_SERVICE_METHOD, + params: validateCreateWritebackPreviewHostServiceParams(input.params), + }; + default: + throw new Error( + `unsupported runtime host service method: ${input.method}`, + ); } - return { - method: CODEX_REFRESH_HOST_SERVICE_METHOD, - params: validateCodexRefreshHostServiceParams(input.params), - }; } async function readFirstLine( @@ -159,12 +303,28 @@ async function writeJsonLine( } } -export async function callRuntimeHostService(input: { - method: typeof CODEX_REFRESH_HOST_SERVICE_METHOD; - params: CodexRefreshHostServiceParams; +function connectRuntimeHostService(endpoint: string): Promise { + if (!endpoint.startsWith("tcp://")) { + return Deno.connect({ transport: "unix", path: endpoint }); + } + const url = new URL(endpoint); + const port = Number(url.port); + if (!url.hostname || !Number.isInteger(port) || port <= 0) { + throw new Error(`invalid Workloop host service TCP endpoint: ${endpoint}`); + } + return Deno.connect({ + hostname: url.hostname, + port, + transport: "tcp", + }); +} + +export async function callRuntimeHostServiceRaw(input: { + method: RuntimeHostServiceMethod; + params: RuntimeHostServiceParams; socketPath?: string | null; token?: string | null; -}): Promise { +}): Promise { const socketPath = input.socketPath?.trim() || Deno.env.get(RUNTIME_HOST_SERVICE_SOCKET_ENV)?.trim(); const token = input.token?.trim() || @@ -179,7 +339,7 @@ export async function callRuntimeHostService(input: { token, type: "request", }; - const conn = await Deno.connect({ transport: "unix", path: socketPath }); + const conn = await connectRuntimeHostService(socketPath); try { await writeJsonLine(conn.writable, request); const line = await readFirstLine(conn.readable); @@ -192,7 +352,7 @@ export async function callRuntimeHostService(input: { if (response.error) { throw new Error(`${response.error.code}: ${response.error.message}`); } - return validateCodexRefreshHostServiceResult(response.result); + return response.result; } finally { try { conn.close(); @@ -201,3 +361,14 @@ export async function callRuntimeHostService(input: { } } } + +export async function callRuntimeHostService(input: { + method: typeof CODEX_REFRESH_HOST_SERVICE_METHOD; + params: CodexRefreshHostServiceParams; + socketPath?: string | null; + token?: string | null; +}): Promise { + return validateCodexRefreshHostServiceResult( + await callRuntimeHostServiceRaw(input), + ); +}