openclaw / src /agents /embedded-agent-runner /cli-backend-dispatch.ts
SaylorTwift's picture
SaylorTwift HF Staff
Add files using upload-large-folder tool
eb3f11e verified
Raw History Blame Contribute Delete
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,
},
},
};
}