headroom_3 / plugins /openclaw /src /engine.ts
chopratejas's picture
OpenClaw plugin fixes, telemetry fix, cost tracker improvement, tokenBudget support
0b469f5
Raw History Blame
7.28 kB
/**
* HeadroomContextEngine — ContextEngine implementation for OpenClaw.
*
* Compresses tool outputs and conversation context using the Headroom proxy.
* Zero LLM calls — all compression is algorithmic (SmartCrusher, ContentRouter, etc.)
*/
/* eslint-disable @typescript-eslint/no-explicit-any */
import { compress } from "headroom-ai";
import { ProxyManager, type ProxyManagerConfig, type ProxyManagerLogger } from "./proxy-manager.js";
import { agentToOpenAI, openAIToAgent } from "./convert.js";
export interface HeadroomEngineConfig extends ProxyManagerConfig {
enabled?: boolean;
}
export class HeadroomContextEngine {
readonly info = {
id: "headroom",
name: "Headroom Context Compression",
version: "0.1.0",
ownsCompaction: true,
};
private proxyManager: ProxyManager;
private proxyUrl: string | null = null;
private config: HeadroomEngineConfig;
private logger: ProxyManagerLogger;
private stats = {
totalCompressions: 0,
totalTokensSaved: 0,
totalTokensBefore: 0,
compactions: 0,
};
constructor(config: HeadroomEngineConfig = {}, logger?: ProxyManagerLogger) {
this.config = config;
this.logger = logger ?? {
info: (m) => console.log(`[headroom] ${m}`),
warn: (m) => console.warn(`[headroom] ${m}`),
error: (m) => console.error(`[headroom] ${m}`),
debug: () => {},
};
this.proxyManager = new ProxyManager(config, this.logger);
}
// === ContextEngine Lifecycle ===
async bootstrap(params: {
sessionId: string;
sessionKey?: string;
sessionFile: string;
}): Promise<{ bootstrapped: boolean; reason?: string }> {
if (this.config.enabled === false) {
return { bootstrapped: false, reason: "disabled" };
}
try {
this.proxyUrl = await this.proxyManager.start();
this.logger.info(`Engine bootstrapped (proxy: ${this.proxyUrl})`);
return { bootstrapped: true };
} catch (error) {
this.logger.error(`Bootstrap failed: ${error}`);
return { bootstrapped: false, reason: String(error) };
}
}
async ingest(params: {
sessionId: string;
message: any;
isHeartbeat?: boolean;
}): Promise<{ ingested: boolean }> {
// No-op: OpenClaw's runtime stores messages. We don't need a separate store.
return { ingested: true };
}
async ingestBatch?(params: {
sessionId: string;
messages: any[];
isHeartbeat?: boolean;
}): Promise<{ ingestedCount: number }> {
return { ingestedCount: params.messages.length };
}
/**
* Assemble context for the model — THE CORE HOOK.
*
* Converts AgentMessage[] → OpenAI format → compress() → AgentMessage[]
*/
async assemble(params: {
sessionId: string;
messages: any[];
tokenBudget?: number;
model?: string;
prompt?: string;
}): Promise<{
messages: any[];
estimatedTokens: number;
systemPromptAddition?: string;
}> {
if (!this.proxyUrl || this.config.enabled === false) {
// Fallback: return messages unchanged
return { messages: params.messages, estimatedTokens: 0 };
}
try {
// Convert AgentMessage → OpenAI format
const openaiMessages = agentToOpenAI(params.messages);
// Compress via proxy — pass tokenBudget so RollingWindow enforces it
const result = await compress(openaiMessages, {
model: params.model ?? "claude-sonnet-4-5",
baseUrl: this.proxyUrl,
fallback: true,
tokenBudget: params.tokenBudget,
});
if (!result.compressed || result.tokensSaved === 0) {
return { messages: params.messages, estimatedTokens: result.tokensBefore };
}
// Convert back to AgentMessage format
const compressedAgentMessages = openAIToAgent(result.messages);
// Track stats
this.stats.totalCompressions++;
this.stats.totalTokensSaved += result.tokensSaved;
this.stats.totalTokensBefore += result.tokensBefore;
this.logger.debug(
`Assembled: ${result.tokensBefore} → ${result.tokensAfter} tokens (saved ${result.tokensSaved})`,
);
return {
messages: compressedAgentMessages,
estimatedTokens: result.tokensAfter,
systemPromptAddition:
result.tokensSaved > 100
? `[Context compressed by Headroom: ${result.tokensSaved} tokens saved. Use headroom_retrieve with the hash to get full details.]`
: undefined,
};
} catch (error) {
this.logger.error(`Assemble failed: ${error}`);
// Graceful fallback: return original messages
return { messages: params.messages, estimatedTokens: 0 };
}
}
/**
* Compact context — zero-cost alternative to LLM summarization.
*
* Calls compress() with the token budget, which triggers:
* - SmartCrusher: aggressive JSON compression (70-90% on tool outputs)
* - Kompress: ModernBERT text compression (40-60% on assistant text)
* - RollingWindow: drops oldest messages if still over budget
* - CCR: stores originals for retrieval via headroom_retrieve tool
*
* Zero LLM calls. All algorithmic.
*/
async compact(params: {
sessionId: string;
sessionFile: string;
tokenBudget?: number;
force?: boolean;
runtimeContext?: any;
}): Promise<{
ok: boolean;
compacted: boolean;
reason?: string;
result?: {
tokensBefore: number;
tokensAfter?: number;
};
}> {
if (!this.proxyUrl) {
return { ok: false, compacted: false, reason: "Proxy not available" };
}
// Read current messages from session file if available
// For now, compact() works in tandem with assemble() — the next assemble()
// call will compress with the token budget. When compact() is called
// independently, we report success since our pipeline handles it.
//
// TODO: Read session file, extract messages, call compress() with tokenBudget,
// write back compacted messages.
this.stats.compactions++;
this.logger.info(
`Compact called (budget: ${params.tokenBudget ?? "none"}, force: ${params.force ?? false})`,
);
return {
ok: true,
compacted: true,
reason: "Headroom applies SmartCrusher + Kompress + RollingWindow on next assemble()",
};
}
async afterTurn?(params: {
sessionId: string;
messages: any[];
prePromptMessageCount: number;
isHeartbeat?: boolean;
}): Promise<void> {
// Optional: could log stats or trigger learning
}
async prepareSubagentSpawn?(params: {
parentSessionKey: string;
childSessionKey: string;
ttlMs?: number;
}): Promise<{ rollback: () => Promise<void> } | undefined> {
// Subagent context is compressed naturally via assemble()
return undefined;
}
async onSubagentEnded?(params: {
childSessionKey: string;
reason: string;
}): Promise<void> {
// No-op
}
async dispose(): Promise<void> {
await this.proxyManager.stop();
this.logger.info(
`Engine disposed. Stats: ${this.stats.totalCompressions} compressions, ` +
`${this.stats.totalTokensSaved} tokens saved`,
);
}
// --- Public API ---
getStats() {
return { ...this.stats };
}
getProxyUrl(): string | null {
return this.proxyUrl;
}
}