From 08b9c4c3e140566cec6842062b8359968a98798d Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 5 Aug 2026 07:22:43 -0700 Subject: [PATCH] refactor(lobster): consolidate embedded workflow execution (#119630) --- extensions/lobster/src/lobster-runner.ts | 142 +++++------------- extensions/lobster/src/lobster-taskflow.ts | 164 ++++++--------------- extensions/lobster/src/lobster-tool.ts | 46 ++---- 3 files changed, 92 insertions(+), 260 deletions(-) diff --git a/extensions/lobster/src/lobster-runner.ts b/extensions/lobster/src/lobster-runner.ts index f0d93604518e..15a4808927ec 100644 --- a/extensions/lobster/src/lobster-runner.ts +++ b/extensions/lobster/src/lobster-runner.ts @@ -45,32 +45,19 @@ type EmbeddedToolContext = { stdout?: NodeJS.WritableStream; stderr?: NodeJS.WritableStream; signal?: AbortSignal; - registry?: unknown; - llmAdapters?: Record; }; 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; }; -type LoadEmbeddedToolRuntime = () => Promise; - 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 { + 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 { - 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; -} - function createEmbeddedToolContext( params: LobsterRunnerParams, signal?: AbortSignal, @@ -282,7 +226,7 @@ async function loadEmbeddedToolRuntimeFromPackage(): Promise Promise; }): LobsterRunner { const loadRuntime = options?.loadRuntime ?? loadEmbeddedToolRuntimeFromPackage; let runtimePromise: Promise | undefined; @@ -305,19 +249,15 @@ export function createEmbeddedLobsterRunner(options?: { let args: Record | undefined; if (parsedArgsJson) { try { - args = parseWorkflowArgs(parsedArgsJson); + args = JSON.parse(parsedArgsJson) as Record; } 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, + }), ); }); }, diff --git a/extensions/lobster/src/lobster-taskflow.ts b/extensions/lobster/src/lobster-taskflow.ts index 70a5599ae497..813eaadf4d4d 100644 --- a/extensions/lobster/src/lobster-taskflow.ts +++ b/extensions/lobster/src/lobster-taskflow.ts @@ -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["tasks"]["managedFlows"]["bindSession"] >; @@ -107,23 +107,13 @@ function toJsonLike(value: unknown, seen = new WeakSet()): JsonLike { } function buildApprovalWaitState(envelope: Extract): 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) { - return new Error(envelope.error.message); +async function executeManagedLobsterFlow( + params: Pick, + flow: FlowRecord, + failureFlowId = flow.flowId, +): Promise { + 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); } diff --git a/extensions/lobster/src/lobster-tool.ts b/extensions/lobster/src/lobster-tool.ts index ade472496d41..0af418e2e3e9 100644 --- a/extensions/lobster/src/lobster-tool.ts +++ b/extensions/lobster/src/lobster-tool.ts @@ -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["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): 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 {