Download src/agents/embedded-agent-runner/cli-backend-dispatch.ts from SaylorTwift/openclaw: direct link, hf CLI and curl.
- Browser
- Download file 12.6 kB
-
https://huggingface.co/SaylorTwift/openclaw/resolve/main/src/agents/embedded-agent-runner/cli-backend-dispatch.ts
- Command line
-
hf download hf://SaylorTwift/openclaw/src/agents/embedded-agent-runner/cli-backend-dispatch.ts
-
curl -L -o cli-backend-dispatch.ts https://huggingface.co/SaylorTwift/openclaw/resolve/main/src/agents/embedded-agent-runner/cli-backend-dispatch.ts
12.6 kB
| /** | |
| * Opt-in CLI-backend dispatch for one-shot embedded runs. | |
| * | |
| * Embedded runs targeting a CLI runtime provider normally fall through to the | |
| * openclaw harness and call the provider API directly with that runtime's | |
| * credentials (`cli_runtime_passthrough_openclaw`). Anthropic routes direct | |
| * anthropic-messages calls on subscription OAuth tokens to metered "extra | |
| * usage" billing: without extra-usage balance the passthrough fails closed | |
| * with a billing error, and with it the run silently draws paid usage instead | |
| * of the plan limits the CLI runtime was configured for. Callers that | |
| * tolerate CLI latency opt in via `cliBackendDispatch: "subscription-auth"` | |
| * to run through the CLI backend on plan limits instead. | |
| */ | |
| import { isRecord } from "@openclaw/normalization-core/record-coerce"; | |
| import type { SessionTranscriptRuntimeTarget } from "../../config/sessions/session-accessor.js"; | |
| import { onAgentEventForRun } from "../../infra/agent-events.js"; | |
| import { createSubsystemLogger } from "../../logging/subsystem.js"; | |
| import { resolvePreparedRunAdmission } from "../admitted-run-context.js"; | |
| import { stripOpenClawMcpToolPrefix } from "../cli-runner/tool-policy.js"; | |
| import { normalizeToolPolicyName } from "../tool-policy.js"; | |
| import { isToolResultError } from "../tool-result-error.js"; | |
| import { resolveEmbeddedCliBackendDispatchEligibility } from "./cli-backend-dispatch-eligibility.js"; | |
| import { createCliDispatchTranscriptRecorder } from "./cli-backend-dispatch-transcript.js"; | |
| import type { RunEmbeddedAgentInternalParams } from "./run/internal-params.js"; | |
| import type { RunEmbeddedAgentParams } from "./run/params.js"; | |
| import type { EmbeddedAgentRunResult } from "./types.js"; | |
| const log = createSubsystemLogger("agents/embedded-cli-dispatch"); | |
| type CliBackendDispatchParams = RunEmbeddedAgentInternalParams & { | |
| sessionTarget: SessionTranscriptRuntimeTarget; | |
| }; | |
| type EmbeddedCliBackendDispatch = { | |
| provider: string; | |
| sessionFile: string; | |
| /** Named loopback allowlist; the dispatch gate guarantees it is non-empty. */ | |
| toolsAllow: string[]; | |
| }; | |
| /** | |
| * Runs the embedded turn through the CLI backend when the opt-in dispatch | |
| * gate matches; returns undefined so the caller continues on the native path. | |
| */ | |
| export async function runEmbeddedAgentViaCliBackendIfEligible( | |
| params: CliBackendDispatchParams, | |
| ): Promise<EmbeddedAgentRunResult | undefined> { | |
| const dispatch = resolveEmbeddedCliBackendDispatch(params); | |
| return dispatch ? await runEmbeddedAgentViaCliBackend(params, dispatch) : undefined; | |
| } | |
| /** Applies the opt-in and transcript-path gates on top of shared eligibility. */ | |
| function resolveEmbeddedCliBackendDispatch( | |
| params: RunEmbeddedAgentParams, | |
| ): EmbeddedCliBackendDispatch | undefined { | |
| if (params.cliBackendDispatch !== "subscription-auth") { | |
| return undefined; | |
| } | |
| // The one-shot bridge cannot carry authenticated source-channel delivery | |
| // context; private source replies must stay with their embedded owner. | |
| if (params.sourceReplyDeliveryMode === "message_tool_only") { | |
| return undefined; | |
| } | |
| // The CLI runner needs the caller-owned transcript path; runs without one | |
| // stay on the passthrough where session targets are resolved internally. | |
| const sessionFile = params.sessionFile?.trim(); | |
| if (!sessionFile) { | |
| return undefined; | |
| } | |
| const toolsAllow = resolveDispatchableToolsAllow(params); | |
| if (!toolsAllow) { | |
| return undefined; | |
| } | |
| const eligibility = resolveEmbeddedCliBackendDispatchEligibility(params); | |
| return eligibility ? { provider: eligibility.provider, sessionFile, toolsAllow } : undefined; | |
| } | |
| /** | |
| * Fail closed on tool policy: dispatch only runs whose embedded tool state the | |
| * CLI bridge can express faithfully — a non-empty named allowlist bounded by | |
| * the loopback grant. Deny-all (`[]`), wildcards, absent allowlists, and | |
| * flag-based restrictions (`disableTools`, `modelRun`) keep the embedded | |
| * passthrough so no closed state silently widens on the CLI surface; full | |
| * translation can arrive with the first caller that needs it (#57326). | |
| */ | |
| function resolveDispatchableToolsAllow(params: RunEmbeddedAgentParams): string[] | undefined { | |
| if (params.disableTools || params.modelRun) { | |
| return undefined; | |
| } | |
| if (!params.toolsAllow || params.toolsAllow.length === 0) { | |
| return undefined; | |
| } | |
| const names = params.toolsAllow.map((name) => normalizeToolPolicyName(name)); | |
| if (names.some((name) => !name || name === "*" || name.includes("*"))) { | |
| return undefined; | |
| } | |
| return [...new Set(names)]; | |
| } | |
| /** Runs an opted-in embedded run through the CLI backend as a one-shot turn. */ | |
| async function runEmbeddedAgentViaCliBackend( | |
| params: CliBackendDispatchParams, | |
| dispatch: EmbeddedCliBackendDispatch, | |
| ): Promise<EmbeddedAgentRunResult> { | |
| const { runCliAgent } = await import("../cli-runner.runtime.js"); | |
| const admittedRunContext = await resolvePreparedRunAdmission({ | |
| runId: params.runId, | |
| runtimeKind: "embedded", | |
| admittedRunContext: params.admittedRunContext, | |
| preparedRunAdmission: params.preparedRunAdmission, | |
| }); | |
| // The dispatch gate guarantees a non-empty named allowlist; translate it to | |
| // the selectable-backend surface: no native tools, only the listed loopback | |
| // MCP tools. The MCP list also bounds the loopback grant server-side (tools | |
| // outside it can be neither listed nor called) and makes prepare serve the | |
| // loopback exclusively, so the message tool and user/plugin MCP servers stay | |
| // unreachable, matching disableMessageTool intent. | |
| const cliToolAvailability = { | |
| native: [] as [], | |
| openClaw: dispatch.toolsAllow, | |
| }; | |
| const onAgentToolResult = params.onAgentToolResult; | |
| const { storePath, expectedLifecycleRevision, expectedWriterRunId } = params.sessionTarget; | |
| // Durable turns mirror CLI output for transcript readers and timeout salvage. | |
| // Detached runs may borrow the identity without owning its transcript. | |
| const transcript = | |
| params.sessionManager || params.sessionPersistence === "detached" | |
| ? undefined | |
| : createCliDispatchTranscriptRecorder({ | |
| sessionId: params.sessionId, | |
| sessionKey: params.sessionKey, | |
| agentId: params.agentId, | |
| storePath, | |
| sessionFile: dispatch.sessionFile, | |
| runId: params.runId, | |
| prompt: params.prompt, | |
| provider: dispatch.provider, | |
| model: params.model, | |
| cwd: params.cwd ?? params.workspaceDir, | |
| config: params.config, | |
| expectedLifecycleRevision, | |
| expectedWriterRunId, | |
| ...(params.senderIsOwner !== undefined ? { senderIsOwner: params.senderIsOwner } : {}), | |
| }); | |
| // CLI tool results arrive as agent events with transport-prefixed MCP | |
| // names; strip and normalize so observers and transcript records see the | |
| // same tool names and soft-error signal the native embedded path reports. | |
| const unsubscribe = onAgentEventForRun(params.runId, (evt) => { | |
| if (evt.runId !== params.runId) { | |
| return; | |
| } | |
| if (evt.stream === "assistant" && typeof evt.data.text === "string") { | |
| transcript?.noteAssistantText(evt.data.text); | |
| return; | |
| } | |
| if (evt.stream !== "tool") { | |
| return; | |
| } | |
| const phase = evt.data.phase; | |
| if (phase !== "start" && phase !== "result") { | |
| return; | |
| } | |
| const rawName = typeof evt.data.name === "string" ? evt.data.name : ""; | |
| if (!rawName) { | |
| return; | |
| } | |
| const toolName = normalizeToolPolicyName(stripOpenClawMcpToolPrefix(rawName)); | |
| const toolCallId = typeof evt.data.toolCallId === "string" ? evt.data.toolCallId : undefined; | |
| if (phase === "start") { | |
| transcript?.noteToolEvent({ | |
| phase, | |
| toolName, | |
| toolCallId, | |
| args: isRecord(evt.data.args) ? evt.data.args : undefined, | |
| }); | |
| return; | |
| } | |
| const isError = evt.data.isError === true || isToolResultError(evt.data.result); | |
| const resultContentSource = evt.data.resultContentSource === "network" ? "network" : undefined; | |
| transcript?.noteToolEvent({ | |
| phase, | |
| toolName, | |
| toolCallId, | |
| result: evt.data.result, | |
| isError, | |
| ...(resultContentSource ? { resultContentSource } : {}), | |
| }); | |
| onAgentToolResult?.({ | |
| toolName, | |
| result: evt.data.result, | |
| isError, | |
| }); | |
| }); | |
| // The killed CLI child can take seconds to settle after a timeout abort, | |
| // while the caller's partial-text salvage reads the session file within a | |
| // short grace window; flush the latest snapshot the moment abort fires. | |
| const flushOnAbort = () => transcript?.flushAssistantSnapshot(); | |
| params.abortSignal?.addEventListener("abort", flushOnAbort, { once: true }); | |
| // Reply/cron callers advance lifecycle state and arm execution-phase | |
| // watchdogs on this signal; dispatched runs emit it at the same | |
| // post-admission boundary where the native path does. | |
| params.onExecutionStarted?.( | |
| params.lifecycleGeneration !== undefined | |
| ? { lifecycleGeneration: params.lifecycleGeneration } | |
| : undefined, | |
| ); | |
| log.info( | |
| `dispatching embedded run through CLI backend: runId=${params.runId} provider=${dispatch.provider} model=${params.model ?? ""}`, | |
| ); | |
| let finalAssistantText: string | undefined; | |
| try { | |
| const result = await runCliAgent({ | |
| admittedRunContext, | |
| sessionManager: params.sessionManager, | |
| sessionId: params.sessionId, | |
| sessionKey: params.sessionKey, | |
| sessionTarget: params.sessionTarget, | |
| expectedLifecycleRevision, | |
| expectedWriterRunId, | |
| chatType: params.chatType, | |
| agentId: params.agentId, | |
| storePath, | |
| trigger: params.trigger, | |
| sessionFile: dispatch.sessionFile, | |
| workspaceDir: params.workspaceDir, | |
| agentDir: params.agentDir, | |
| config: params.config, | |
| prompt: params.prompt, | |
| imagePrompt: params.prompt, | |
| images: params.images, | |
| imageOrder: params.imageOrder, | |
| media: params.media, | |
| provider: dispatch.provider, | |
| model: params.model, | |
| ...(params.requestedRouteResolution === "resolved" && params.provider && params.model | |
| ? { requesterModel: { provider: params.provider, model: params.model } } | |
| : {}), | |
| authProfileId: params.authProfileId, | |
| modelHasVision: params.modelHasVision, | |
| contextWindow: params.contextWindow, | |
| thinkLevel: params.thinkLevel, | |
| fastMode: params.fastMode, | |
| fastModeStartedAtMs: params.fastModeStartedAtMs, | |
| fastModeAutoOnSeconds: params.fastModeAutoOnSeconds, | |
| timeoutMs: params.timeoutMs, | |
| runTimeoutOverrideMs: params.runTimeoutOverrideMs ?? params.timeoutMs, | |
| runId: params.runId, | |
| lifecycleGeneration: params.lifecycleGeneration, | |
| lane: params.lane, | |
| extraSystemPrompt: params.extraSystemPrompt, | |
| messageChannel: params.messageChannel, | |
| messageProvider: params.messageProvider, | |
| bootstrapContextMode: params.bootstrapContextMode, | |
| bootstrapContextRunKind: params.bootstrapContextRunKind, | |
| abortSignal: params.abortSignal, | |
| onBlockReply: params.onBlockReply, | |
| onPartialReply: params.onPartialReply, | |
| onExecutionPhase: params.onExecutionPhase, | |
| cliToolAvailability, | |
| // One-shot helper run: fresh CLI process, no warm live session left | |
| // behind, and no implicit message sends without an explicit target. | |
| disableCliLiveSession: true, | |
| cleanupCliLiveSessionOnRunEnd: true, | |
| requireExplicitMessageTarget: true, | |
| cleanupBundleMcpOnRunEnd: params.cleanupBundleMcpOnRunEnd, | |
| }); | |
| finalAssistantText = result.payloads?.find( | |
| (payload) => payload.isReasoning !== true && typeof payload.text === "string", | |
| )?.text; | |
| return withoutCliSessionBinding(result); | |
| } finally { | |
| params.abortSignal?.removeEventListener("abort", flushOnAbort); | |
| unsubscribe(); | |
| // Flush before the promise settles: timeout salvage reads the session | |
| // file as soon as the caller observes the rejection. | |
| await transcript?.finalize(finalAssistantText); | |
| } | |
| } | |
| /** Dispatch runs own no session entry, so a returned CLI binding has no owner to persist it. */ | |
| function withoutCliSessionBinding(result: EmbeddedAgentRunResult): EmbeddedAgentRunResult { | |
| const agentMeta = result.meta.agentMeta; | |
| if (!agentMeta?.cliSessionBinding && agentMeta?.clearCliSessionBinding !== true) { | |
| return result; | |
| } | |
| return { | |
| ...result, | |
| meta: { | |
| ...result.meta, | |
| agentMeta: { | |
| ...agentMeta, | |
| cliSessionBinding: undefined, | |
| clearCliSessionBinding: undefined, | |
| }, | |
| }, | |
| }; | |
| } | |