From 2eb9c7ebd72732933deae58f44f43c2771fc741e Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Tue, 14 Jul 2026 14:15:44 -0700 Subject: [PATCH] refactor(extensions): privatize small plugin internals (#107774) * refactor(llama-cpp): privatize embedding internals * refactor(parallel): privatize MCP response helpers * refactor(logbook): privatize analysis parsers * refactor(crabbox): narrow worker provider exports * refactor(acpx): privatize process reaper internals * chore(deadcode): refresh unused-export baseline --- extensions/acpx/src/process-reaper.test.ts | 44 +--- extensions/acpx/src/process-reaper.ts | 4 +- .../src/crabbox-worker-provider.test.ts | 13 +- .../crabbox/src/crabbox-worker-provider.ts | 4 +- extensions/llama-cpp/index.test.ts | 129 ++++++------ .../llama-cpp/src/embedding-provider.ts | 6 +- extensions/logbook/src/analyze.test.ts | 97 ++++----- extensions/logbook/src/analyze.ts | 4 +- .../src/parallel-mcp-search.runtime.test.ts | 188 ++++++------------ .../src/parallel-mcp-search.runtime.ts | 6 +- scripts/deadcode-exports.baseline.mjs | 12 -- 11 files changed, 200 insertions(+), 307 deletions(-) diff --git a/extensions/acpx/src/process-reaper.test.ts b/extensions/acpx/src/process-reaper.test.ts index 1df6e90f49b..a022cfb8d29 100644 --- a/extensions/acpx/src/process-reaper.test.ts +++ b/extensions/acpx/src/process-reaper.test.ts @@ -5,9 +5,7 @@ import { OPENCLAW_ACPX_LEASE_ID_ARG, OPENCLAW_GATEWAY_INSTANCE_ID_ARG } from "./ import { cleanupOpenClawOwnedAcpxProcessTree, isOpenClawLeaseAwareAcpxProcessCommand, - isOpenClawOwnedAcpxProcessCommand, reapStaleOpenClawOwnedAcpxOrphans, - type AcpxProcessInfo, } from "./process-reaper.js"; const WRAPPER_ROOT = "/tmp/openclaw-state/acpx"; @@ -23,6 +21,9 @@ const LOCAL_NODE_MODULES_CODEX_PLATFORM_COMMAND = path.resolve( "node_modules/@zed-industries/codex-acp-linux-x64/bin/codex-acp", ); +type CleanupDeps = NonNullable[0]["deps"]>; +type AcpxProcessInfo = Awaited>>[number]; + function cleanupDeps(processes: AcpxProcessInfo[]) { const killed: Array<{ pid: number; signal: NodeJS.Signals }> = []; return { @@ -52,27 +53,6 @@ function collectMatching( } describe("process reaper", () => { - it("recognizes generated Codex and Claude wrappers only under the configured root", () => { - expect( - isOpenClawOwnedAcpxProcessCommand({ - command: CODEX_WRAPPER_COMMAND, - wrapperRoot: WRAPPER_ROOT, - }), - ).toBe(true); - expect( - isOpenClawOwnedAcpxProcessCommand({ - command: CLAUDE_WRAPPER_COMMAND, - wrapperRoot: WRAPPER_ROOT, - }), - ).toBe(true); - expect( - isOpenClawOwnedAcpxProcessCommand({ - command: "node /tmp/other/codex-acp-wrapper.mjs", - wrapperRoot: WRAPPER_ROOT, - }), - ).toBe(false); - }); - it("only treats generated wrappers as launch-lease aware", () => { expect( isOpenClawLeaseAwareAcpxProcessCommand({ @@ -88,24 +68,6 @@ describe("process reaper", () => { ); }); - it("recognizes OpenClaw plugin-runtime-deps ACP adapter children", () => { - expect(isOpenClawOwnedAcpxProcessCommand({ command: PLUGIN_DEPS_CODEX_COMMAND })).toBe(true); - expect(isOpenClawOwnedAcpxProcessCommand({ command: "npx @zed-industries/codex-acp" })).toBe( - false, - ); - }); - - it("recognizes plugin-local ACP adapter package paths without trusting arbitrary installs", () => { - expect(isOpenClawOwnedAcpxProcessCommand({ command: LOCAL_NODE_MODULES_CODEX_COMMAND })).toBe( - true, - ); - expect( - isOpenClawOwnedAcpxProcessCommand({ - command: "node /tmp/other-project/node_modules/@zed-industries/codex-acp/bin/codex-acp.js", - }), - ).toBe(false); - }); - it("kills an owned recorded process tree children first", async () => { const { deps, killed } = cleanupDeps([ { pid: 100, ppid: 1, command: CODEX_WRAPPER_COMMAND }, diff --git a/extensions/acpx/src/process-reaper.ts b/extensions/acpx/src/process-reaper.ts index 6e80fa53870..d5cd55279e9 100644 --- a/extensions/acpx/src/process-reaper.ts +++ b/extensions/acpx/src/process-reaper.ts @@ -32,7 +32,7 @@ const ACP_PACKAGE_MARKERS = [ ]; /** Minimal process-table row used by ACPX cleanup. */ -export type AcpxProcessInfo = { +type AcpxProcessInfo = { pid: number; ppid: number; command: string; @@ -166,7 +166,7 @@ function liveCommandMatchesLeaseIdentity(params: { } /** Check whether a command is owned by OpenClaw ACPX runtime packages or wrappers. */ -export function isOpenClawOwnedAcpxProcessCommand(params: { +function isOpenClawOwnedAcpxProcessCommand(params: { command: string | undefined; wrapperRoot?: string; }): boolean { diff --git a/extensions/crabbox/src/crabbox-worker-provider.test.ts b/extensions/crabbox/src/crabbox-worker-provider.test.ts index a4cd535b7b9..c55819d5d58 100644 --- a/extensions/crabbox/src/crabbox-worker-provider.test.ts +++ b/extensions/crabbox/src/crabbox-worker-provider.test.ts @@ -2,12 +2,8 @@ import path from "node:path"; import type { WorkerProfile } from "openclaw/plugin-sdk/plugin-entry"; import type { SpawnResult } from "openclaw/plugin-sdk/process-runtime"; import { describe, expect, it, vi } from "vitest"; -import { - createCrabboxWorkerProvider, - type CrabboxCommandRunner, - resolveCrabboxBinary, - resolveOpenClawRoot, -} from "./crabbox-worker-provider.js"; +import { resolveCrabboxBinary } from "./crabbox-worker-profile.js"; +import { createCrabboxWorkerProvider, resolveOpenClawRoot } from "./crabbox-worker-provider.js"; const LEASE_ID = "cbx_012345abcdef"; const FALLBACK_LEASE_ID = "cbx_20260711123456123456"; @@ -25,6 +21,11 @@ const PROFILE = { idleTimeout: "60m", }; +type CrabboxWorkerProviderDependencies = NonNullable< + Parameters[0] +>; +type CrabboxCommandRunner = NonNullable; + function commandResult(overrides: Partial = {}): SpawnResult { return { stdout: "", diff --git a/extensions/crabbox/src/crabbox-worker-provider.ts b/extensions/crabbox/src/crabbox-worker-provider.ts index ce1373747f1..39dc5725f16 100644 --- a/extensions/crabbox/src/crabbox-worker-provider.ts +++ b/extensions/crabbox/src/crabbox-worker-provider.ts @@ -18,7 +18,7 @@ import { resolveCrabboxBinary, } from "./crabbox-worker-profile.js"; -export { resolveCrabboxBinary, resolveOpenClawRoot } from "./crabbox-worker-profile.js"; +export { resolveOpenClawRoot } from "./crabbox-worker-profile.js"; const CRABBOX_WORKER_PROVIDER_ID = "crabbox"; const CRABBOX_KEY_REF_PROVIDER = "crabbox"; @@ -52,7 +52,7 @@ const UNUSABLE_PROVISION_STATES = new Set([...DESTROYED_STATES, "deleting", "fai const LEASE_ID_PATTERN = /^(?:cbx_|tbx_)[A-Za-z0-9][A-Za-z0-9_-]{0,127}$/u; const LEASE_TOKEN_IN_OUTPUT_PATTERN = /^leased\s+(\S{1,128})(?=\s|$)/mu; -export type CrabboxCommandRunner = typeof runCommandWithTimeout; +type CrabboxCommandRunner = typeof runCommandWithTimeout; type LeaseCommandContext = { binary: string; diff --git a/extensions/llama-cpp/index.test.ts b/extensions/llama-cpp/index.test.ts index bcfa5692eac..a88c5f3b697 100644 --- a/extensions/llama-cpp/index.test.ts +++ b/extensions/llama-cpp/index.test.ts @@ -22,12 +22,23 @@ vi.mock("openclaw/plugin-sdk/memory-core-host-engine-embeddings", () => ({ })); import llamaCppPlugin from "./index.js"; -import { - DEFAULT_LLAMA_CPP_EMBEDDING_MODEL, - createLlamaCppMemoryEmbeddingProvider, - formatLlamaCppSetupError, - llamaCppEmbeddingProviderAdapter, -} from "./src/embedding-provider.js"; +import { llamaCppEmbeddingProviderAdapter } from "./src/embedding-provider.js"; + +const DEFAULT_LLAMA_CPP_EMBEDDING_MODEL = + "hf:ggml-org/embeddinggemma-300m-qat-q8_0-GGUF/embeddinggemma-300m-qat-Q8_0.gguf"; +type AdapterCreateOptions = Parameters[0]; +type MemoryCreateTestOptions = AdapterCreateOptions & { + fallback?: "none"; + outputDimensionality?: number; +}; + +async function createLlamaCppMemoryEmbeddingProvider(options: MemoryCreateTestOptions) { + const { fallback: _fallback, outputDimensionality, ...adapterOptions } = options; + return await llamaCppEmbeddingProviderAdapter.create({ + ...adapterOptions, + dimensions: outputDimensionality, + }); +} afterEach(() => { clearEmbeddingProviders(); @@ -148,16 +159,13 @@ describe("llama.cpp provider plugin", () => { embedBatch: vi.fn(), }); - const result = await createLlamaCppMemoryEmbeddingProvider( - { - config: {}, - provider: "local", - fallback: "none", - model: DEFAULT_LLAMA_CPP_EMBEDDING_MODEL, - outputDimensionality: 512, - }, - { nodeLlamaCppImportUrl: "file:///plugin/node-llama-cpp.js" }, - ); + const result = await createLlamaCppMemoryEmbeddingProvider({ + config: {}, + provider: "local", + fallback: "none", + model: DEFAULT_LLAMA_CPP_EMBEDDING_MODEL, + outputDimensionality: 512, + }); const resolvedIdentity = llamaCppEmbeddingProviderAdapter.resolveIndexIdentity?.({ config: {}, provider: "local", @@ -197,16 +205,13 @@ describe("llama.cpp provider plugin", () => { embedBatch: vi.fn(), }); - const result = await createLlamaCppMemoryEmbeddingProvider( - { - config: {}, - provider: "local", - fallback: "none", - model: modelPath, - local: { modelPath }, - }, - { nodeLlamaCppImportUrl: "file:///plugin/node-llama-cpp.js" }, - ); + const result = await createLlamaCppMemoryEmbeddingProvider({ + config: {}, + provider: "local", + fallback: "none", + model: modelPath, + local: { modelPath }, + }); expect(result.provider?.model).toBe(DEFAULT_LLAMA_CPP_EMBEDDING_MODEL); expect(result.runtime?.cacheKeyData).toEqual({ @@ -265,7 +270,7 @@ describe("llama.cpp provider plugin", () => { local: { modelPath }, }), { - nodeLlamaCppImportUrl: "file:///plugin/node-llama-cpp.js", + nodeLlamaCppImportUrl: expect.stringContaining("node-llama-cpp"), }, ); }); @@ -283,16 +288,13 @@ describe("llama.cpp provider plugin", () => { embedBatch: vi.fn(), }); - const result = await createLlamaCppMemoryEmbeddingProvider( - { - config: {}, - provider: "local", - fallback: "none", - model: modelPath, - local: { modelPath }, - }, - { nodeLlamaCppImportUrl: "file:///plugin/node-llama-cpp.js" }, - ); + const result = await createLlamaCppMemoryEmbeddingProvider({ + config: {}, + provider: "local", + fallback: "none", + model: modelPath, + local: { modelPath }, + }); expect(result.provider?.model).toBe(modelPath); expect(result.runtime?.cacheKeyData).toEqual({ @@ -316,16 +318,13 @@ describe("llama.cpp provider plugin", () => { embedBatch: vi.fn(), }); - const result = await createLlamaCppMemoryEmbeddingProvider( - { - config: {}, - provider: "local", - fallback: "none", - model: modelPath, - local: { modelPath }, - }, - { nodeLlamaCppImportUrl: "file:///plugin/node-llama-cpp.js" }, - ); + const result = await createLlamaCppMemoryEmbeddingProvider({ + config: {}, + provider: "local", + fallback: "none", + model: modelPath, + local: { modelPath }, + }); expect(result.provider?.model).toBe(modelPath); expect(result.runtime).not.toHaveProperty("indexIdentityAliases"); @@ -341,16 +340,13 @@ describe("llama.cpp provider plugin", () => { embedBatch: vi.fn(), }); - const result = await createLlamaCppMemoryEmbeddingProvider( - { - config: {}, - provider: "local", - fallback: "none", - model: DEFAULT_LLAMA_CPP_EMBEDDING_MODEL, - local: { modelPath: DEFAULT_LLAMA_CPP_EMBEDDING_MODEL, modelCacheDir }, - }, - { nodeLlamaCppImportUrl: "file:///plugin/node-llama-cpp.js" }, - ); + const result = await createLlamaCppMemoryEmbeddingProvider({ + config: {}, + provider: "local", + fallback: "none", + model: DEFAULT_LLAMA_CPP_EMBEDDING_MODEL, + local: { modelPath: DEFAULT_LLAMA_CPP_EMBEDDING_MODEL, modelCacheDir }, + }); expect(result.provider?.model).toBe(DEFAULT_LLAMA_CPP_EMBEDDING_MODEL); expect(result.runtime?.cacheKeyData).toEqual({ @@ -432,16 +428,13 @@ describe("llama.cpp provider plugin", () => { embedBatch: vi.fn(), }); - const result = await createLlamaCppMemoryEmbeddingProvider( - { - config: {}, - provider: "local", - fallback: "none", - model: modelPath, - local: { modelPath, modelCacheDir }, - }, - { nodeLlamaCppImportUrl: "file:///plugin/node-llama-cpp.js" }, - ); + const result = await createLlamaCppMemoryEmbeddingProvider({ + config: {}, + provider: "local", + fallback: "none", + model: modelPath, + local: { modelPath, modelCacheDir }, + }); expect(result.provider?.model).toBe(DEFAULT_LLAMA_CPP_EMBEDDING_MODEL); expect(result.runtime?.indexIdentityAliases).toEqual([ @@ -467,7 +460,7 @@ describe("llama.cpp provider plugin", () => { code: "ERR_MODULE_NOT_FOUND", }); - expect(formatLlamaCppSetupError(err)).toContain( + expect(llamaCppEmbeddingProviderAdapter.formatSetupError?.(err)).toContain( "openclaw plugins install @openclaw/llama-cpp-provider", ); }); diff --git a/extensions/llama-cpp/src/embedding-provider.ts b/extensions/llama-cpp/src/embedding-provider.ts index 800fb3e1b70..2ac362ea31d 100644 --- a/extensions/llama-cpp/src/embedding-provider.ts +++ b/extensions/llama-cpp/src/embedding-provider.ts @@ -29,7 +29,7 @@ type LlamaCppEmbeddingProviderRuntimeOptions = { const LLAMA_CPP_EMBEDDING_PROVIDER_ID = "local"; const LOCAL_EMBEDDING_RUNTIME_FACTS = Symbol.for("openclaw.localEmbeddingRuntimeFacts"); -export const DEFAULT_LLAMA_CPP_EMBEDDING_MODEL = +const DEFAULT_LLAMA_CPP_EMBEDDING_MODEL = "hf:ggml-org/embeddinggemma-300m-qat-q8_0-GGUF/embeddinggemma-300m-qat-Q8_0.gguf"; const DEFAULT_LLAMA_CPP_EMBEDDING_MODEL_CACHE_FILE_NAME = "hf_ggml-org_embeddinggemma-300m-qat-Q8_0.gguf"; @@ -133,7 +133,7 @@ function formatErrorMessage(err: unknown): string { return String(err); } -export function formatLlamaCppSetupError(err: unknown): string { +function formatLlamaCppSetupError(err: unknown): string { const detail = formatErrorMessage(err); const missing = isNodeLlamaCppMissing(err); return [ @@ -195,7 +195,7 @@ function adaptMemoryEmbeddingProvider(provider: MemoryEmbeddingProvider): Embedd return adapted; } -export async function createLlamaCppMemoryEmbeddingProvider( +async function createLlamaCppMemoryEmbeddingProvider( options: MemoryEmbeddingProviderCreateOptions, runtimeOptions: LlamaCppEmbeddingProviderRuntimeOptions = {}, ): Promise { diff --git a/extensions/logbook/src/analyze.test.ts b/extensions/logbook/src/analyze.test.ts index c621e46a145..ad876325f04 100644 --- a/extensions/logbook/src/analyze.test.ts +++ b/extensions/logbook/src/analyze.test.ts @@ -2,8 +2,6 @@ import { execFileSync } from "node:child_process"; import { expectDefined } from "@openclaw/normalization-core"; import { describe, expect, it } from "vitest"; import { - clockToMs, - extractJsonPayload, parseCardsJson, parseObservationSegments, pickKeyframeId, @@ -15,52 +13,13 @@ import { const DAY = "2026-07-03"; const dayMs = (clock: string) => { - const ms = clockToMs(DAY, clock); - if (ms === null) { + const ms = new Date(`${DAY}T${clock}`).getTime(); + if (!Number.isFinite(ms)) { throw new Error(`bad clock ${clock}`); } return ms; }; -describe("clockToMs", () => { - it("parses 24h and 12h clocks on the local day", () => { - expect(dayMs("13:05:30") - dayMs("13:05:00")).toBe(30_000); - expect(dayMs("1:05 pm")).toBe(dayMs("13:05:00")); - expect(dayMs("12:00 am")).toBe(dayMs("00:00:00")); - }); - - it("rejects malformed input", () => { - expect(clockToMs(DAY, "25:00:00")).toBeNull(); - expect(clockToMs(DAY, "13:05 pm")).toBeNull(); - expect(clockToMs(DAY, "00:05 am")).toBeNull(); - expect(clockToMs(DAY, "half past nine")).toBeNull(); - expect(clockToMs("not-a-day", "10:00:00")).toBeNull(); - }); - - it("preserves local wall-clock time across a DST transition", () => { - const moduleUrl = new URL("./analyze.ts", import.meta.url).href; - const output = execFileSync( - process.execPath, - [ - "--import", - "tsx", - "--eval", - `const { clockToMs } = await import(${JSON.stringify(moduleUrl)}); process.stdout.write(JSON.stringify([clockToMs("2026-03-08", "10:00:00"), clockToMs("2026-03-08", "02:30:00")]));`, - ], - { encoding: "utf8", env: { ...process.env, TZ: "America/New_York" } }, - ); - - expect(JSON.parse(output)).toEqual([Date.UTC(2026, 2, 8, 14), null]); - }); -}); - -describe("extractJsonPayload", () => { - it("strips fences and surrounding prose", () => { - expect(extractJsonPayload('```json\n{"a":1}\n```')).toBe('{"a":1}'); - expect(extractJsonPayload('Here you go:\n[{"a":1}]\nHope that helps!')).toBe('[{"a":1}]'); - }); -}); - describe("parseObservationSegments", () => { const startMs = dayMs("10:00:00"); const endMs = dayMs("10:15:00"); @@ -81,6 +40,54 @@ describe("parseObservationSegments", () => { it("returns empty on unparseable output", () => { expect(parseObservationSegments({ raw: "no json here", day: DAY, startMs, endMs })).toEqual([]); }); + + it("parses 12h clocks from fenced model output", () => { + const raw = [ + "Here you go:", + "```json", + JSON.stringify({ + segments: [{ start: "1:05 pm", end: "1:05:30 pm", description: "Reviewing the timeline" }], + }), + "```", + "Hope that helps!", + ].join("\n"); + const segments = parseObservationSegments({ + raw, + day: DAY, + startMs: dayMs("13:00:00"), + endMs: dayMs("13:10:00"), + }); + expect(segments).toHaveLength(1); + const segment = expectDefined(segments[0], "12h observation segment"); + expect(segment.endMs - segment.startMs).toBe(30_000); + }); + + it.each([ + [DAY, "25:00:00"], + [DAY, "13:05 pm"], + [DAY, "00:05 am"], + [DAY, "half past nine"], + ["not-a-day", "10:00:00"], + ])("rejects malformed clock %s %s", (day, clock) => { + const raw = JSON.stringify([{ start: clock, end: "13:06:00", description: "Invalid clock" }]); + expect(parseObservationSegments({ raw, day, startMs, endMs })).toEqual([]); + }); + + it("preserves local wall-clock time across a DST transition", () => { + const moduleUrl = new URL("./analyze.ts", import.meta.url).href; + const output = execFileSync( + process.execPath, + [ + "--import", + "tsx", + "--eval", + `const { parseObservationSegments } = await import(${JSON.stringify(moduleUrl)}); const parse = (start, end) => parseObservationSegments({ raw: JSON.stringify([{ start, end, description: "x" }]), day: "2026-03-08", startMs: 0, endMs: Number.MAX_SAFE_INTEGER }); process.stdout.write(JSON.stringify([parse("10:00:00", "10:00:01")[0]?.startMs ?? null, parse("02:30:00", "03:30:00").length]));`, + ], + { encoding: "utf8", env: { ...process.env, TZ: "America/New_York" } }, + ); + + expect(JSON.parse(output)).toEqual([Date.UTC(2026, 2, 8, 14), 0]); + }); }); describe("parseCardsJson", () => { @@ -101,9 +108,9 @@ describe("parseCardsJson", () => { it("accepts a valid card array and normalizes fields", () => { const result = parseCardsJson({ - raw: JSON.stringify([ + raw: `Here you go:\n${JSON.stringify([ card({ category: "CODING", appSites: { primary: "https://GitHub.com/openclaw" } }), - ]), + ])}\nHope that helps!`, day: DAY, windowStartMs, windowEndMs, diff --git a/extensions/logbook/src/analyze.ts b/extensions/logbook/src/analyze.ts index 6b04d43937b..c842c5a240b 100644 --- a/extensions/logbook/src/analyze.ts +++ b/extensions/logbook/src/analyze.ts @@ -15,7 +15,7 @@ export const MAX_FRAMES_PER_CALL = 16; type ParsedSegment = { startMs: number; endMs: number; text: string }; /** Parses "HH:MM:SS" (or "H:MM", with optional am/pm) on a local day into epoch ms. */ -export function clockToMs(day: string, clock: string): number | null { +function clockToMs(day: string, clock: string): number | null { const match = /^\s*(\d{1,2}):(\d{2})(?::(\d{2}))?\s*(am|pm)?\s*$/i.exec(clock); if (!match) { return null; @@ -60,7 +60,7 @@ export function clockToMs(day: string, clock: string): number | null { } /** Strips code fences and extracts the outermost JSON array/object from model text. */ -export function extractJsonPayload(raw: string): string { +function extractJsonPayload(raw: string): string { const cleaned = raw.replaceAll("```json", "").replaceAll("```", "").trim(); const firstBracket = cleaned.search(/[[{]/); if (firstBracket < 0) { diff --git a/extensions/parallel/src/parallel-mcp-search.runtime.test.ts b/extensions/parallel/src/parallel-mcp-search.runtime.test.ts index 9cf45026431..79ec326d8ee 100644 --- a/extensions/parallel/src/parallel-mcp-search.runtime.test.ts +++ b/extensions/parallel/src/parallel-mcp-search.runtime.test.ts @@ -32,12 +32,7 @@ vi.mock("openclaw/plugin-sdk/provider-web-search", async (importOriginal) => { }; }); -import { - extractMcpToolPayload, - iterMcpMessages, - runParallelMcpSearch, - selectMcpEnvelope, -} from "./parallel-mcp-search.runtime.js"; +import { runParallelMcpSearch } from "./parallel-mcp-search.runtime.js"; function jsonResponse(body: unknown, headers?: Record): Response { return new Response(JSON.stringify(body), { @@ -46,6 +41,13 @@ function jsonResponse(body: unknown, headers?: Record): Response }); } +function rawResponse(body: string, contentType: string): Response { + return new Response(body, { + status: 200, + headers: { "Content-Type": contentType }, + }); +} + function cancelTrackedResponse( text: string, init: ResponseInit, @@ -83,129 +85,69 @@ function requireEndpointCall(index: number): EndpointCall { return expectDefined(endpointMockState.calls[index], `Parallel MCP endpoint call ${index}`); } -function boundaryJsonPayload(base: Record): { - payload: Record; - truncatedJson: string; -} { - const empty = { ...base, detail: "" }; - const jsonPrefix = JSON.stringify(empty).slice(0, -2); - const detailPrefix = "x".repeat(499 - jsonPrefix.length); - return { - payload: { ...base, detail: `${detailPrefix}😀tail` }, - truncatedJson: `${jsonPrefix}${detailPrefix}`, - }; -} - -function thrownMessage(run: () => unknown): string { - try { - run(); - } catch (error) { - return error instanceof Error ? error.message : String(error); - } - throw new Error("Expected call to throw."); -} - -describe("iterMcpMessages", () => { - it("parses a single JSON object body", () => { - expect(iterMcpMessages('{"id":"a","result":{}}')).toEqual([{ id: "a", result: {} }]); - }); - - it("flattens a JSON array batch", () => { - expect(iterMcpMessages('[{"id":"a"},{"id":"b"}]')).toEqual([{ id: "a" }, { id: "b" }]); - }); - - it("parses SSE events with concatenated data lines", () => { - const sse = [ - "event: message", - 'data: {"id":"a",', - 'data: "result":{}}', - "", - 'data: {"id":"b"}', - "", - ].join("\n"); - expect(iterMcpMessages(sse)).toEqual([{ id: "a", result: {} }, { id: "b" }]); - }); - - it("skips unparseable chunks and empty bodies", () => { - expect(iterMcpMessages("")).toEqual([]); - expect(iterMcpMessages("not json")).toEqual([]); - expect(iterMcpMessages("data: {bad json}\n\n")).toEqual([]); - }); -}); - -describe("selectMcpEnvelope", () => { - it("returns the message whose id matches, skipping notifications", () => { - const body = [ - '{"jsonrpc":"2.0","method":"notifications/progress"}', - '{"jsonrpc":"2.0","id":"other","result":{"n":1}}', - '{"jsonrpc":"2.0","id":"want","result":{"n":2}}', - ] - .map((line) => `data: ${line}`) - .join("\n\n"); - expect(selectMcpEnvelope(body, "want")).toMatchObject({ id: "want", result: { n: 2 } }); - }); - - it("falls back to the last result-bearing message when no id matches", () => { - const body = '{"id":"x","result":{"first":true}}\n'; - expect(selectMcpEnvelope(body, "missing")).toMatchObject({ result: { first: true } }); - }); - - it("returns {} when there is no result or error message", () => { - expect(selectMcpEnvelope('{"method":"notifications/initialized"}', "any")).toEqual({}); - }); -}); - -describe("extractMcpToolPayload", () => { - it("prefers structuredContent", () => { - expect(extractMcpToolPayload({ result: { structuredContent: { results: [1] } } })).toEqual({ - results: [1], - }); - }); - - it("parses the first JSON-parseable text block", () => { - expect( - extractMcpToolPayload({ - result: { - content: [ - { type: "text", text: "not json" }, - { type: "text", text: '{"ok":true}' }, - ], - }, - }), - ).toEqual({ ok: true }); - }); - - it("throws on a JSON-RPC error", () => { - const { payload, truncatedJson } = boundaryJsonPayload({ code: -1, message: "boom" }); - - expect(thrownMessage(() => extractMcpToolPayload({ error: payload }))).toBe( - `Parallel MCP error: ${truncatedJson}`, - ); - }); - - it("throws on a tool-level isError", () => { - const { payload, truncatedJson } = boundaryJsonPayload({ isError: true }); - - expect(thrownMessage(() => extractMcpToolPayload({ result: payload }))).toBe( - `Parallel MCP tool error: ${truncatedJson}`, - ); - }); - - it("throws when there is no parseable content", () => { - const { payload, truncatedJson } = boundaryJsonPayload({ content: [] }); - - expect(thrownMessage(() => extractMcpToolPayload({ result: payload }))).toBe( - `Parallel MCP returned no parseable content: ${truncatedJson}`, - ); - }); -}); - describe("runParallelMcpSearch", () => { beforeEach(() => { endpointMockState.calls = []; endpointMockState.responses = []; }); + it("handles SSE notifications, multiline events, JSON batches, and structured payloads", async () => { + endpointMockState.responses.push( + rawResponse( + [ + 'data: {"jsonrpc":"2.0","method":"notifications/progress"}', + "", + 'data: {"jsonrpc":"2.0","id":"ignored",', + 'data: "result":{"protocolVersion":"2025-06-18"}}', + "", + ].join("\n"), + "text/event-stream", + ), + jsonResponse({ jsonrpc: "2.0" }), + jsonResponse([ + { jsonrpc: "2.0", method: "notifications/progress" }, + { + jsonrpc: "2.0", + id: "ignored", + result: { + structuredContent: { + search_id: "search_sse", + results: [{ url: "https://example.com", title: "Example", excerpts: ["hi"] }], + }, + }, + }, + ]), + ); + + await expect( + runParallelMcpSearch({ searchQueries: ["test"], maxResults: 5 }), + ).resolves.toMatchObject({ + search_id: "search_sse", + results: [{ url: "https://example.com", title: "Example" }], + }); + }); + + it.each([ + [{ error: { code: -1, message: "boom" } }, "Parallel MCP error"], + [{ result: { isError: true } }, "Parallel MCP tool error"], + [{ result: { content: [] } }, "Parallel MCP returned no parseable content"], + ])("surfaces bounded tool-envelope failures", async (envelope, expectedPrefix) => { + const detail = `${"x".repeat(600)}😀tail`; + const detailedEnvelope = + "error" in envelope + ? { error: { ...envelope.error, detail } } + : { result: { ...envelope.result, detail } }; + endpointMockState.responses.push( + jsonResponse({ result: { protocolVersion: "2025-06-18" } }), + jsonResponse({}), + jsonResponse(detailedEnvelope), + ); + + await expect(runParallelMcpSearch({ searchQueries: ["test"], maxResults: 5 })).rejects.toThrow( + expectedPrefix, + ); + }); + it("runs the 3-step handshake and maps results into the REST-compatible shape", async () => { endpointMockState.responses.push( jsonResponse( diff --git a/extensions/parallel/src/parallel-mcp-search.runtime.ts b/extensions/parallel/src/parallel-mcp-search.runtime.ts index c26e850f090..6c2f6dd66e9 100644 --- a/extensions/parallel/src/parallel-mcp-search.runtime.ts +++ b/extensions/parallel/src/parallel-mcp-search.runtime.ts @@ -72,7 +72,7 @@ function mcpHeaders(params: { * responses into a JSON array, so arrays are flattened. Unparseable chunks and * non-`data` SSE fields (`event:`/`id:`/comments) are skipped. */ -export function iterMcpMessages(text: string): JsonRpcMessage[] { +function iterMcpMessages(text: string): JsonRpcMessage[] { const out: JsonRpcMessage[] = []; const emit = (payload: unknown): void => { if (Array.isArray(payload)) { @@ -132,7 +132,7 @@ export function iterMcpMessages(text: string): JsonRpcMessage[] { * `id` matches. Falls back to the last result/error-bearing message if no id * matches; `{}` if none is present. */ -export function selectMcpEnvelope(text: string, requestId: string): JsonRpcMessage { +function selectMcpEnvelope(text: string, requestId: string): JsonRpcMessage { let fallback: JsonRpcMessage = {}; for (const msg of iterMcpMessages(text)) { if (!("result" in msg || "error" in msg)) { @@ -153,7 +153,7 @@ export function selectMcpEnvelope(text: string, requestId: string): JsonRpcMessa * scans text blocks for the first JSON-parseable one. Throws on a JSON-RPC * error or a tool-level `isError`. */ -export function extractMcpToolPayload(envelope: JsonRpcMessage): McpToolPayload { +function extractMcpToolPayload(envelope: JsonRpcMessage): McpToolPayload { if ("error" in envelope) { throw new Error( `Parallel MCP error: ${truncateUtf16Safe(JSON.stringify(envelope.error), 500)}`, diff --git a/scripts/deadcode-exports.baseline.mjs b/scripts/deadcode-exports.baseline.mjs index 565994254ec..33870af70ca 100644 --- a/scripts/deadcode-exports.baseline.mjs +++ b/scripts/deadcode-exports.baseline.mjs @@ -2,8 +2,6 @@ // New entries fail CI. After deleting dead code, run `pnpm deadcode:exports:update`. // Do not add entries to avoid fixing new findings. export const KNIP_UNUSED_EXPORT_BASELINE = [ - "extensions/acpx/src/process-reaper.ts: AcpxProcessInfo", - "extensions/acpx/src/process-reaper.ts: isOpenClawOwnedAcpxProcessCommand", "extensions/canvas/src/host/a2ui.ts: createA2uiHttpRequestHandler", "extensions/canvas/src/host/a2ui.ts: injectCanvasRuntime", "extensions/canvas/src/host/a2ui.ts: isA2uiPath", @@ -11,8 +9,6 @@ export const KNIP_UNUSED_EXPORT_BASELINE = [ "extensions/codex/src/session-catalog.ts: CODEX_TERMINAL_RESUME_COMMAND", "extensions/codex/src/session-upstream-activity.ts: checkCodexUpstreamActivity (upstream)", "extensions/codex/src/session-upstream-activity.ts: classifyCodexUpstreamTurns", - "extensions/crabbox/src/crabbox-worker-provider.ts: CrabboxCommandRunner", - "extensions/crabbox/src/crabbox-worker-provider.ts: resolveCrabboxBinary", "extensions/diagnostics-prometheus/src/service.ts: testApi", "extensions/diffs/src/browser.ts: resetSharedBrowserStateForTests", "extensions/diffs/src/shiki-curated-languages.ts: bundledLanguages", @@ -42,12 +38,7 @@ export const KNIP_UNUSED_EXPORT_BASELINE = [ "extensions/googlechat/src/monitor.ts: testing", "extensions/googlechat/src/targets.ts: resolveGoogleChatSpaceChatType", "extensions/imessage/src/monitor-reply-cache.ts: resetIMessageShortIdState", - "extensions/llama-cpp/src/embedding-provider.ts: createLlamaCppMemoryEmbeddingProvider", - "extensions/llama-cpp/src/embedding-provider.ts: DEFAULT_LLAMA_CPP_EMBEDDING_MODEL", - "extensions/llama-cpp/src/embedding-provider.ts: formatLlamaCppSetupError", "extensions/lmstudio/src/stream.ts: resetLmstudioPreloadCooldownForTest", - "extensions/logbook/src/analyze.ts: clockToMs", - "extensions/logbook/src/analyze.ts: extractJsonPayload", "extensions/matrix/src/approval-reactions.ts: clearMatrixApprovalReactionTargetsForTest", "extensions/matrix/src/matrix/client/config.ts: setMatrixAuthClientDepsForTest", "extensions/matrix/src/matrix/monitor/handler.ts: MatrixRetryableInboundError", @@ -75,9 +66,6 @@ export const KNIP_UNUSED_EXPORT_BASELINE = [ "extensions/ollama/src/web-search-provider.ts: testing", "extensions/openshell/src/backend.ts: ENSURE_OPEN_SHELL_REMOTE_REAL_DIRECTORY_SCRIPT", "extensions/openshell/src/backend.ts: PINNED_REMOTE_PATH_MUTATION_SCRIPT", - "extensions/parallel/src/parallel-mcp-search.runtime.ts: extractMcpToolPayload", - "extensions/parallel/src/parallel-mcp-search.runtime.ts: iterMcpMessages", - "extensions/parallel/src/parallel-mcp-search.runtime.ts: selectMcpEnvelope", "extensions/qa-channel/src/inbound.ts: isHttpMediaUrl", "extensions/qa-lab/src/gateway-process-boundary.ts: __testing", "extensions/qa-lab/src/qa-agent-workspace.ts: __testing",