refactor(lobster): consolidate embedded workflow execution (#119630)

This commit is contained in:
Peter Steinberger
2026-08-05 07:22:43 -07:00
committed by GitHub
parent 1b4a60be15
commit 08b9c4c3e1
3 changed files with 92 additions and 260 deletions
+40 -102
View File
@@ -45,32 +45,19 @@ type EmbeddedToolContext = {
stdout?: NodeJS.WritableStream;
stderr?: NodeJS.WritableStream;
signal?: AbortSignal;
registry?: unknown;
llmAdapters?: Record<string, unknown>;
};
type EmbeddedToolEnvelope = {
protocolVersion?: number;
ok: boolean;
status?: "ok" | "needs_approval" | "needs_input" | "cancelled";
output?: unknown[];
requiresApproval?: {
type?: "approval_request";
prompt: string;
items: unknown[];
preview?: string;
resumeToken?: string;
approvalId?: string;
} | null;
requiresInput?: {
prompt: string;
schema?: unknown;
items?: unknown[];
resumeToken?: string;
approvalId?: string;
} | null;
error?: {
type?: string;
message: string;
};
};
@@ -86,14 +73,10 @@ type EmbeddedToolRuntime = {
token?: string;
approvalId?: string;
approved?: boolean;
response?: unknown;
cancel?: boolean;
ctx?: EmbeddedToolContext;
}) => Promise<EmbeddedToolEnvelope>;
};
type LoadEmbeddedToolRuntime = () => Promise<EmbeddedToolRuntime>;
const workflowExts = new Set([".lobster", ".yaml", ".yml", ".json"]);
function normalizeForCwdSandbox(p: string): string {
@@ -136,65 +119,33 @@ function createLimitedSink(maxBytes: number, label: "stdout" | "stderr") {
});
}
function normalizeEnvelope(envelope: EmbeddedToolEnvelope): LobsterEnvelope {
if (envelope.ok) {
if (envelope.status === "needs_input") {
return {
ok: false,
error: {
type: "unsupported_status",
message: "Lobster input requests are not supported by the OpenClaw Lobster tool yet",
},
};
}
return {
ok: true,
status: envelope.status ?? "ok",
output: Array.isArray(envelope.output) ? envelope.output : [],
requiresApproval: envelope.requiresApproval
? {
type: "approval_request",
prompt: envelope.requiresApproval.prompt,
items: envelope.requiresApproval.items,
...(envelope.requiresApproval.resumeToken
? { resumeToken: envelope.requiresApproval.resumeToken }
: {}),
...(envelope.requiresApproval.approvalId
? { approvalId: envelope.requiresApproval.approvalId }
: {}),
}
: null,
};
function normalizeEnvelope(envelope: EmbeddedToolEnvelope): Extract<LobsterEnvelope, { ok: true }> {
if (!envelope.ok) {
throw new Error(envelope.error?.message ?? "lobster runtime failed");
}
if (envelope.status === "needs_input") {
throw new Error("Lobster input requests are not supported by the OpenClaw Lobster tool yet");
}
return {
ok: false,
error: {
type: envelope.error?.type,
message: envelope.error?.message ?? "lobster runtime failed",
},
ok: true,
status: envelope.status ?? "ok",
output: Array.isArray(envelope.output) ? envelope.output : [],
requiresApproval: envelope.requiresApproval
? {
type: "approval_request",
prompt: envelope.requiresApproval.prompt,
items: envelope.requiresApproval.items,
...(envelope.requiresApproval.resumeToken
? { resumeToken: envelope.requiresApproval.resumeToken }
: {}),
...(envelope.requiresApproval.approvalId
? { approvalId: envelope.requiresApproval.approvalId }
: {}),
}
: null,
};
}
function throwOnErrorEnvelope(envelope: LobsterEnvelope): Extract<LobsterEnvelope, { ok: true }> {
if (envelope.ok) {
return envelope;
}
throw new Error(envelope.error.message);
}
async function resolveWorkflowFile(candidate: string, cwd: string) {
const resolved = path.isAbsolute(candidate) ? candidate : path.resolve(cwd, candidate);
const fileStat = await stat(resolved);
if (!fileStat.isFile()) {
throw new Error("Workflow path is not a file");
}
const ext = path.extname(resolved).toLowerCase();
if (!workflowExts.has(ext)) {
throw new Error("Workflow file must end in .lobster, .yaml, .yml, or .json");
}
return resolved;
}
function isMissingPathError(error: unknown) {
return (
typeof error === "object" &&
@@ -204,32 +155,25 @@ function isMissingPathError(error: unknown) {
);
}
function hasWorkflowFileExtension(candidate: string) {
return workflowExts.has(path.extname(candidate).toLowerCase());
}
async function detectWorkflowFile(candidate: string, cwd: string) {
const trimmed = candidate.trim();
if (!trimmed || trimmed.includes("|") || !hasWorkflowFileExtension(trimmed)) {
if (!trimmed || trimmed.includes("|") || !workflowExts.has(path.extname(trimmed).toLowerCase())) {
return null;
}
if (!/\s/.test(trimmed)) {
return await resolveWorkflowFile(trimmed, cwd);
}
const resolved = path.isAbsolute(trimmed) ? trimmed : path.resolve(cwd, trimmed);
try {
return await resolveWorkflowFile(trimmed, cwd);
if (!(await stat(resolved)).isFile()) {
throw new Error("Workflow path is not a file");
}
return resolved;
} catch (error) {
if (isMissingPathError(error)) {
if (/\s/.test(trimmed) && isMissingPathError(error)) {
return null;
}
throw error;
}
}
function parseWorkflowArgs(argsJson: string) {
return JSON.parse(argsJson) as Record<string, unknown>;
}
function createEmbeddedToolContext(
params: LobsterRunnerParams,
signal?: AbortSignal,
@@ -282,7 +226,7 @@ async function loadEmbeddedToolRuntimeFromPackage(): Promise<EmbeddedToolRuntime
}
export function createEmbeddedLobsterRunner(options?: {
loadRuntime?: LoadEmbeddedToolRuntime;
loadRuntime?: () => Promise<EmbeddedToolRuntime>;
}): LobsterRunner {
const loadRuntime = options?.loadRuntime ?? loadEmbeddedToolRuntimeFromPackage;
let runtimePromise: Promise<EmbeddedToolRuntime> | undefined;
@@ -305,19 +249,15 @@ export function createEmbeddedLobsterRunner(options?: {
let args: Record<string, unknown> | undefined;
if (parsedArgsJson) {
try {
args = parseWorkflowArgs(parsedArgsJson);
args = JSON.parse(parsedArgsJson) as Record<string, unknown>;
} catch {
throw new Error("run --args-json must be valid JSON");
}
}
return throwOnErrorEnvelope(
normalizeEnvelope(await runtime.runToolRequest({ filePath, args, ctx })),
);
return normalizeEnvelope(await runtime.runToolRequest({ filePath, args, ctx }));
}
return throwOnErrorEnvelope(
normalizeEnvelope(await runtime.runToolRequest({ pipeline, ctx })),
);
return normalizeEnvelope(await runtime.runToolRequest({ pipeline, ctx }));
}
const token = params.token?.trim() ?? "";
@@ -329,15 +269,13 @@ export function createEmbeddedLobsterRunner(options?: {
throw new Error("approve required");
}
return throwOnErrorEnvelope(
normalizeEnvelope(
await runtime.resumeToolRequest({
...(token ? { token } : {}),
...(approvalId ? { approvalId } : {}),
approved: params.approve,
ctx,
}),
),
return normalizeEnvelope(
await runtime.resumeToolRequest({
...(token ? { token } : {}),
...(approvalId ? { approvalId } : {}),
approved: params.approve,
ctx,
}),
);
});
},
+43 -121
View File
@@ -2,7 +2,7 @@
import type { OpenClawPluginApi } from "../runtime-api.js";
import type { LobsterEnvelope, LobsterRunner, LobsterRunnerParams } from "./lobster-runner.js";
type JsonLike =
export type JsonLike =
| null
| boolean
| number
@@ -12,7 +12,7 @@ type JsonLike =
[key: string]: JsonLike;
};
type BoundTaskFlow = ReturnType<
export type BoundTaskFlow = ReturnType<
NonNullable<OpenClawPluginApi["runtime"]>["tasks"]["managedFlows"]["bindSession"]
>;
@@ -107,23 +107,13 @@ function toJsonLike(value: unknown, seen = new WeakSet<object>()): JsonLike {
}
function buildApprovalWaitState(envelope: Extract<LobsterEnvelope, { ok: true }>): JsonLike {
if (!envelope.requiresApproval) {
return {
kind: "lobster_approval",
prompt: "",
items: [],
} satisfies LobsterApprovalWaitState;
}
const approval = envelope.requiresApproval;
return {
kind: "lobster_approval",
prompt: envelope.requiresApproval.prompt,
items: envelope.requiresApproval.items.map((item) => toJsonLike(item)),
...(envelope.requiresApproval.resumeToken
? { resumeToken: envelope.requiresApproval.resumeToken }
: {}),
...(envelope.requiresApproval.approvalId
? { approvalId: envelope.requiresApproval.approvalId }
: {}),
prompt: approval ? approval.prompt : "",
items: approval ? approval.items.map((item) => toJsonLike(item)) : [],
...(approval?.resumeToken ? { resumeToken: approval.resumeToken } : {}),
...(approval?.approvalId ? { approvalId: approval.approvalId } : {}),
} satisfies LobsterApprovalWaitState;
}
@@ -134,31 +124,52 @@ function applyEnvelopeToFlow(params: {
waitingStep: string;
}): MutationResult {
const { taskFlow, flow, envelope, waitingStep } = params;
const flowMutation = { flowId: flow.flowId, expectedRevision: flow.revision };
if (!envelope.ok) {
return taskFlow.fail({
flowId: flow.flowId,
expectedRevision: flow.revision,
});
return taskFlow.fail(flowMutation);
}
if (envelope.status === "needs_approval") {
return taskFlow.setWaiting({
flowId: flow.flowId,
expectedRevision: flow.revision,
...flowMutation,
currentStep: waitingStep,
waitJson: buildApprovalWaitState(envelope),
});
}
return taskFlow.finish({
flowId: flow.flowId,
expectedRevision: flow.revision,
});
return taskFlow.finish(flowMutation);
}
function buildEnvelopeError(envelope: Extract<LobsterEnvelope, { ok: false }>) {
return new Error(envelope.error.message);
async function executeManagedLobsterFlow(
params: Pick<RunManagedLobsterFlowParams, "taskFlow" | "runner" | "runnerParams" | "waitingStep">,
flow: FlowRecord,
failureFlowId = flow.flowId,
): Promise<ManagedLobsterFlowResult> {
try {
const envelope = await params.runner.run(params.runnerParams);
const mutation = applyEnvelopeToFlow({
taskFlow: params.taskFlow,
flow,
envelope,
waitingStep: params.waitingStep ?? "await_lobster_approval",
});
if (!envelope.ok) {
return { ok: false, flow, mutation, error: new Error(envelope.error.message) };
}
return { ok: true, envelope, flow, mutation };
} catch (error) {
const err = error instanceof Error ? error : new Error(String(error));
try {
const mutation = params.taskFlow.fail({
flowId: failureFlowId,
expectedRevision: flow.revision,
});
return { ok: false, flow, mutation, error: err };
} catch {
return { ok: false, flow, error: err };
}
}
}
export async function runManagedLobsterFlow(
@@ -174,55 +185,9 @@ export async function runManagedLobsterFlow(
? params.taskFlow.tryCreateManaged(createFlowParams)
: params.taskFlow.createManaged(createFlowParams);
if (!flow) {
return {
ok: false,
error: new Error("TaskFlow persistence failed."),
};
}
try {
const envelope = await params.runner.run(params.runnerParams);
const mutation = applyEnvelopeToFlow({
taskFlow: params.taskFlow,
flow,
envelope,
waitingStep: params.waitingStep ?? "await_lobster_approval",
});
if (!envelope.ok) {
return {
ok: false,
flow,
mutation,
error: buildEnvelopeError(envelope),
};
}
return {
ok: true,
envelope,
flow,
mutation,
};
} catch (error) {
const err = error instanceof Error ? error : new Error(String(error));
try {
const mutation = params.taskFlow.fail({
flowId: flow.flowId,
expectedRevision: flow.revision,
});
return {
ok: false,
flow,
mutation,
error: err,
};
} catch {
return {
ok: false,
flow,
error: err,
};
}
return { ok: false, error: new Error("TaskFlow persistence failed.") };
}
return await executeManagedLobsterFlow(params, flow);
}
export async function resumeManagedLobsterFlow(
@@ -242,48 +207,5 @@ export async function resumeManagedLobsterFlow(
error: new Error(`TaskFlow resume failed: ${resumed.code}`),
};
}
try {
const envelope = await params.runner.run(params.runnerParams);
const mutation = applyEnvelopeToFlow({
taskFlow: params.taskFlow,
flow: resumed.flow,
envelope,
waitingStep: params.waitingStep ?? "await_lobster_approval",
});
if (!envelope.ok) {
return {
ok: false,
flow: resumed.flow,
mutation,
error: buildEnvelopeError(envelope),
};
}
return {
ok: true,
envelope,
flow: resumed.flow,
mutation,
};
} catch (error) {
const err = error instanceof Error ? error : new Error(String(error));
try {
const mutation = params.taskFlow.fail({
flowId: params.flowId,
expectedRevision: resumed.flow.revision,
});
return {
ok: false,
flow: resumed.flow,
mutation,
error: err,
};
} catch {
return {
ok: false,
flow: resumed.flow,
error: err,
};
}
}
return await executeManagedLobsterFlow(params, resumed.flow, params.flowId);
}
+9 -37
View File
@@ -17,25 +17,13 @@ import {
type LobsterRunnerParams,
} from "./lobster-runner.js";
import {
type BoundTaskFlow,
type JsonLike,
type ManagedLobsterFlowResult,
resumeManagedLobsterFlow,
runManagedLobsterFlow,
} from "./lobster-taskflow.js";
type BoundTaskFlow = ReturnType<
NonNullable<OpenClawPluginApi["runtime"]>["tasks"]["managedFlows"]["bindSession"]
>;
type JsonLike =
| null
| boolean
| number
| string
| JsonLike[]
| {
[key: string]: JsonLike;
};
type LobsterToolOptions = {
runner?: LobsterRunner;
taskFlow?: BoundTaskFlow;
@@ -56,13 +44,6 @@ type ManagedFlowResumeParams = {
waitingStep?: string;
};
type ManagedFlowSuccessResult = {
ok: true;
envelope: unknown;
flow: unknown;
mutation: unknown;
};
function readOptionalTrimmedString(value: unknown, fieldName: string): string | undefined {
if (value === undefined) {
return undefined;
@@ -203,17 +184,15 @@ function parseResumeFlowParams(params: Record<string, unknown>): ManagedFlowResu
};
}
function formatManagedFlowResult(result: ManagedFlowSuccessResult) {
const envelope =
result.envelope && typeof result.envelope === "object" && !Array.isArray(result.envelope)
? result.envelope
: { envelope: result.envelope };
const details = {
...envelope,
function resolveManagedFlowToolResult(result: ManagedLobsterFlowResult) {
if (!result.ok) {
throw result.error;
}
return jsonResult({
...result.envelope,
flow: result.flow,
mutation: result.mutation,
};
return jsonResult(details);
});
}
function requireTaskFlowRuntime(taskFlow: BoundTaskFlow | undefined, action: "run" | "resume") {
@@ -223,13 +202,6 @@ function requireTaskFlowRuntime(taskFlow: BoundTaskFlow | undefined, action: "ru
return taskFlow;
}
function resolveManagedFlowToolResult(result: ManagedLobsterFlowResult) {
if (!result.ok) {
throw result.error;
}
return formatManagedFlowResult(result);
}
export function createLobsterTool(api: OpenClawPluginApi, options?: LobsterToolOptions) {
const runner = options?.runner ?? createEmbeddedLobsterRunner();
return {