mirror of
https://github.com/volcengine/OpenViking.git
synced 2026-09-30 09:17:51 +08:00
386 lines
12 KiB
TypeScript
386 lines
12 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 { loadConfigFromModuleUrl, type OVConfig } from "./config.js";
|
|
import { OVClient } from "./client.js";
|
|
import { RecallManager } from "./recall.js";
|
|
import { RecallLedger } from "./shared/recall-ledger.mjs";
|
|
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 = loadConfigFromModuleUrl(import.meta.url);
|
|
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,
|
|
// The ledger keeps request prefixes byte-stable for provider prompt
|
|
// caches (#4137); it is per pi session and opened once the id is known.
|
|
config.recallLedger ? new RecallLedger() : null,
|
|
);
|
|
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();
|
|
recall.openLedger(piSessionId);
|
|
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();
|
|
|
|
// The context hook omits persisted entry ids, but its user messages are a
|
|
// deep copy of the active SessionManager context. Associate those objects
|
|
// with stable ids before takeover may filter the array; retained messages
|
|
// keep object identity through that transform.
|
|
const userEntryIds = ctx.sessionManager.buildContextEntries()
|
|
.filter((entry: any) => entry?.type === "message" && entry.message?.role === "user")
|
|
.map((entry: any) => entry.id as string);
|
|
const messageIds = new WeakMap<object, string>();
|
|
let userIndex = 0;
|
|
for (const message of event.messages as any[]) {
|
|
if (message?.role !== "user") continue;
|
|
const entryId = userEntryIds[userIndex++];
|
|
if (entryId && typeof message === "object") {
|
|
messageIds.set(message, entryId);
|
|
}
|
|
}
|
|
|
|
const afterTakeover = config.takeoverEnabled
|
|
? takeover.transformContext(event.messages as any)
|
|
: event.messages;
|
|
const messages = recall.injectRecall(
|
|
afterTakeover,
|
|
(message) => messageIds.get(message) ?? null,
|
|
);
|
|
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 commitResult = config.takeoverEnabled ? null : await sync.commit();
|
|
const ok = config.takeoverEnabled
|
|
? await takeover.commitAndAdvance()
|
|
: commitResult !== null;
|
|
if (ok) {
|
|
ctx.ui.notify(
|
|
"OpenViking: committed successfully" +
|
|
(commitResult?.trace_id ? ` (trace_id=${commitResult.trace_id})` : ""),
|
|
"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.
|
|
}
|
|
}
|