diff --git a/CHANGELOG.md b/CHANGELOG.md index eccd86b7b71..059edbeec08 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,8 @@ - **mcp (RTK):** expose the RTK tool-output **learn/discover** workflow as two new MCP tools so an agent can grow the RTK filter catalog without leaving the protocol. `omniroute_rtk_discover` analyzes recently captured raw tool output (`discoverRepeatedNoise` / `suggestFilter`) and returns candidate noise patterns plus a suggested filter; `omniroute_rtk_learn` lists the captured command samples (`listRtkCommandSamples`) and resolves a command to its RTK filter id (`commandToId`). Both are read-only (scope `read:compression`), wrap the existing RTK discovery primitives (no new logic in the engine), and log to the MCP audit trail. Regression guard: `tests/unit/compression/rtk-mcp-tools.test.ts` (4). gaps v3.8.42 — T07. +- **compression (LLM tier):** add an **opt-in, default-off LLM-tier compression engine** (`llm`) that condenses the prose of non-system messages via a pluggable chat-completion backend. It mirrors the `llmlingua` engine's contract but is **safe by construction**: the default backend is a **no-op pass-through** (the engine never mutates the payload until an operator both enables it *and* wires a real backend via `setLlmCompressorBackend()`), it is **not** part of the default stacked pipeline, `enabled` defaults to `false`, fenced code blocks and `system` messages are never sent to the model, and every backend error **fails open** (the original segment/body is kept, never thrown). A `minTokens` floor skips small prompts. The real production backend is intentionally a VPS-validated follow-up (Hard Rule #18), exactly as the `llmlingua` worker backend is gated. New `open-sse/services/compression/engines/llm/index.ts`. Regression guard: `tests/unit/compression/llm-compressor-engine.test.ts` (8). gaps v3.8.42 — T05/C3. + ### 🔧 Bug Fixes - **providers (OpenAI-compatible):** Codex MCP / `tool_search` deferred discovery (and `apply_patch`) now works through a **Custom OpenAI-compatible provider**. When such a provider received a Responses-API-shaped request that carried MCP / `tool_search` tools, OmniRoute downgraded it to `/chat/completions`, which drops the deferred tool-discovery mechanism — so the MCP namespaces never surfaced to the model and `apply_patch` was mis-handled as a JSON tool. The executor now detects a Responses-shaped request (`input` / `previous_response_id` / `max_output_tokens` / `reasoning`) that carries `namespace` / `tool_search*` tools and routes it to the upstream `/responses` endpoint natively instead of downgrading (it can also be forced via `providerSpecificData._omnirouteForceResponsesUpstream`). This is a distinct code path from the official Codex OAuth backend (#3033 / #4539, which the earlier fix never touched). Regression guard: `tests/unit/executor-default-base.test.ts`. Thanks to @KooshaPari for the fix. ([#5483](https://github.com/diegosouzapw/OmniRoute/issues/5483)) diff --git a/open-sse/services/compression/engines/index.ts b/open-sse/services/compression/engines/index.ts index 4c9c9f2a2d9..6b4faa390f7 100644 --- a/open-sse/services/compression/engines/index.ts +++ b/open-sse/services/compression/engines/index.ts @@ -7,6 +7,7 @@ import { ccrEngine } from "./ccr/index.ts"; import { llmlinguaEngine } from "./llmlingua/index.ts"; import { ionizerEngine } from "./ionizer/index.ts"; import { relevanceEngine } from "./relevance/index.ts"; +import { llmCompressorEngine } from "./llm/index.ts"; let registered = false; @@ -30,6 +31,7 @@ export function registerBuiltinCompressionEngines(): void { { id: "llmlingua", engine: llmlinguaEngine }, { id: "ionizer", engine: ionizerEngine }, { id: "relevance", engine: relevanceEngine }, + { id: "llm", engine: llmCompressorEngine }, ]; for (const { id, engine } of engines) { diff --git a/open-sse/services/compression/engines/llm/index.ts b/open-sse/services/compression/engines/llm/index.ts new file mode 100644 index 00000000000..0ead7e5e75e --- /dev/null +++ b/open-sse/services/compression/engines/llm/index.ts @@ -0,0 +1,346 @@ +/** + * LLM-tier async compression engine (T05/C3 — opt-in, default-off). + * + * A generic "LLM compressor" tier: it can condense the prose of non-system messages via + * a pluggable LLM backend. It mirrors the `llmlingua` engine's contract exactly, but the + * backend is a full chat-completion model rather than a local ONNX classifier. + * + * ## Safe by construction (default-off) + * - The DEFAULT backend is a no-op (`text => text`), so out of the box the engine NEVER + * mutates the payload — it returns the body unchanged (`compressed:false`). + * - It is NOT part of the default stacked pipeline; an operator must add it explicitly. + * - `enabled` defaults to `false` in the config schema. + * A real LLM backend is wired via `setLlmCompressorBackend()` (the same injection seam the + * tests use). The real production model is intentionally left as a VPS-validated follow-up + * (Hard Rule #18), exactly as the `llmlingua` worker backend is gated. + * + * ## Fail-open everywhere (never throws, never corrupts) + * 1. Backend rejection/error per prose segment → segment kept as-is. + * 2. Any unexpected error in `applyAsync` → original body, `compressed:false`, no throw. + * + * ## Code-block protection (inviolable) + * Fenced code blocks (and other preserved constructs) are tombstoned via + * `extractPreservedBlocks` and re-stitched verbatim — the engine physically never passes + * code to the backend. System messages are never touched. + */ + +import { createCompressionStats, estimateCompressionTokens } from "../../stats.ts"; +import { extractPreservedBlocks } from "../../preservation.ts"; +import type { + CompressionEngine, + CompressionEngineApplyOptions, + EngineConfigField, + EngineValidationResult, +} from "../types.ts"; +import type { CompressionResult } from "../../types.ts"; + +// ─── backend abstraction ────────────────────────────────────────────────────── + +export interface LlmCompressorBackendOptions { + model?: string; + compressionRate?: number; +} + +/** + * A backend takes a prose text segment (+ optional config) and returns a compressed + * version. Any rejection/error MUST be caught by the caller; the engine fail-opens. + */ +export type LlmCompressorBackend = ( + text: string, + opts?: LlmCompressorBackendOptions +) => Promise; + +/** The default production backend: a no-op that returns the text unchanged (pass-through). */ +const noopBackend: LlmCompressorBackend = async (text) => text; + +/** Module-level injectable backend (null = use the no-op default). */ +let _backend: LlmCompressorBackend | null = null; + +/** Override the backend — for tests and for wiring a real LLM. `null` restores the no-op. */ +export function setLlmCompressorBackend(b: LlmCompressorBackend | null): void { + _backend = b; +} + +function resolveBackend(): LlmCompressorBackend { + return _backend ?? noopBackend; +} + +// ─── prose / code splitting (code is never sent to the backend) ───────────────── + +interface TextSegment { + kind: "prose" | "preserved"; + text: string; +} + +function splitProseAndPreserved(text: string): TextSegment[] { + const { text: withPlaceholders, blocks } = extractPreservedBlocks(text); + if (blocks.length === 0) return [{ kind: "prose", text }]; + + const placeholderToOriginal = new Map(blocks.map((b) => [b.placeholder, b.content])); + const escaped = blocks.map((b) => b.placeholder.replace(/[.*+?^${}()|[\]\\]/g, "\\$&")); + const splitRe = new RegExp(`(${escaped.join("|")})`, "g"); + + const segments: TextSegment[] = []; + for (const part of withPlaceholders.split(splitRe)) { + if (!part) continue; + const original = placeholderToOriginal.get(part); + segments.push( + original !== undefined ? { kind: "preserved", text: original } : { kind: "prose", text: part } + ); + } + return segments; +} + +// ─── message processing ───────────────────────────────────────────────────────── + +type MessageLike = { + role?: string; + content?: string | Array>; + [key: string]: unknown; +}; + +async function compressProseText( + text: string, + backend: LlmCompressorBackend, + opts?: LlmCompressorBackendOptions +): Promise<{ text: string; didCompress: boolean }> { + if (!text.trim()) return { text, didCompress: false }; + try { + const compressed = await backend(text, opts); + if (typeof compressed === "string" && compressed.length < text.length) { + return { text: compressed, didCompress: true }; + } + return { text, didCompress: false }; + } catch { + return { text, didCompress: false }; + } +} + +async function compressMessageText( + text: string, + backend: LlmCompressorBackend, + opts?: LlmCompressorBackendOptions +): Promise<{ text: string; didCompress: boolean }> { + const segments = splitProseAndPreserved(text); + let anyCompressed = false; + const parts: string[] = []; + for (const seg of segments) { + if (seg.kind === "preserved") { + parts.push(seg.text); + } else { + const { text: out, didCompress } = await compressProseText(seg.text, backend, opts); + parts.push(out); + if (didCompress) anyCompressed = true; + } + } + return { text: parts.join(""), didCompress: anyCompressed }; +} + +async function processMessages( + messages: MessageLike[], + backend: LlmCompressorBackend, + opts?: LlmCompressorBackendOptions +): Promise<{ messages: MessageLike[]; compressedCount: number }> { + let compressedCount = 0; + const result: MessageLike[] = []; + for (const msg of messages) { + if (msg.role === "system") { + result.push({ ...msg }); + continue; + } + try { + if (typeof msg.content === "string") { + const { text, didCompress } = await compressMessageText(msg.content, backend, opts); + if (didCompress) { + compressedCount++; + result.push({ ...msg, content: text }); + } else { + result.push({ ...msg }); + } + } else if (Array.isArray(msg.content)) { + let changed = false; + const newContent: Array> = []; + for (const part of msg.content) { + if (part["type"] === "text" && typeof part["text"] === "string") { + const { text, didCompress } = await compressMessageText( + part["text"] as string, + backend, + opts + ); + if (didCompress) { + changed = true; + compressedCount++; + newContent.push({ ...part, text }); + } else { + newContent.push(part); + } + } else { + newContent.push(part); + } + } + result.push(changed ? { ...msg, content: newContent } : { ...msg }); + } else { + result.push({ ...msg }); + } + } catch { + result.push({ ...msg }); + } + } + return { messages: result, compressedCount }; +} + +// ─── config schema ──────────────────────────────────────────────────────────── + +const LLM_COMPRESSOR_SCHEMA: EngineConfigField[] = [ + // Default OFF: this tier costs an extra model call and mutates the payload, so it is + // opt-in (Hard Rule #20 spirit) — never on by default. + { key: "enabled", type: "boolean", label: "Enabled", defaultValue: false }, + { key: "model", type: "string", label: "Compression model", defaultValue: "" }, + { + key: "minTokens", + type: "number", + label: "Min tokens (floor)", + defaultValue: 2000, + min: 0, + max: 100000, + }, + { + key: "compressionRate", + type: "number", + label: "Compression rate (keep ratio)", + defaultValue: 0.5, + min: 0.1, + max: 0.9, + }, +]; + +function validateLlmCompressorConfig(config: Record): EngineValidationResult { + const errors: string[] = []; + if (config["enabled"] !== undefined && typeof config["enabled"] !== "boolean") { + errors.push("enabled must be a boolean"); + } + if (config["model"] !== undefined && typeof config["model"] !== "string") { + errors.push("model must be a string"); + } + if (config["minTokens"] !== undefined) { + const v = config["minTokens"]; + if (typeof v !== "number" || Number.isNaN(v) || v < 0) errors.push("minTokens must be a number >= 0"); + } + if (config["compressionRate"] !== undefined) { + const v = config["compressionRate"]; + if (typeof v !== "number" || Number.isNaN(v) || v < 0.1 || v > 0.9) { + errors.push("compressionRate must be a number between 0.1 and 0.9"); + } + } + return { valid: errors.length === 0, errors }; +} + +// ─── engine export ────────────────────────────────────────────────────────────── + +const ENGINE_ID = "llm"; + +export const llmCompressorEngine: CompressionEngine = { + id: ENGINE_ID, + name: "LLM Compressor (opt-in)", + description: + "Opt-in LLM-tier compression: condenses the prose of non-system messages via a pluggable " + + "chat-completion backend. Default-off and a no-op until an operator both enables it and wires " + + "a real backend; fenced code blocks and system messages are never sent to the model. " + + "Fail-opens on any backend error.", + icon: "robot", + targets: ["messages"], + stackable: true, + // Runs after llmlingua (35) but before ultra (40); semantic LLM rewriting is most useful + // once cheaper structural/semantic passes have already reduced the prose. + stackPriority: 38, + metadata: { + id: ENGINE_ID, + name: "LLM Compressor (opt-in)", + description: + "Opt-in LLM-tier prose compression via a pluggable backend. Default-off / no-op; " + + "code blocks and system messages are protected; fail-open on backend error.", + inputScope: "messages", + targetLatencyMs: 1500, + supportsPreview: false, + stable: true, + }, + + /** Synchronous pass-through — the real work is async-only (`applyAsync`). */ + apply(body: Record): CompressionResult { + return { body, compressed: false, stats: null }; + }, + + async applyAsync( + body: Record, + options?: CompressionEngineApplyOptions + ): Promise { + const stepConfig = options?.stepConfig ?? {}; + // Opt-in: only runs when explicitly enabled. + if (stepConfig["enabled"] !== true) { + return { body, compressed: false, stats: null }; + } + + const messages = body["messages"]; + if (!Array.isArray(messages) || messages.length === 0) { + return { body, compressed: false, stats: null }; + } + + const minTokens = + typeof stepConfig["minTokens"] === "number" ? (stepConfig["minTokens"] as number) : 2000; + if (minTokens > 0) { + const nonSystemText = (messages as MessageLike[]) + .filter((m) => m.role !== "system") + .map((m) => (typeof m.content === "string" ? m.content : JSON.stringify(m.content ?? ""))) + .join("\n"); + if (estimateCompressionTokens(nonSystemText) < minTokens) { + return { body, compressed: false, stats: null }; + } + } + + const backendOpts: LlmCompressorBackendOptions = { + model: typeof stepConfig["model"] === "string" ? (stepConfig["model"] as string) : undefined, + compressionRate: + typeof stepConfig["compressionRate"] === "number" + ? (stepConfig["compressionRate"] as number) + : undefined, + }; + + try { + const backend = resolveBackend(); + const start = performance.now(); + const { messages: newMessages, compressedCount } = await processMessages( + messages as MessageLike[], + backend, + backendOpts + ); + if (compressedCount === 0) { + return { body, compressed: false, stats: null }; + } + const newBody: Record = { ...body, messages: newMessages }; + const durationMs = Math.round(performance.now() - start); + const stats = createCompressionStats( + body, + newBody, + "stacked", + [ENGINE_ID], + [`llm-compressed-${compressedCount}-messages`], + durationMs + ); + return { body: newBody, compressed: true, stats }; + } catch { + return { body, compressed: false, stats: null }; + } + }, + + compress(body: Record, config?: Record): CompressionResult { + return this.apply(body, { stepConfig: config ?? {} }); + }, + + getConfigSchema(): EngineConfigField[] { + return LLM_COMPRESSOR_SCHEMA; + }, + + validateConfig(config: Record): EngineValidationResult { + return validateLlmCompressorConfig(config); + }, +}; diff --git a/tests/unit/compression/llm-compressor-engine.test.ts b/tests/unit/compression/llm-compressor-engine.test.ts new file mode 100644 index 00000000000..98535e18883 --- /dev/null +++ b/tests/unit/compression/llm-compressor-engine.test.ts @@ -0,0 +1,110 @@ +import { describe, it, afterEach } from "node:test"; +import assert from "node:assert/strict"; +import { + llmCompressorEngine, + setLlmCompressorBackend, + type LlmCompressorBackend, +} from "../../../open-sse/services/compression/engines/llm/index.ts"; + +// T05/C3 — opt-in LLM-tier compressor. Default-off / no-op; fail-open; code & system protected. + +// A fake backend that halves prose length (always shorter for len >= 2). +const halveBackend: LlmCompressorBackend = async (text) => + text.slice(0, Math.max(1, Math.floor(text.length / 2))); + +const throwingBackend: LlmCompressorBackend = async () => { + throw new Error("backend boom"); +}; + +function body(messages: unknown[]) { + return { messages } as Record; +} + +afterEach(() => setLlmCompressorBackend(null)); + +describe("llmCompressorEngine (T05/C3)", () => { + it("is a pass-through with the default (no-op) backend even when enabled", async () => { + const input = body([{ role: "user", content: "some prose that could be compressed" }]); + const out = await llmCompressorEngine.applyAsync!(input, { + stepConfig: { enabled: true, minTokens: 0 }, + }); + assert.equal(out.compressed, false); + assert.equal(out.body, input, "default no-op backend must not mutate the body"); + }); + + it("does nothing when not enabled, even with a compressing backend", async () => { + setLlmCompressorBackend(halveBackend); + const input = body([{ role: "user", content: "lots of prose here to compress" }]); + const off = await llmCompressorEngine.applyAsync!(input, { stepConfig: { minTokens: 0 } }); + assert.equal(off.compressed, false); + assert.equal(off.body, input); + }); + + it("compresses prose when enabled with a real backend, protecting code + system messages", async () => { + setLlmCompressorBackend(halveBackend); + const code = "```js\nconst secret = 42;\n```"; + const input = body([ + { role: "system", content: "You are a careful assistant. Do not change me." }, + { + role: "user", + content: `Here is a long explanation in prose.\n\n${code}\n\nAnd more prose after.`, + }, + ]); + const out = await llmCompressorEngine.applyAsync!(input, { + stepConfig: { enabled: true, minTokens: 0 }, + }); + assert.equal(out.compressed, true); + const msgs = (out.body as { messages: Array<{ role: string; content: string }> }).messages; + // System message untouched. + assert.equal(msgs[0].content, "You are a careful assistant. Do not change me."); + // The fenced code block survives verbatim (never sent to the backend). + assert.ok(msgs[1].content.includes(code), "code block must be preserved verbatim"); + // The user message got shorter overall. + assert.ok(msgs[1].content.length < (input.messages as { content: string }[])[1].content.length); + }); + + it("fail-opens: a throwing backend leaves the body unchanged", async () => { + setLlmCompressorBackend(throwingBackend); + const input = body([{ role: "user", content: "prose that the backend will choke on" }]); + const out = await llmCompressorEngine.applyAsync!(input, { + stepConfig: { enabled: true, minTokens: 0 }, + }); + assert.equal(out.compressed, false); + assert.equal( + (out.body as { messages: { content: string }[] }).messages[0].content, + "prose that the backend will choke on" + ); + }); + + it("respects the minTokens floor (skips small prompts)", async () => { + setLlmCompressorBackend(halveBackend); + const input = body([{ role: "user", content: "tiny" }]); + const out = await llmCompressorEngine.applyAsync!(input, { + stepConfig: { enabled: true, minTokens: 2000 }, + }); + assert.equal(out.compressed, false); + }); + + it("sync apply is always a no-op pass-through", () => { + const input = body([{ role: "user", content: "anything" }]); + const out = llmCompressorEngine.apply(input); + assert.equal(out.compressed, false); + assert.equal(out.body, input); + }); + + it("validateConfig accepts valid config and rejects bad fields", () => { + assert.equal(llmCompressorEngine.validateConfig({ enabled: false }).valid, true); + assert.equal( + llmCompressorEngine.validateConfig({ enabled: true, compressionRate: 0.5, minTokens: 1000 }).valid, + true + ); + assert.equal(llmCompressorEngine.validateConfig({ enabled: "yes" }).valid, false); + assert.equal(llmCompressorEngine.validateConfig({ compressionRate: 2 }).valid, false); + assert.equal(llmCompressorEngine.validateConfig({ minTokens: -1 }).valid, false); + }); + + it("defaults to disabled in its config schema (opt-in)", () => { + const enabledField = llmCompressorEngine.getConfigSchema!().find((f) => f.key === "enabled"); + assert.equal(enabledField?.defaultValue, false); + }); +});