import { z } from "zod" import { automationInvokedTriggerDefinition } from "@automate.ax/catalog/triggers/core-invocation" import { defineAction } from "../../automation/actions" import { serializedIntegrationAccountDefinition } from "../../automation/integrations" import { getCurrentHookScopePath } from "../../automation/runtime" import { correlate, filter, getCurrentContextPrerequisite, getCurrentSignalPrerequisites, group, outcome, scope, withoutSignalPrerequisites, withContextPrerequisite, withPrerequisites, } from "../../automation/signal-operators" import type { DurationString } from "../../automation/signal-protocol" import { FailedSignalError, transform } from "../../automation/signal-protocol" import { createSubscription } from "../../automation/subscription" const CODEX_AUTH_URL = "https://auth.openai.com/api/accounts/v1/user-auth-credential/whoami" const CODEX_API_BASE_URL = "https://chatgpt.com/backend-api/" const CODEX_TASK_URL_BASE = "https://chatgpt.com/codex/tasks/" const CODEX_POLL_ENTRYPOINT_PREFIX = "__$codex-poll:" const CODEX_REQUEST_TIMEOUT_MS = 50_000 const CODEX_USER_AGENT = "automate.ax-codex" const CODEX_SECRET_SCHEMA = z.object({ apiKey: z.string().min(1) }) const CODEX_CLOUD_TASK_ID_SCHEMA = z.string().min(1) const CODEX_CLOUD_TURN_ID_SCHEMA = z.string().min(1) const CODEX_CLOUD_TASK_URL_SCHEMA = z.url() const CODEX_CLOUD_ACCOUNT_SELECTION_SCHEMA = z.union([ z.string().min(1), z.object({ [serializedIntegrationAccountDefinition]: z.object({ binding: z.string().min(1), serviceId: z.literal("codex"), }), }), ]) const CODEX_ACCESS_TOKEN_METADATA_SCHEMA = z.object({ chatgpt_account_id: z.string().min(1), chatgpt_account_is_fedramp: z.boolean(), }) const CODEX_ERROR_RESPONSE_SCHEMA = z.looseObject({ detail: z.string().optional(), error: z.looseObject({ message: z.string().optional() }).optional(), message: z.string().optional(), }) const CODEX_TEXT_CONTENT_SCHEMA = z.looseObject({ content_type: z.string().optional(), text: z.string().optional(), }) const CODEX_TEXT_CONTENT_PART_SCHEMA = z.union([ z.string(), CODEX_TEXT_CONTENT_SCHEMA, ]) const CODEX_TURN_ITEM_SCHEMA = z.looseObject({ content: CODEX_TEXT_CONTENT_PART_SCHEMA.array().optional(), diff: z.string().optional(), output_diff: z .looseObject({ diff: z.string().optional(), }) .optional(), role: z.string().optional(), type: z.string().optional(), }) const CODEX_WORKLOG_MESSAGE_SCHEMA = z.looseObject({ author: z.looseObject({ role: z.string().optional() }).optional(), content: z .looseObject({ parts: CODEX_TEXT_CONTENT_PART_SCHEMA.array().optional(), }) .optional(), }) const CODEX_TURN_SCHEMA = z.looseObject({ attempt_placement: z.number().int().nullable().optional(), branch: z.string().nullable().optional(), created_at: z.number().nonnegative().optional(), environment_id: z.string().nullable().optional(), error: z .looseObject({ code: z.string().nullable().optional(), message: z.string().nullable().optional(), }) .nullable() .optional(), id: z.string().optional(), input_items: CODEX_TURN_ITEM_SCHEMA.array().optional(), output_items: CODEX_TURN_ITEM_SCHEMA.array().optional(), pull_request_data: z .looseObject({ url: z.url().optional() }) .nullable() .optional(), pull_request_status: z.string().nullable().optional(), previous_turn_id: z.string().nullable().optional(), role: z.string().optional(), sibling_turn_ids: z.string().array().nullable().optional(), turn_status: z.string().optional(), worklog: z .looseObject({ messages: CODEX_WORKLOG_MESSAGE_SCHEMA.array().optional() }) .nullable() .optional(), }) const CODEX_TASK_STATUS_DISPLAY_SCHEMA = z .looseObject({ branch_name: z.string().optional(), environment_label: z.string().optional(), latest_turn_status_display: z .looseObject({ cancellation_requested_at: z .number() .nonnegative() .nullable() .optional(), diff_stats: z .looseObject({ files_modified: z.number().int().nullable().optional(), lines_added: z.number().int().nullable().optional(), lines_removed: z.number().int().nullable().optional(), }) .nullable() .optional(), sibling_turn_ids: z.string().array().optional(), turn_id: z.string().optional(), turn_status: z.string().optional(), }) .optional(), state: z.string().optional(), }) .optional() const CODEX_PULL_REQUEST_AUTHOR_RESPONSE_SCHEMA = z.looseObject({ avatar_url: z.url().nullable().optional(), email: z.string().nullable().optional(), login: z.string(), name: z.string().nullable().optional(), }) const CODEX_PULL_REQUEST_RESPONSE_SCHEMA = z.looseObject({ additions: z.number().int().nonnegative().optional(), base: z.string(), base_sha: z.string().nullable().optional(), body: z.string().nullable().optional(), changed_files: z.number().int().nonnegative().optional(), commits: z.number().int().nonnegative().optional(), deletions: z.number().int().nonnegative().optional(), draft: z.boolean(), head: z.string(), head_sha: z.string().nullable().optional(), merge_commit_sha: z.string().nullable().optional(), mergeable: z.boolean().nullable().optional(), merged: z.boolean(), number: z.number().int().positive(), state: z.string(), title: z.string(), url: z.url(), user: CODEX_PULL_REQUEST_AUTHOR_RESPONSE_SCHEMA.nullable().optional(), }) const CODEX_EXTERNAL_PULL_REQUEST_RESPONSE_SCHEMA = z.looseObject({ assistant_turn_id: z.string(), codex_updated_sha: z.string().nullable().optional(), id: z.string(), pull_request: CODEX_PULL_REQUEST_RESPONSE_SCHEMA, }) const CODEX_TASK_RESPONSE_SCHEMA = z.looseObject({ archived: z.boolean().optional(), created_at: z.number().nonnegative().optional(), current_turn_id: z.string().nullable().optional(), external_pull_requests: CODEX_EXTERNAL_PULL_REQUEST_RESPONSE_SCHEMA.array().optional(), id: CODEX_CLOUD_TASK_ID_SCHEMA, task_status_display: CODEX_TASK_STATUS_DISPLAY_SCHEMA, title: z.string().optional(), }) const CODEX_CLOUD_TASK_LIST_RESPONSE_SCHEMA = z.object({ cursor: z.string().nullable().optional(), items: z .looseObject({ archived: z.boolean().optional(), created_at: z.number().nonnegative().optional(), id: CODEX_CLOUD_TASK_ID_SCHEMA, pull_requests: CODEX_EXTERNAL_PULL_REQUEST_RESPONSE_SCHEMA.array().optional(), task_status_display: CODEX_TASK_STATUS_DISPLAY_SCHEMA, title: z.string(), updated_at: z.number().nonnegative().optional(), }) .array(), }) const CODEX_CLOUD_TASK_DETAILS_RESPONSE_SCHEMA = z.looseObject({ current_assistant_turn: CODEX_TURN_SCHEMA.nullable().optional(), current_diff_task_turn: CODEX_TURN_SCHEMA.nullable().optional(), current_user_turn: CODEX_TURN_SCHEMA.nullable().optional(), task: CODEX_TASK_RESPONSE_SCHEMA, }) const CODEX_CLOUD_TASK_CREATE_RESPONSE_SCHEMA = z.looseObject({ id: CODEX_CLOUD_TASK_ID_SCHEMA.optional(), task: CODEX_TASK_RESPONSE_SCHEMA.optional(), turn: CODEX_TURN_SCHEMA.optional(), }) const CODEX_CLOUD_TURNS_RESPONSE_SCHEMA = z.looseObject({ current_turn_id: z.string().nullable().optional(), turn_mapping: z.record( z.string(), z.looseObject({ children: z.string().array().optional(), id: z.string(), parent: z.string().nullable().optional(), turn: CODEX_TURN_SCHEMA, }), ), }) const CODEX_CLOUD_ENVIRONMENT_RESPONSE_SCHEMA = z.looseObject({ id: z.string().min(1), is_pinned: z.boolean().optional(), label: z.string().nullable().optional(), task_count: z.number().int().nonnegative().optional(), }) const CODEX_CLOUD_ENVIRONMENT_PAGE_RESPONSE_SCHEMA = z.looseObject({ cursor: z.string().nullable().optional(), items: CODEX_CLOUD_ENVIRONMENT_RESPONSE_SCHEMA.array(), }) const CODEX_CLOUD_MUTATION_RESPONSE_SCHEMA = z.looseObject({ success: z.literal(true), }) const CODEX_CLOUD_MAX_POLLS_SCHEMA = z .number() .int() .min(1) .max(120) .default(36) /** States returned for Codex cloud tasks. */ export const CODEX_CLOUD_TASK_STATUS_SCHEMA = z.enum([ "pending", "ready", "applied", "error", "unknown", ]) /** States returned for one Codex cloud task attempt. */ export const CODEX_CLOUD_TURN_STATUS_SCHEMA = z.enum([ "pending", "inProgress", "completed", "failed", "cancelled", "unknown", ]) /** Latest-attempt diff summary reported for a Codex cloud task. */ export const CODEX_CLOUD_DIFF_SUMMARY_SCHEMA = z.object({ /** Number of files changed by the latest attempt. */ filesChanged: z.number().int().nonnegative(), /** Number of lines added by the latest attempt. */ linesAdded: z.number().int().nonnegative(), /** Number of lines removed by the latest attempt. */ linesRemoved: z.number().int().nonnegative(), }) /** GitHub user attached to a Codex-created pull request. */ export const CODEX_CLOUD_PULL_REQUEST_AUTHOR_SCHEMA = z.object({ avatarUrl: z.url().nullable(), email: z.string().nullable(), login: z.string(), name: z.string().nullable(), }) /** Pull request created from a Codex cloud task. */ export const CODEX_CLOUD_PULL_REQUEST_SCHEMA = z.object({ additions: z.number().int().nonnegative().nullable(), author: CODEX_CLOUD_PULL_REQUEST_AUTHOR_SCHEMA.nullable(), baseBranch: z.string(), baseSha: z.string().nullable(), body: z.string().nullable(), changedFiles: z.number().int().nonnegative().nullable(), codexUpdatedSha: z.string().nullable(), commits: z.number().int().nonnegative().nullable(), deletions: z.number().int().nonnegative().nullable(), draft: z.boolean(), headBranch: z.string(), headSha: z.string().nullable(), id: z.string(), mergeCommitSha: z.string().nullable(), mergeable: z.boolean().nullable(), merged: z.boolean(), number: z.number().int().positive(), state: z.string(), title: z.string(), turnId: z.string(), url: z.url(), }) /** Codex cloud task returned by the remote task service. */ export const CODEX_CLOUD_TASK_SCHEMA = z.object({ /** Whether the task has been archived. */ archived: z.boolean(), /** Number of requested assistant attempts, when reported. */ attemptTotal: z.number().int().positive().nullable(), /** Git branch or ref used by the current task turn. */ branch: z.string().nullable(), /** Stable cloud environment identifier, when reported. */ environmentId: z.string().nullable(), /** Human-readable cloud environment label, when reported. */ environmentLabel: z.string().nullable(), /** Stable Codex cloud task identifier. */ id: CODEX_CLOUD_TASK_ID_SCHEMA, /** Whether this task has created a pull request. */ isReview: z.boolean(), /** Pull requests created from this task. */ pullRequests: CODEX_CLOUD_PULL_REQUEST_SCHEMA.array(), /** Current task state. */ status: CODEX_CLOUD_TASK_STATUS_SCHEMA, /** Diff statistics reported for the latest attempt. */ summary: CODEX_CLOUD_DIFF_SUMMARY_SCHEMA, /** Task title generated by Codex. */ title: z.string(), /** Time at which the task was last updated. */ updatedAt: z.date().nullable(), /** ChatGPT URL for the Codex cloud task. */ url: CODEX_CLOUD_TASK_URL_SCHEMA, }) /** One assistant attempt within a Codex cloud task. */ export const CODEX_CLOUD_TASK_ATTEMPT_SCHEMA = z.object({ attemptNumber: z.number().int().positive().nullable(), createdAt: z.date().nullable(), diff: z.string().nullable(), error: z.string().nullable(), messages: z.string().array(), status: CODEX_CLOUD_TURN_STATUS_SCHEMA, turnId: z.string(), }) /** One user instruction and its assistant attempts in a Codex cloud task. */ export const CODEX_CLOUD_TASK_TURN_SCHEMA = z.object({ createdAt: z.date().nullable(), id: CODEX_CLOUD_TURN_ID_SCHEMA, prompt: z.string().nullable(), attempts: CODEX_CLOUD_TASK_ATTEMPT_SCHEMA.array(), }) /** Detailed Codex cloud task including its latest prompt, output, and diff. */ export const CODEX_CLOUD_TASK_DETAILS_SCHEMA = CODEX_CLOUD_TASK_SCHEMA.extend({ attempts: CODEX_CLOUD_TASK_ATTEMPT_SCHEMA.array(), currentTurnId: CODEX_CLOUD_TURN_ID_SCHEMA.nullable(), prompt: z.string().nullable(), turns: CODEX_CLOUD_TASK_TURN_SCHEMA.array(), }) /** Paginated result returned when listing Codex cloud tasks. */ export const CODEX_CLOUD_TASK_PAGE_SCHEMA = z.object({ /** Cursor for the next page, or `null` when no next page is available. */ cursor: z.string().nullable(), /** Recent Codex cloud tasks matching the requested environment. */ tasks: CODEX_CLOUD_TASK_SCHEMA.array(), }) /** One preconfigured Codex cloud environment. */ export const CODEX_CLOUD_ENVIRONMENT_SCHEMA = z.object({ id: z.string(), isPinned: z.boolean(), label: z.string().nullable(), taskCount: z.number().int().nonnegative(), }) /** Paginated result returned when listing Codex cloud environments. */ export const CODEX_CLOUD_ENVIRONMENT_PAGE_SCHEMA = z.object({ cursor: z.string().nullable(), environments: CODEX_CLOUD_ENVIRONMENT_SCHEMA.array(), }) /** Result returned after Codex accepts a cloud task. */ export const CODEX_CLOUD_TASK_SUBMISSION_SCHEMA = z.object({ /** Stable Codex cloud task identifier. */ id: CODEX_CLOUD_TASK_ID_SCHEMA, /** ChatGPT URL for the new Codex cloud task. */ url: CODEX_CLOUD_TASK_URL_SCHEMA, }) /** Result returned after Codex accepts a follow-up or PR request. */ export const CODEX_CLOUD_TASK_TURN_REFERENCE_SCHEMA = z.object({ taskId: CODEX_CLOUD_TASK_ID_SCHEMA, turnId: CODEX_CLOUD_TURN_ID_SCHEMA, url: CODEX_CLOUD_TASK_URL_SCHEMA, }) /** Lists preconfigured remote Codex cloud environments. */ export const listCodexCloudEnvironments = defineAction( "List Codex cloud environments", ) .describe("Lists remote environments configured for Codex cloud.") .account("codex") .input( z.object({ /** Pagination cursor returned by an earlier request. */ cursor: z.string().min(1).optional(), /** Maximum number of environments to return. */ limit: z.number().int().min(1).max(50).default(50), }), ) .output(CODEX_CLOUD_ENVIRONMENT_PAGE_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => { const url = new URL("wham/environments/search", CODEX_API_BASE_URL) url.searchParams.set("limit", String(input.limit)) if (input.cursor) url.searchParams.set("cursor", input.cursor) const page = CODEX_CLOUD_ENVIRONMENT_PAGE_RESPONSE_SCHEMA.parse( await (await requestCodexCloud(account.secret, url)).json(), ) return { cursor: page.cursor ?? null, environments: page.items.map((environment) => ({ id: environment.id, isPinned: environment.is_pinned ?? false, label: environment.label ?? null, taskCount: environment.task_count ?? 0, })), } }) /** * Starts a remote task in a preconfigured Codex cloud environment. * * Automate.ax only submits the task. Codex performs the coding work in its * OpenAI-hosted cloud environment. */ export const startCodexCloudTask = defineAction("Start Codex cloud task") .describe("Starts a remote task in a Codex cloud environment.") .account("codex") .input( z.object({ /** Number of independent assistant attempts Codex should run. */ attempts: z.number().int().min(1).max(4).default(1), /** Git branch or ref Codex should check out before starting. */ branch: z.string().trim().min(1), /** Stable identifier of a preconfigured Codex cloud environment. */ environmentId: z.string().trim().min(1), /** Whether Codex should answer or make changes. */ mode: z.enum(["code", "ask"]).default("code"), /** Coding task for the remote Codex agent. */ prompt: z.string().trim().min(1), }), ) .output(CODEX_CLOUD_TASK_SUBMISSION_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { const response = CODEX_CLOUD_TASK_CREATE_RESPONSE_SCHEMA.parse( await ( await requestCodexCloud(account.secret, "wham/tasks", { body: JSON.stringify({ input_items: [codexPromptItem(input.prompt)], ...(input.attempts > 1 && { metadata: { best_of_n: input.attempts }, }), new_task: { branch: input.branch, environment_id: input.environmentId, run_environment_in_qa_mode: input.mode === "ask", }, }), method: "POST", }) ).json(), ) const id = response.task?.id ?? response.id if (!id) throw new Error("Codex did not return a cloud task identifier.") return { id, url: codexTaskUrl(id) } }) /** Lists recent Codex cloud tasks. */ export const listCodexCloudTasks = defineAction("List Codex cloud tasks") .describe("Lists recent remote tasks from Codex cloud.") .account("codex") .input( z.object({ /** Pagination cursor returned by an earlier request. */ cursor: z.string().min(1).optional(), /** Optional cloud environment identifier used to filter tasks. */ environmentId: z.string().trim().min(1).optional(), /** Whether to list current, archived, or all tasks. */ filter: z.enum(["current", "archived", "all"]).default("current"), /** Maximum number of tasks to return. */ limit: z.number().int().min(1).max(20).default(20), }), ) .output(CODEX_CLOUD_TASK_PAGE_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => { const url = new URL("wham/tasks/list", CODEX_API_BASE_URL) url.searchParams.set("limit", String(input.limit)) url.searchParams.set("task_filter", input.filter) if (input.cursor) url.searchParams.set("cursor", input.cursor) if (input.environmentId) { url.searchParams.set("environment_id", input.environmentId) } const page = CODEX_CLOUD_TASK_LIST_RESPONSE_SCHEMA.parse( await (await requestCodexCloud(account.secret, url)).json(), ) return { cursor: page.cursor ?? null, tasks: page.items.map((task) => codexTask({ archived: task.archived ?? false, createdAt: task.created_at, pullRequests: task.pull_requests, statusDisplay: task.task_status_display, title: task.title, updatedAt: task.updated_at, id: task.id, }), ), } }) /** Gets one Codex cloud task with its latest output, diff, attempts, and PRs. */ export const getCodexCloudTask = defineAction("Get Codex cloud task") .describe("Gets detailed output and state for one Codex cloud task.") .account("codex") .input( z.object({ /** Stable Codex cloud task identifier. */ taskId: CODEX_CLOUD_TASK_ID_SCHEMA, }), ) .output(CODEX_CLOUD_TASK_DETAILS_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input }) => getNormalizedCodexTaskDetails(account.secret, input.taskId), ) /** Continues an existing Codex cloud task with another remote turn. */ export const continueCodexCloudTask = defineAction("Continue Codex cloud task") .describe("Starts a follow-up turn on an existing Codex cloud task.") .account("codex") .input( z.object({ /** Whether Codex should answer or make changes. */ mode: z.enum(["code", "ask"]).default("code"), /** Follow-up instruction for the remote agent. */ prompt: z.string().trim().min(1), /** Stable Codex cloud task identifier. */ taskId: CODEX_CLOUD_TASK_ID_SCHEMA, /** Assistant turn to continue. Defaults to the task's current turn. */ turnId: CODEX_CLOUD_TURN_ID_SCHEMA.optional(), }), ) .output(CODEX_CLOUD_TASK_TURN_REFERENCE_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { const turnId = input.turnId ?? (await getCodexTaskDetails(account.secret, input.taskId)).task .current_turn_id if (!turnId) throw new Error("Codex task does not have a turn to continue.") const response = CODEX_CLOUD_TASK_CREATE_RESPONSE_SCHEMA.parse( await ( await requestCodexCloud(account.secret, "wham/tasks", { body: JSON.stringify({ follow_up: { environment_mode: input.mode, task_id: input.taskId, turn_id: turnId, }, input_items: [codexPromptItem(input.prompt)], }), method: "POST", }) ).json(), ) const nextTurnId = response.turn?.id ?? response.task?.current_turn_id if (!nextTurnId) { throw new Error("Codex did not return the new cloud turn identifier.") } return { taskId: input.taskId, turnId: nextTurnId, url: codexTaskUrl(input.taskId), } }) /** Cancels the active turn for one Codex cloud task. */ export const cancelCodexCloudTask = defineAction("Cancel Codex cloud task") .describe("Requests cancellation of the active Codex cloud task turn.") .account("codex") .input(z.object({ taskId: CODEX_CLOUD_TASK_ID_SCHEMA })) .output(z.void()) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { await mutateCodexTask(account.secret, input.taskId, "cancel") }) /** Archives one Codex cloud task. */ export const archiveCodexCloudTask = defineAction("Archive Codex cloud task") .describe("Archives one Codex cloud task.") .account("codex") .input(z.object({ taskId: CODEX_CLOUD_TASK_ID_SCHEMA })) .output(z.void()) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { await mutateCodexTask(account.secret, input.taskId, "archive") }) /** Restores one archived Codex cloud task. */ export const unarchiveCodexCloudTask = defineAction( "Unarchive Codex cloud task", ) .describe("Restores one archived Codex cloud task.") .account("codex") .input(z.object({ taskId: CODEX_CLOUD_TASK_ID_SCHEMA })) .output(z.void()) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { await mutateCodexTask(account.secret, input.taskId, "recover") }) /** Starts pull-request creation for a completed Codex cloud task turn. */ export const createCodexCloudPullRequest = defineAction( "Create Codex cloud pull request", ) .describe("Creates a pull request from a completed Codex cloud task turn.") .account("codex") .input( z.object({ /** Create the pull request as a draft. */ draft: z.boolean().default(false), /** Stable Codex cloud task identifier. */ taskId: CODEX_CLOUD_TASK_ID_SCHEMA, /** Completed assistant turn containing the diff. */ turnId: CODEX_CLOUD_TURN_ID_SCHEMA.optional(), }), ) .output(CODEX_CLOUD_TASK_TURN_REFERENCE_SCHEMA) .retry({ replaySafety: "unsafe" }) .handler(async ({ account, input }) => { const turnId = input.turnId ?? (await getCodexTaskDetails(account.secret, input.taskId)).task .current_turn_id if (!turnId) throw new Error("Codex task does not have a turn to publish.") await requestCodexCloud( account.secret, `wham/tasks/${encodeURIComponent(input.taskId)}/turns/${encodeURIComponent(turnId)}/pr`, { body: JSON.stringify(input.draft ? { mode: "draft" } : {}), method: "POST", }, ) return { taskId: input.taskId, turnId, url: codexTaskUrl(input.taskId) } }) const CODEX_CLOUD_POLL_DELIVERY_SCHEMA = z.object({ account: CODEX_CLOUD_ACCOUNT_SELECTION_SCHEMA, entrypoint: z.string().startsWith(CODEX_POLL_ENTRYPOINT_PREFIX), key: z.string().min(1), maxPolls: CODEX_CLOUD_MAX_POLLS_SCHEMA, poll: z.number().int().positive(), pollInterval: z.union([z.string().min(1), z.number().nonnegative()]), taskId: CODEX_CLOUD_TASK_ID_SCHEMA, }) const CODEX_CLOUD_POLL_START_SCHEMA = z.object({ key: z.string().min(1) }) const CODEX_CLOUD_POLL_RESULT_SCHEMA = z.discriminatedUnion("status", [ z.object({ key: z.string().min(1), status: z.literal("pending") }), z.object({ key: z.string().min(1), status: z.literal("completed"), task: CODEX_CLOUD_TASK_DETAILS_SCHEMA, }), ]) const scheduleCodexCloudPollAction = defineAction("Wait for Codex cloud task") .input(CODEX_CLOUD_POLL_DELIVERY_SCHEMA.omit({ key: true, poll: true })) .output(CODEX_CLOUD_POLL_START_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ input, runtime }) => { const key = `${runtime.contextId}:${input.entrypoint}` await runtime.invokeAutomation({ automation: runtime.automationId, entrypoint: input.entrypoint, payload: { ...input, key, poll: 1 }, }) return { key } }) const pollCodexCloudTaskAction = defineAction("Check Codex cloud task") .account("codex") .input(CODEX_CLOUD_POLL_DELIVERY_SCHEMA) .output(CODEX_CLOUD_POLL_RESULT_SCHEMA) .retry({ replaySafety: "safe" }) .handler(async ({ account, input, runtime }) => { const task = await getNormalizedCodexTaskDetails( account.secret, input.taskId, ) if (task.status !== "pending") { return { key: input.key, status: "completed", task } as const } if (input.poll === input.maxPolls) { throw new Error( `Codex cloud task ${input.taskId} remained pending after ${input.maxPolls} checks at ${input.pollInterval} intervals.`, ) } await runtime.invokeAutomation({ automation: runtime.automationId, entrypoint: input.entrypoint, payload: { ...input, poll: input.poll + 1 }, runIn: input.pollInterval, }) return { key: input.key, status: "pending" } as const }) /** Input for waiting on a previously submitted Codex cloud task. */ export interface CodexCloudTaskWaitInput { /** Maximum number of status checks before the wait fails. Defaults to 36. */ maxPolls?: number /** Durable delay between status checks. Defaults to `"5m"`. */ pollInterval?: DurationString | number /** Stable task ID or compatible task signal. */ taskId: Extract< Parameters[0], { taskId: unknown } >["taskId"] } /** Input for starting and durably waiting on a Codex cloud task. */ export type RunCodexCloudTaskInput = Extract< Parameters[0], { prompt: unknown } > & Pick /** Account selection for composed Codex cloud task helpers. */ export type CodexCloudTaskAccountOptions = NonNullable< Parameters[1] > /** * Waits for a Codex cloud task through bounded, durable status checks. * * No action invocation remains open between checks. Each interval uses a * durable delivery root and a correlated child context. The returned signal * emits the first non-pending task. * * @param input - Task identity, interval, and polling budget. * @param options - Optional Codex account selection. */ export function waitForCodexCloudTask( input: CodexCloudTaskWaitInput, options?: CodexCloudTaskAccountOptions, ) { const maxPolls = CODEX_CLOUD_MAX_POLLS_SCHEMA.parse(input.maxPolls) const prerequisites = getCurrentSignalPrerequisites() const contextPrerequisite = getCurrentContextPrerequisite() return group( { name: "Wait for Codex cloud task", presentation: "hidden" }, () => withoutSignalPrerequisites(() => { const entrypoint = `${CODEX_POLL_ENTRYPOINT_PREFIX}${getCurrentHookScopePath().join(".")}` const delivery = createSubscription< z.output >(automationInvokedTriggerDefinition, { entrypoint }, undefined, { inferActionBoundary: false, }) const schedule = () => scheduleCodexCloudPollAction({ account: options?.account ?? "default", entrypoint, maxPolls, pollInterval: input.pollInterval ?? "5m", taskId: input.taskId, }) const scheduleWithContext = () => contextPrerequisite ? withContextPrerequisite(contextPrerequisite, schedule) : schedule() // The starter remains separate because its output anchors terminal correlation. const scheduled = prerequisites ? scope(() => withPrerequisites(prerequisites, scheduleWithContext)) : scheduleWithContext() const terminal = filter( transform( [ delivery, outcome( pollCodexCloudTaskAction(delivery, { account: delivery.account, }), ), ], ({ key }, result) => ({ key, result }), ), ({ result }) => result.status !== "succeeded" || result.value.status === "completed", ) return transform( [ correlate([ scheduled.keyBy(({ key }) => key), terminal.keyBy(({ key }) => key), ]), terminal, ], (_matched, { result }) => { if (result.status === "failed") { throw new FailedSignalError(result.failure) } if (result.status === "closed") { throw new Error("Codex cloud polling closed without a result.") } if (result.value.status !== "completed") { throw new Error("Codex cloud polling completed without a task.") } return result.value.task }, ) }), ) } /** * Starts a Codex cloud task and continues after it reaches a terminal state. * * @param input - Task submission fields and durable polling configuration. * @param options - Optional Codex account selection. */ export function runCodexCloudTask( input: RunCodexCloudTaskInput, options?: CodexCloudTaskAccountOptions, ) { const { maxPolls, pollInterval, ...submission } = input return waitForCodexCloudTask( { ...(maxPolls === undefined ? {} : { maxPolls }), ...(pollInterval === undefined ? {} : { pollInterval }), taskId: startCodexCloudTask(submission, options).id, }, options, ) } /** * Sends one request to the undocumented ChatGPT Codex cloud backend. * * @param secret - Resolved Codex integration secret. * @param path - Relative backend path or complete backend URL. * @param init - Fetch options excluding authentication headers. */ async function requestCodexCloud( secret: Record, path: string | URL, init: RequestInit = {}, ) { const headers = await codexAuthHeaders(secret) if (init.body) headers.set("Content-Type", "application/json") headers.set("User-Agent", CODEX_USER_AGENT) const response = await fetch( typeof path === "string" ? new URL(path, CODEX_API_BASE_URL) : path, { ...init, headers, signal: init.signal ?? AbortSignal.timeout(CODEX_REQUEST_TIMEOUT_MS), }, ) if (!response.ok) throw await codexApiError("request", response) return response } /** * Retrieves and normalizes one detailed Codex task response. * * @param secret - Resolved Codex integration secret. * @param taskId - Stable Codex cloud task identifier. */ async function getNormalizedCodexTaskDetails( secret: Record, taskId: string, ) { const [details, turnGraph] = await Promise.all([ getCodexTaskDetails(secret, taskId), getCodexTaskTurns(secret, taskId), ]) const assistantTurn = details.current_assistant_turn const diffTurn = details.current_diff_task_turn const turns = codexTaskTurns( turnGraph, assistantTurn?.id ? { diffTurn, turnId: assistantTurn.id } : undefined, ) const currentTurn = turns.find((turn) => turn.attempts.some( (attempt) => attempt.turnId === (assistantTurn?.id ?? details.task.current_turn_id), ), ) return { ...codexTask({ archived: details.task.archived ?? false, branch: assistantTurn?.branch ?? diffTurn?.branch ?? details.task.task_status_display?.branch_name, createdAt: details.task.created_at, environmentId: assistantTurn?.environment_id ?? diffTurn?.environment_id, id: details.task.id, pullRequests: details.task.external_pull_requests, statusDisplay: details.task.task_status_display, title: details.task.title ?? "Untitled Codex task", updatedAt: assistantTurn?.created_at ?? diffTurn?.created_at ?? details.current_user_turn?.created_at, }), attempts: currentTurn?.attempts ?? (assistantTurn?.id ? [codexAttempt(assistantTurn, codexTurnDiff(diffTurn, assistantTurn))] : []), currentTurnId: assistantTurn?.id ?? details.task.current_turn_id ?? null, prompt: currentTurn?.prompt ?? (codexTurnMessages(details.current_user_turn).join("\n\n") || null), turns, } } /** * Resolves one Codex cloud task details response. * * @param secret - Resolved Codex integration secret. * @param taskId - Stable Codex cloud task identifier. */ async function getCodexTaskDetails( secret: Record, taskId: string, ) { return CODEX_CLOUD_TASK_DETAILS_RESPONSE_SCHEMA.parse( await ( await requestCodexCloud( secret, `wham/tasks/${encodeURIComponent(taskId)}`, ) ).json(), ) } /** * Resolves the complete user and assistant turn graph for one cloud task. * * @param secret - Resolved Codex integration secret. * @param taskId - Stable Codex cloud task identifier. */ async function getCodexTaskTurns( secret: Record, taskId: string, ) { return CODEX_CLOUD_TURNS_RESPONSE_SCHEMA.parse( await ( await requestCodexCloud( secret, `wham/tasks/${encodeURIComponent(taskId)}/turns`, ) ).json(), ) } /** * Applies a status-only task mutation. * * @param secret - Resolved Codex integration secret. * @param taskId - Stable Codex cloud task identifier. * @param operation - Provider mutation path. */ async function mutateCodexTask( secret: Record, taskId: string, operation: "archive" | "cancel" | "recover", ) { CODEX_CLOUD_MUTATION_RESPONSE_SCHEMA.parse( await ( await requestCodexCloud( secret, `wham/tasks/${encodeURIComponent(taskId)}/${operation}`, { method: "POST" }, ) ).json(), ) } /** * Resolves the account-routing headers required by Codex access tokens. * * @param secret - Resolved Codex integration secret. */ async function codexAuthHeaders(secret: Record) { const { apiKey } = CODEX_SECRET_SCHEMA.parse(secret) const response = await fetch(CODEX_AUTH_URL, { headers: { Accept: "application/json", Authorization: `Bearer ${apiKey}`, "User-Agent": CODEX_USER_AGENT, }, signal: AbortSignal.timeout(CODEX_REQUEST_TIMEOUT_MS), }) if (!response.ok) throw await codexApiError("authentication", response) const metadata = CODEX_ACCESS_TOKEN_METADATA_SCHEMA.parse( await response.json(), ) return new Headers({ Accept: "application/json", Authorization: `Bearer ${apiKey}`, "ChatGPT-Account-ID": metadata.chatgpt_account_id, "User-Agent": CODEX_USER_AGENT, ...(metadata.chatgpt_account_is_fedramp && { "X-OpenAI-Fedramp": "true", }), }) } /** * Creates a provider error without exposing the access token. * * @param operation - Short description of the failed provider operation. * @param response - Unsuccessful provider response. */ async function codexApiError(operation: string, response: Response) { const detail = codexErrorDetail( (await response.text()).trim(), response.headers.get("content-type"), ) const requestId = response.headers.get("x-request-id") ?? response.headers.get("cf-ray") ?? undefined return new Error( `Codex cloud ${operation} failed with HTTP ${response.status}${detail ? `: ${detail}` : ""}${requestId ? ` (request ${requestId})` : ""}.`, ) } /** * Selects a bounded provider error message without returning HTML pages. * * @param body - Raw provider response body. * @param contentType - Provider response content type. */ function codexErrorDetail(body: string, contentType: string | null) { if (!body) return "" if (contentType?.includes("application/json")) { try { const parsed = CODEX_ERROR_RESPONSE_SCHEMA.safeParse(JSON.parse(body)) if (!parsed.success) return "" return ( ( parsed.data.message ?? parsed.data.detail ?? parsed.data.error?.message ) ?.trim() .slice(0, 512) ?? "" ) } catch { return "" } } return contentType?.startsWith("text/plain") ? body.slice(0, 512) : "" } /** * Normalizes one provider task into the public task shape. * * @param input - Provider task fields to normalize. * @param input.archived - Whether the task is archived. * @param input.branch - Git branch or ref used by the task. * @param input.createdAt - Provider creation timestamp. * @param input.environmentId - Stable cloud environment identifier. * @param input.id - Stable Codex task identifier. * @param input.pullRequests - Pull requests attached to the task. * @param input.statusDisplay - Provider task status display. * @param input.title - Human-readable task title. * @param input.updatedAt - Provider update timestamp. */ function codexTask(input: { archived: boolean branch?: string | null createdAt?: number environmentId?: string | null id: string pullRequests?: z.output[] statusDisplay: z.output title: string updatedAt?: number }) { const pullRequests = (input.pullRequests ?? []).map(codexPullRequest) const updatedAtSeconds = input.updatedAt ?? input.createdAt return { archived: input.archived, attemptTotal: input.statusDisplay?.latest_turn_status_display?.sibling_turn_ids === undefined ? null : input.statusDisplay.latest_turn_status_display.sibling_turn_ids .length + 1, branch: input.branch ?? input.statusDisplay?.branch_name ?? null, environmentId: input.environmentId ?? null, environmentLabel: input.statusDisplay?.environment_label ?? null, id: input.id, isReview: pullRequests.length > 0, pullRequests, status: codexTaskStatus(input.statusDisplay), summary: codexDiffSummary(input.statusDisplay), title: input.title, updatedAt: updatedAtSeconds === undefined ? null : new Date(updatedAtSeconds * 1_000), url: codexTaskUrl(input.id), } } /** * Normalizes one provider pull request. * * @param external - Provider pull-request wrapper. */ function codexPullRequest( external: z.output, ) { const pullRequest = external.pull_request return { additions: pullRequest.additions ?? null, author: pullRequest.user ? { avatarUrl: pullRequest.user.avatar_url ?? null, email: pullRequest.user.email ?? null, login: pullRequest.user.login, name: pullRequest.user.name ?? null, } : null, baseBranch: pullRequest.base, baseSha: pullRequest.base_sha ?? null, body: pullRequest.body ?? null, changedFiles: pullRequest.changed_files ?? null, codexUpdatedSha: external.codex_updated_sha ?? null, commits: pullRequest.commits ?? null, deletions: pullRequest.deletions ?? null, draft: pullRequest.draft, headBranch: pullRequest.head, headSha: pullRequest.head_sha ?? null, id: external.id, mergeCommitSha: pullRequest.merge_commit_sha ?? null, mergeable: pullRequest.mergeable ?? null, merged: pullRequest.merged, number: pullRequest.number, state: pullRequest.state, title: pullRequest.title, turnId: external.assistant_turn_id, url: pullRequest.url, } } /** * Normalizes the provider turn graph into user instructions and attempts. * * @param graph - Complete provider turn graph. * @param current - Current assistant turn and separately reported diff. * @param current.diffTurn - Separately reported current diff turn. * @param current.turnId - Current assistant turn identifier. */ function codexTaskTurns( graph: z.output, current?: { diffTurn: z.output | null | undefined turnId: string }, ) { return Object.values(graph.turn_mapping) .filter(({ turn }) => (turn.input_items ?? []).some((item) => item.role === "user"), ) .map((userNode) => { // Keep the normalized attempts visible beside the user-turn fields below. const attempts = (userNode.children ?? []) .flatMap((childId) => { const assistantNode = graph.turn_mapping[childId] if (!assistantNode) return [] const assistantTurn = { ...assistantNode.turn, id: assistantNode.turn.id ?? assistantNode.id, } return [ codexAttempt( assistantTurn, codexTurnDiff( assistantTurn.id === current?.turnId ? current.diffTurn : undefined, assistantTurn, ), ), ] }) .toSorted( (left, right) => (left.attemptNumber ?? Number.MAX_SAFE_INTEGER) - (right.attemptNumber ?? Number.MAX_SAFE_INTEGER), ) const userTurn = { ...userNode.turn, id: userNode.turn.id ?? userNode.id, } return { attempts, createdAt: userTurn.created_at ? new Date(userTurn.created_at * 1_000) : null, id: userTurn.id, prompt: codexTurnMessages(userTurn).join("\n\n") || null, } }) .toSorted( (left, right) => (left.createdAt?.getTime() ?? 0) - (right.createdAt?.getTime() ?? 0), ) } /** * Normalizes one provider assistant turn. * * @param turn - Provider assistant turn. * @param diff - Unified diff selected for the attempt. */ function codexAttempt( turn: z.output, diff: string | null, ) { const code = turn.error?.code?.trim() const message = turn.error?.message?.trim() return { attemptNumber: turn.attempt_placement === undefined || turn.attempt_placement === null ? null : turn.attempt_placement + 1, createdAt: turn.created_at === undefined ? null : new Date(turn.created_at * 1_000), diff, error: code && message ? `${code}: ${message}` : (code ?? message ?? null), messages: codexTurnMessages(turn), status: codexTurnStatus(turn.turn_status), turnId: z.string().parse(turn.id), } } /** * Maps provider task state into the stable public action state. * * @param display - Provider task status display. */ function codexTaskStatus( display: z.output, ): z.output { if ( display?.state === "pending" || display?.state === "ready" || display?.state === "applied" || display?.state === "error" ) { return display.state } const turnStatus = display?.latest_turn_status_display?.turn_status if (turnStatus === "completed") return "ready" if (turnStatus === "failed" || turnStatus === "cancelled") return "error" if (turnStatus === "in_progress" || turnStatus === "pending") return "pending" return display?.state || turnStatus ? "unknown" : "pending" } /** * Maps one provider turn state into the public turn state. * * @param status - Provider turn status. */ function codexTurnStatus( status: string | undefined, ): z.output { if (status === "in_progress") return "inProgress" if ( status === "pending" || status === "completed" || status === "failed" || status === "cancelled" ) { return status } return "unknown" } /** * Extracts nonnegative diff statistics from the latest task attempt. * * @param display - Provider task status display. */ function codexDiffSummary( display: z.output, ) { const stats = display?.latest_turn_status_display?.diff_stats return { filesChanged: Math.max(0, stats?.files_modified ?? 0), linesAdded: Math.max(0, stats?.lines_added ?? 0), linesRemoved: Math.max(0, stats?.lines_removed ?? 0), } } /** * Extracts assistant or user text messages from one turn. * * @param turn - Provider user or assistant turn. */ function codexTurnMessages( turn: z.output | null | undefined, ) { const output = (turn?.output_items ?? turn?.input_items ?? []).flatMap( (item) => item.type === "message" ? (item.content ?? []).flatMap((part) => { if (typeof part === "string") return part ? [part] : [] return part.content_type === "text" && part.text ? [part.text] : [] }) : [], ) const worklog = (turn?.worklog?.messages ?? []).flatMap((message) => message.author?.role === "assistant" ? (message.content?.parts ?? []).flatMap((part) => { if (typeof part === "string") return part ? [part] : [] return part.content_type === "text" && part.text ? [part.text] : [] }) : [], ) return output.length > 0 ? output : worklog } /** * Extracts a unified diff from the first turn that contains one. * * @param turns - Provider turns in preferred order. */ function codexTurnDiff( ...turns: (z.output | null | undefined)[] ) { for (const turn of turns) { for (const item of turn?.output_items ?? []) { if (item.type === "output_diff" && item.diff) return item.diff if (item.type === "pr" && item.output_diff?.diff) { return item.output_diff.diff } } } return null } /** * Creates one provider-native text prompt item. * * @param prompt - User prompt text. */ function codexPromptItem(prompt: string) { return { content: [{ content_type: "text", text: prompt }], role: "user", type: "message", } } /** * Returns the ChatGPT page for one Codex cloud task. * * @param id - Stable Codex cloud task identifier. */ function codexTaskUrl(id: string) { return `${CODEX_TASK_URL_BASE}${encodeURIComponent(id)}` }