mirror of
https://github.com/volcengine/OpenViking.git
synced 2026-10-01 09:48:03 +08:00
* fix(retrieval): honor tier ceilings and stop cooling unserved recalls Follow-up to #3534, from its post-merge review round. - The abstract-to-overview substitute now applies only to categories whose stored abstract is the whole file body. A resource or skill whose abstract is missing (`processing_mode=vectors_only`) or over the per-entry cap read its body and returned an overview instead, which for a short file is the body almost verbatim — crossing the opt-in deepening boundary those categories are documented to have, and doing it even under an explicit `detail="abstract"`. They now degrade to a bare URI and their body is never read. - A digest reporting `no_relevant` blanks `rendered`, so the client injects nothing, yet those URIs still entered the dedup ledger and were cooled for `dedup_turns` turns. That contradicted the ledger's own bare-URI grace rule and held memories back from the later turn they were relevant to. - Flat retrieval reaches built-in memory types outside the four named ones (`cases`, `patterns`, `tools`, `trajectories`, skill-usage memories) and reported them as an undeclared `memories` category that no tier or penalty table covered, so other-peer hits skipped the score penalty and callers could not pin their tier. The catch-all is now a declared category with both; it stays out of `quotas`, whose buckets it would overlap. Skill-usage memories also stop being misread as the `skills` category. - ZCode, OpenCode and pi own an OV session id but did not forward it, so their recalls silently ran without query expansion or cross-turn dedup. - The context-request deadline covered only the server's 30s rewrite fuse, but the pipeline is serial: expansion, retrieval and budgeting all precede it. 45s covers both fuses and the work between them. - `plugin` config scope and the `/recall` successor example now match what the code actually does. * fix(retrieval): make the context deadline and expansion opt-out reachable Forwarding a session id turns on server-side query expansion, an LLM call with its own 5s fuse, but neither the deadline that was supposed to cover it nor the switch that turns it off reached the two harnesses this PR newly enabled it for. - `contextRequestTimeoutMs()` now derives the deadline from the request body rather than from `cfg` plus a rewrite flag. The body is what states which server stages will run: a session takes the expansion fuse, `rewrite` takes the digest fuse, and a bare retrieval takes neither and keeps the caller's own budget. Reading `cfg` alone could not tell those apart. - OpenCode pinned `timeoutMs: 5000` after spreading the helper's options and pi ignored them entirely, so the helper's deadline was dead code in both. Their own budgets are now defaults rather than ceilings. OpenCode's 5s in particular was shorter than the expansion fuse it had just enabled, so a legal request would have been aborted client-side and dropped back to the path with neither dedup nor expansion. - OpenCode and pi read `OPENVIKING_RECALL_QUERY_EXPANSION` (and `recallQueryExpansion` in their own config files) and set the `configured` flag the shared body builder requires, so the documented opt-out exists where the cost was introduced. - The integration overview no longer implies every harness reads the same environment knobs, and describes the deadline as per-stage rather than rewrite-only.
352 lines
11 KiB
TypeScript
352 lines
11 KiB
TypeScript
/**
|
|
* Pi OpenViking Extension
|
|
*
|
|
* Integrates pi with an OpenViking context database for persistent,
|
|
* cross-session memory. Syncs conversation turns to OV, recalls
|
|
* relevant memories on each prompt, and commits sessions for long-term
|
|
* memory extraction.
|
|
*
|
|
* Design informed by: OpenClaw (synchronous recall), Claude Code plugin
|
|
* (most mature, production-hardened), Hermes (anti-pattern: stale prefetch).
|
|
*/
|
|
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
|
|
import { appendFileSync, mkdirSync } from "node:fs";
|
|
import { dirname } from "node:path";
|
|
import { loadConfig, type OVConfig } from "./config.js";
|
|
import { OVClient } from "./client.js";
|
|
import { RecallManager } from "./recall.js";
|
|
import { SyncManager } from "./sync.js";
|
|
import { buildProfileBlock } from "./shared/profile-inject.mjs";
|
|
import { guardVikingUriToolCall } from "./lib/uri-guard-adapter.mjs";
|
|
import { registerTools } from "./tools.js";
|
|
import { createTakeoverManager } from "./takeover.js";
|
|
|
|
export default async function (pi: ExtensionAPI) {
|
|
// --- Load config ---
|
|
const config = loadConfig(dirname(new URL(import.meta.url).pathname));
|
|
if (!config.enabled) return;
|
|
|
|
// Env overrides
|
|
|
|
// --- Initialize modules ---
|
|
const client = new OVClient(config);
|
|
const sync = new SyncManager(client, config);
|
|
const recall = new RecallManager(client, config, () => sync.sessionId);
|
|
const debugLog = (message: string) => {
|
|
const file = process.env.OV_DEBUG_LOG;
|
|
if (!file) return;
|
|
try {
|
|
mkdirSync(dirname(file), { recursive: true });
|
|
appendFileSync(file, `${new Date().toISOString()} ${message}\n`);
|
|
} catch {
|
|
// Best effort; logging must never affect pi.
|
|
}
|
|
};
|
|
const takeover = createTakeoverManager({ pi, client, sync, config, log: debugLog });
|
|
|
|
// Session state
|
|
let connected = false;
|
|
let bypassed = false;
|
|
let profileBlock = "";
|
|
let archiveOverview = "";
|
|
let toolsRegistered = false;
|
|
let compacted = false;
|
|
let started = false;
|
|
let startPromise: Promise<void> | null = null;
|
|
|
|
// ================================================================
|
|
// Event Handlers
|
|
// ================================================================
|
|
|
|
const start = async (ctx: any): Promise<void> => {
|
|
if (started) return;
|
|
if (startPromise) return startPromise;
|
|
|
|
startPromise = (async () => {
|
|
// Bypass check
|
|
const cwd = process.cwd();
|
|
for (const pattern of config.bypassPatterns) {
|
|
if (matchBypass(cwd, pattern)) {
|
|
bypassed = true;
|
|
started = true;
|
|
return;
|
|
}
|
|
}
|
|
|
|
// Health check
|
|
connected = await client.health();
|
|
if (!connected) {
|
|
if (config.logLevel === "info") {
|
|
ctx.ui.notify("OpenViking: server not reachable", "warning");
|
|
}
|
|
return;
|
|
}
|
|
|
|
// Ensure OV session
|
|
const piSessionId = ctx.sessionManager.getSessionId();
|
|
const ok = await sync.ensureSession(piSessionId);
|
|
if (!ok) {
|
|
if (config.logLevel !== "silent") {
|
|
ctx.ui.notify("OpenViking: failed to create session", "error");
|
|
}
|
|
return;
|
|
}
|
|
await sync.replayPending();
|
|
|
|
// Profile injection
|
|
profileBlock = await buildSessionProfileBlock(client, config);
|
|
|
|
const branch = typeof ctx.sessionManager.getBranch === "function"
|
|
? ctx.sessionManager.getBranch()
|
|
: [];
|
|
if (config.takeoverEnabled) {
|
|
takeover.restore(branch);
|
|
sync.restoreWatermark(takeover.state.syncedEntryCount);
|
|
} else if (sync.sessionId) {
|
|
// Resume rehydration — fetch archive overview if session was previously committed.
|
|
archiveOverview = await fetchArchiveOverview(client, sync.sessionId, config);
|
|
}
|
|
|
|
// Register tools (also needed for pi -c continuations).
|
|
if (!toolsRegistered) {
|
|
registerTools(pi, client, sync);
|
|
toolsRegistered = true;
|
|
}
|
|
updateStatus(ctx, connected, 0, sync.sessionId, config, takeover.state);
|
|
|
|
started = true;
|
|
if (config.logLevel === "info") {
|
|
ctx.ui.notify(`OpenViking connected (${piSessionId.slice(0, 8)}...)`, "info");
|
|
}
|
|
})().finally(() => {
|
|
startPromise = null;
|
|
});
|
|
|
|
return startPromise;
|
|
};
|
|
|
|
// --- session_start ---
|
|
pi.on("session_start", async (event, ctx) => {
|
|
await start(ctx);
|
|
});
|
|
|
|
// --- before_agent_start ---
|
|
pi.on("before_agent_start", async (event, ctx) => {
|
|
// session_start doesn't fire for pi -c continuations.
|
|
await start(ctx);
|
|
|
|
if (!connected || bypassed) return;
|
|
|
|
// Queue recall for the context hook. Pi renders the user message before
|
|
// that hook, so recall latency does not delay the message appearing.
|
|
recall.queueSearch(event.prompt);
|
|
|
|
// Compose system prompt additions
|
|
const parts: string[] = [];
|
|
if (profileBlock) parts.push(profileBlock);
|
|
if (!config.takeoverEnabled && archiveOverview && (compacted || archiveOverview.trim())) {
|
|
parts.push(archiveOverview);
|
|
}
|
|
parts.push("OpenViking tools: viking_search, viking_read, viking_browse, viking_remember, viking_forget, viking_add_resource, viking_archive_expand.");
|
|
|
|
const additions = parts.join("\n\n");
|
|
if (!additions) return;
|
|
|
|
return {
|
|
systemPrompt: event.systemPrompt + "\n\n" + additions,
|
|
};
|
|
});
|
|
|
|
// --- context ---
|
|
pi.on("context", async (event, _ctx) => {
|
|
if (!connected || bypassed) return;
|
|
|
|
// Keep recall synchronous with the provider request so the current prompt
|
|
// still receives current-query memory, without blocking user-message UI.
|
|
await recall.searchPending();
|
|
|
|
const afterTakeover = config.takeoverEnabled
|
|
? takeover.transformContext(event.messages as any)
|
|
: event.messages;
|
|
const messages = recall.injectRecall(afterTakeover);
|
|
return { messages };
|
|
});
|
|
|
|
// --- tool_call ---
|
|
pi.on("tool_call", async (event, _ctx) => {
|
|
const decision = guardVikingUriToolCall(event);
|
|
if (!decision) return;
|
|
return decision;
|
|
});
|
|
|
|
// --- turn_end ---
|
|
pi.on("turn_end", async (event, ctx) => {
|
|
if (!connected || bypassed || !config.syncTurns) return;
|
|
|
|
const branch = ctx.sessionManager.getBranch();
|
|
const result = await sync.syncBranch(branch);
|
|
debugLog(`turn_end: synced ${result.added} entries, ~${result.tokens} tokens`);
|
|
await takeover.onTurnSynced(result.tokens);
|
|
updateStatus(ctx, connected, result.added, sync.sessionId, config, takeover.state);
|
|
});
|
|
|
|
// --- session_before_compact ---
|
|
pi.on("session_before_compact", async (event, _ctx) => {
|
|
if (!connected || bypassed) return;
|
|
|
|
if (config.takeoverEnabled) {
|
|
const prep = (event as any)?.preparation ?? {};
|
|
return await takeover.handleBeforeCompact({
|
|
firstKeptEntryId: prep.firstKeptEntryId,
|
|
tokensBefore: prep.tokensBefore ?? 0,
|
|
});
|
|
}
|
|
|
|
const archiveId = await sync.commit();
|
|
compacted = true;
|
|
|
|
// Cache archive overview for rehydration after compaction
|
|
if (archiveId && sync.sessionId) {
|
|
archiveOverview = await fetchArchiveOverview(
|
|
client, sync.sessionId, config,
|
|
);
|
|
}
|
|
// Return nothing → pi proceeds with default compaction
|
|
});
|
|
|
|
// --- session_shutdown ---
|
|
pi.on("session_shutdown", async (_event, ctx) => {
|
|
if (!connected || bypassed) return;
|
|
|
|
await sync.shutdown();
|
|
if (config.takeoverEnabled) {
|
|
await takeover.shutdown();
|
|
} else {
|
|
await sync.commit();
|
|
}
|
|
});
|
|
|
|
// --- agent_end ---
|
|
pi.on("agent_end", async (_event, _ctx) => {
|
|
recall.invalidate();
|
|
});
|
|
|
|
// ================================================================
|
|
// Commands
|
|
// ================================================================
|
|
|
|
pi.registerCommand("viking", {
|
|
description: "OpenViking status and manual operations. Use 'commit' to force a sync.",
|
|
handler: async (args, ctx) => {
|
|
if (!connected) {
|
|
ctx.ui.notify("OpenViking: not connected", "warning");
|
|
return;
|
|
}
|
|
|
|
if (args?.trim() === "commit") {
|
|
await sync.shutdown();
|
|
const ok = config.takeoverEnabled
|
|
? await takeover.commitAndAdvance()
|
|
: (await sync.commit()) !== null;
|
|
if (ok) {
|
|
ctx.ui.notify("OpenViking: committed successfully", "info");
|
|
} else {
|
|
ctx.ui.notify("OpenViking: commit failed", "error");
|
|
}
|
|
return;
|
|
}
|
|
|
|
// Status
|
|
const sid = sync.sessionId ?? "none";
|
|
const t = takeover.state;
|
|
const takeoverInfo = config.takeoverEnabled
|
|
? ` | takeover: ${t.coveredUserTurns}/${t.lastSeenUserTurns} turns archived, ~${t.pendingTokens} tokens pending`
|
|
: "";
|
|
ctx.ui.notify(
|
|
`OpenViking: ${connected ? "connected" : "disconnected"} | session: ${sid.slice(0, 12)}...${takeoverInfo}`,
|
|
"info",
|
|
);
|
|
},
|
|
});
|
|
}
|
|
|
|
// ================================================================
|
|
// Helper Functions
|
|
// ================================================================
|
|
|
|
/** Simple bypass pattern matching (prefix and glob). */
|
|
function matchBypass(cwd: string, pattern: string): boolean {
|
|
if (pattern.startsWith("*")) {
|
|
return cwd.endsWith(pattern.slice(1));
|
|
}
|
|
if (pattern.endsWith("*")) {
|
|
return cwd.startsWith(pattern.slice(0, -1));
|
|
}
|
|
return cwd === pattern || cwd.startsWith(pattern + "/");
|
|
}
|
|
|
|
/** Build the <openviking-context> profile block. */
|
|
async function buildSessionProfileBlock(
|
|
client: OVClient, config: OVConfig,
|
|
): Promise<string> {
|
|
try {
|
|
const profile = await buildProfileBlock(
|
|
(path: string, init?: any, options?: any) => client.fetchJSON(path, init, 10000),
|
|
config.profileTokenBudget,
|
|
config.peerId,
|
|
);
|
|
if (!profile?.block) return "";
|
|
return [
|
|
'<openviking-context source="session-start">',
|
|
profile.block,
|
|
"</openviking-context>",
|
|
].join("\n");
|
|
} catch {
|
|
return "";
|
|
}
|
|
}
|
|
|
|
/** Fetch archive overview for rehydration using the session context API. */
|
|
async function fetchArchiveOverview(
|
|
client: OVClient, sessionId: string, config: OVConfig,
|
|
): Promise<string> {
|
|
try {
|
|
const ctx = await client.getSessionContext(sessionId, config.resumeContextBudget);
|
|
if (!ctx || !ctx.latest_archive_overview) return "";
|
|
|
|
return [
|
|
'<openviking-context source="session-archive">',
|
|
"<session-archive>",
|
|
ctx.latest_archive_overview,
|
|
"</session-archive>",
|
|
"</openviking-context>",
|
|
].join("\n");
|
|
} catch {
|
|
return "";
|
|
}
|
|
}
|
|
|
|
function updateStatus(
|
|
ctx: any,
|
|
connected: boolean,
|
|
added: number,
|
|
sessionId: string | null,
|
|
config: OVConfig,
|
|
takeoverState?: { pendingTokens?: number; coveredUserTurns?: number },
|
|
): void {
|
|
const setter = ctx?.ui?.setStatus;
|
|
if (typeof setter !== "function") return;
|
|
const threshold = config.takeoverEnabled
|
|
? config.takeoverTokenThreshold
|
|
: config.commitTokenThreshold;
|
|
const pending = config.takeoverEnabled && takeoverState
|
|
? ` · ctx ${takeoverState.coveredUserTurns ?? 0} · ~${takeoverState.pendingTokens ?? 0}/${threshold}`
|
|
: ` · ✎ ${threshold}`;
|
|
const status = `${connected ? "OV ✓" : "OV ✗"} · ↩${added}${pending} · ${sessionId ? sessionId.slice(0, 12) : "none"}`;
|
|
try {
|
|
setter(status);
|
|
} catch {
|
|
// Best effort; pi API shape may vary across fast-moving versions.
|
|
}
|
|
}
|