File size: 12,555 Bytes
eb3f11e | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 | /**
* 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,
},
},
};
}
|