refactor: consolidate agent/editor plugins under integrations/ (#5491)

Co-authored-by: Claude <noreply@anthropic.com>
This commit is contained in:
Kartik
2026-06-12 10:31:35 +05:30
committed by GitHub
parent c676c2c458
commit 2c796d144f
295 changed files with 94 additions and 81 deletions
@@ -0,0 +1,70 @@
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
import type MemoryClient from "mem0ai";
import type { Mem0Config, ScopeContext } from "../types.ts";
import { DEFAULT_CUSTOM_CATEGORIES } from "../types.ts";
import { resolveAddParams } from "../memory/scoping.ts";
import { captureEvent } from "../telemetry.ts";
interface MessageLike {
role: string;
content?: unknown;
}
function extractText(content: unknown): string | null {
if (typeof content === "string") return content;
if (Array.isArray(content)) {
const texts = content
.filter((b: any) => b.type === "text" && typeof b.text === "string")
.map((b: any) => b.text);
return texts.length > 0 ? texts.join("\n") : null;
}
return null;
}
export function extractConversation(
messages: MessageLike[],
): Array<{ role: "user" | "assistant"; content: string }> {
const result: Array<{ role: "user" | "assistant"; content: string }> = [];
for (const msg of messages) {
if (msg.role !== "user" && msg.role !== "assistant") continue;
const text = extractText(msg.content);
if (!text) continue;
result.push({ role: msg.role as "user" | "assistant", content: text });
}
return result;
}
export function setupAutoCapture(
pi: ExtensionAPI,
mem0: MemoryClient,
config: Mem0Config,
getScopeCtx: () => ScopeContext,
telemetryCtx?: { apiKey?: string },
): void {
if (!config.autoCapture) return;
pi.on("agent_end", async (event) => {
const messages = event.messages ?? [];
const conversation = extractConversation(messages);
if (conversation.length === 0) return;
const scopeCtx = getScopeCtx();
const addParams = resolveAddParams("project", scopeCtx);
try {
await mem0.add(conversation, {
...addParams,
customCategories: DEFAULT_CUSTOM_CATEGORIES,
});
captureEvent("pi.capture.auto", { success: true, message_count: conversation.length }, telemetryCtx);
} catch (err: unknown) {
captureEvent("pi.capture.auto", {
success: false,
error_type: err instanceof Error ? err.name : "unknown",
}, telemetryCtx);
console.error("[mem0] auto-capture failed:", err);
}
});
}
@@ -0,0 +1,297 @@
import { describe, it, expect, vi, beforeEach } from "vitest";
import { registerCommands } from "./commands.ts";
import type { Mem0Config, ScopeContext } from "./types.ts";
vi.mock("./telemetry.ts", () => ({
captureCommandEvent: vi.fn(),
}));
vi.mock("./dream/index.ts", () => ({
acquireDreamLock: vi.fn(() => true),
}));
vi.mock("./dream/prompt.ts", () => ({
DREAM_PROTOCOL: "dream protocol text",
}));
function makeMem0() {
return {
search: vi.fn(),
delete: vi.fn(),
add: vi.fn(),
get: vi.fn(),
getAll: vi.fn(),
update: vi.fn(),
} as any;
}
function makePi() {
const commands = new Map<string, { handler: (args: string, ctx: any) => Promise<void> }>();
return {
registerCommand: vi.fn((name: string, opts: any) => {
commands.set(name, opts);
}),
sendMessage: vi.fn(),
_commands: commands,
_invoke: (name: string, args: string, ctx: any) => commands.get(name)!.handler(args, ctx),
};
}
function makeCtx(confirmResult = true) {
return {
hasUI: true,
ui: {
notify: vi.fn(),
confirm: vi.fn(async () => confirmResult),
select: vi.fn(),
input: vi.fn(),
},
};
}
const defaultConfig: Mem0Config = {
apiKey: "test-key",
userId: "test-user",
autoCapture: false,
defaultScope: "project",
contextInjection: false,
dream: { enabled: false, auto: false, minHours: 24, minSessions: 5, minMemories: 20 },
};
const scopeCtx: ScopeContext = { userId: "test-user", appId: "test-app", runId: "test-run" };
describe("registerCommands", () => {
let pi: ReturnType<typeof makePi>;
let mem0: ReturnType<typeof makeMem0>;
beforeEach(() => {
pi = makePi();
mem0 = makeMem0();
registerCommands(pi as any, mem0, defaultConfig, () => scopeCtx);
});
it("registers all expected commands", () => {
const names = [...pi._commands.keys()];
expect(names).toContain("mem0-remember");
expect(names).toContain("mem0-forget");
expect(names).toContain("mem0-search");
expect(names).toContain("mem0-tour");
expect(names).toContain("mem0-dream");
expect(names).toContain("mem0-pin");
expect(names).toContain("mem0-scope");
expect(names).toContain("mem0-status");
});
describe("/mem0-forget", () => {
it("shows warning when no query provided", async () => {
const ctx = makeCtx();
await pi._invoke("mem0-forget", "", ctx);
expect(ctx.ui.notify).toHaveBeenCalledWith("Usage: /mem0-forget <query>", "warning");
expect(mem0.search).not.toHaveBeenCalled();
});
it("notifies when no memories match", async () => {
const ctx = makeCtx();
mem0.search.mockResolvedValue({ results: [] });
await pi._invoke("mem0-forget", "old preference", ctx);
expect(ctx.ui.notify).toHaveBeenCalledWith("No matching memories found.", "info");
});
it("asks for confirmation before deleting a single match", async () => {
const ctx = makeCtx(true);
mem0.search.mockResolvedValue({ results: [{ id: "abc-123", memory: "test mem" }] });
mem0.delete.mockResolvedValue({ message: "Deleted" });
await pi._invoke("mem0-forget", "test", ctx);
expect(ctx.ui.confirm).toHaveBeenCalledWith(
"Delete this memory?",
expect.stringContaining("test mem"),
);
expect(mem0.delete).toHaveBeenCalledWith("abc-123");
});
it("does not delete when user cancels confirmation", async () => {
const ctx = makeCtx(false);
mem0.search.mockResolvedValue({ results: [{ id: "abc-123", memory: "test mem" }] });
await pi._invoke("mem0-forget", "test", ctx);
expect(ctx.ui.confirm).toHaveBeenCalled();
expect(mem0.delete).not.toHaveBeenCalled();
expect(ctx.ui.notify).toHaveBeenCalledWith("Cancelled.", "info");
});
it("uses select UI for multiple matches and deletes chosen memory", async () => {
const ctx = makeCtx();
mem0.search.mockResolvedValue({
results: [
{ id: "id-1", memory: "mem one" },
{ id: "id-2", memory: "mem two" },
],
});
mem0.delete.mockResolvedValue({ message: "Deleted" });
ctx.ui.select = vi.fn(async (_title: string, options: string[]) => options[1]);
await pi._invoke("mem0-forget", "test", ctx);
expect(ctx.ui.select).toHaveBeenCalledWith(
"Which memory should I delete?",
expect.arrayContaining([
expect.stringContaining("mem one"),
expect.stringContaining("mem two"),
]),
);
expect(mem0.delete).toHaveBeenCalledWith("id-2");
});
it("does not delete when user cancels select", async () => {
const ctx = makeCtx();
ctx.ui.select = vi.fn(async () => undefined);
mem0.search.mockResolvedValue({
results: [
{ id: "id-1", memory: "mem one" },
{ id: "id-2", memory: "mem two" },
],
});
await pi._invoke("mem0-forget", "test", ctx);
expect(mem0.delete).not.toHaveBeenCalled();
expect(ctx.ui.notify).toHaveBeenCalledWith("Cancelled.", "info");
});
});
describe("/mem0-pin", () => {
it("uses update to pin in-place, preserving memory ID", async () => {
const ctx = makeCtx(true);
mem0.search.mockResolvedValue({ results: [{ id: "abc-123", memory: "important fact" }] });
mem0.update.mockResolvedValue([]);
await pi._invoke("mem0-pin", "important", ctx);
expect(ctx.ui.confirm).toHaveBeenCalledWith(
"Pin this memory?",
expect.stringContaining("important fact"),
);
expect(mem0.update).toHaveBeenCalledWith("abc-123", { text: "[PINNED] important fact" });
expect(mem0.add).not.toHaveBeenCalled();
expect(mem0.delete).not.toHaveBeenCalled();
});
it("does not pin when user cancels", async () => {
const ctx = makeCtx(false);
mem0.search.mockResolvedValue({ results: [{ id: "abc-123", memory: "fact" }] });
await pi._invoke("mem0-pin", "fact", ctx);
expect(mem0.update).not.toHaveBeenCalled();
});
it("skips already-pinned memories", async () => {
const ctx = makeCtx();
mem0.search.mockResolvedValue({ results: [{ id: "abc-123", memory: "[PINNED] fact" }] });
await pi._invoke("mem0-pin", "fact", ctx);
expect(ctx.ui.confirm).not.toHaveBeenCalled();
expect(mem0.add).not.toHaveBeenCalled();
expect(ctx.ui.notify).toHaveBeenCalledWith("Already pinned.", "info");
});
it("uses select UI for multiple matches and pins chosen memory", async () => {
const ctx = makeCtx();
mem0.search.mockResolvedValue({
results: [
{ id: "id-1", memory: "fact one" },
{ id: "id-2", memory: "fact two" },
],
});
mem0.update.mockResolvedValue([]);
ctx.ui.select = vi.fn(async (_title: string, options: string[]) => options[1]);
await pi._invoke("mem0-pin", "fact", ctx);
expect(ctx.ui.select).toHaveBeenCalledWith(
"Which memory should I pin?",
expect.arrayContaining([
expect.stringContaining("fact one"),
expect.stringContaining("fact two"),
]),
);
expect(mem0.update).toHaveBeenCalledWith("id-2", { text: "[PINNED] fact two" });
});
it("does not pin when user cancels select", async () => {
const ctx = makeCtx();
ctx.ui.select = vi.fn(async () => undefined);
mem0.search.mockResolvedValue({
results: [
{ id: "id-1", memory: "fact one" },
{ id: "id-2", memory: "fact two" },
],
});
await pi._invoke("mem0-pin", "fact", ctx);
expect(mem0.update).not.toHaveBeenCalled();
expect(ctx.ui.notify).toHaveBeenCalledWith("Cancelled.", "info");
});
});
describe("/mem0-search", () => {
it("always performs semantic search", async () => {
const ctx = makeCtx();
mem0.search.mockResolvedValue({ results: [{ id: "id-1", memory: "result" }] });
await pi._invoke("mem0-search", "my preferences", ctx);
expect(mem0.search).toHaveBeenCalledWith("my preferences", expect.any(Object));
expect(pi.sendMessage).toHaveBeenCalledWith(
expect.objectContaining({ customType: "mem0-search" }),
);
});
it("uses semantic search even for hex-looking strings", async () => {
const ctx = makeCtx();
mem0.search.mockResolvedValue({ results: [] });
await pi._invoke("mem0-search", "abcd1234", ctx);
expect(mem0.search).toHaveBeenCalledWith("abcd1234", expect.any(Object));
expect(mem0.getAll).not.toHaveBeenCalled();
expect(mem0.get).not.toHaveBeenCalled();
});
it("shows empty results message", async () => {
const ctx = makeCtx();
mem0.search.mockResolvedValue({ results: [] });
await pi._invoke("mem0-search", "nonexistent", ctx);
expect(pi.sendMessage).toHaveBeenCalledWith(
expect.objectContaining({ content: "No memories found." }),
);
});
});
describe("/mem0-remember", () => {
it("stores a memory verbatim", async () => {
const ctx = makeCtx();
mem0.add.mockResolvedValue({ message: "Memory stored." });
await pi._invoke("mem0-remember", "I prefer dark mode", ctx);
expect(mem0.add).toHaveBeenCalledWith(
[{ role: "user", content: "I prefer dark mode" }],
expect.objectContaining({ infer: false }),
);
});
it("shows warning when no text provided", async () => {
const ctx = makeCtx();
await pi._invoke("mem0-remember", " ", ctx);
expect(ctx.ui.notify).toHaveBeenCalledWith("Usage: /mem0-remember <text>", "warning");
});
});
});
@@ -0,0 +1,291 @@
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
import type MemoryClient from "mem0ai";
import type { Mem0Config, ScopeContext, Scope } from "./types.ts";
import { DEFAULT_CUSTOM_CATEGORIES } from "./types.ts";
import { resolveSearchFilters, resolveAddParams } from "./memory/scoping.ts";
import { formatMemoryList, formatMemoryCompact, groupByCategory } from "./memory/formatting.ts";
import { DREAM_PROTOCOL } from "./dream/prompt.ts";
import { acquireDreamLock } from "./dream/index.ts";
import { CONFIG_DIR } from "./config/index.ts";
import { captureCommandEvent } from "./telemetry.ts";
export function registerCommands(
pi: ExtensionAPI,
mem0: MemoryClient,
config: Mem0Config,
getScopeCtx: () => ScopeContext,
telemetryCtx?: { apiKey?: string },
): void {
// ── /mem0-remember ──────────────────────────────────────────────────
pi.registerCommand("mem0-remember", {
description: "Store a memory verbatim (no inference)",
handler: async (args, ctx) => {
const text = args?.trim();
if (!text) {
ctx.ui.notify("Usage: /mem0-remember <text>", "warning");
return;
}
const scopeCtx = getScopeCtx();
const addParams = resolveAddParams(config.defaultScope, scopeCtx);
const result = await mem0.add(
[{ role: "user", content: text }],
{ ...addParams, customCategories: DEFAULT_CUSTOM_CATEGORIES, infer: false },
);
const msg = (result as { message?: string }).message ?? "Memory stored.";
captureCommandEvent("mem0-remember", {}, telemetryCtx);
ctx.ui.notify(msg, "info");
},
});
// ── /mem0-forget ────────────────────────────────────────────────────
pi.registerCommand("mem0-forget", {
description: "Delete memories matching a natural language query",
handler: async (args, ctx) => {
const query = args?.trim();
if (!query) {
ctx.ui.notify("Usage: /mem0-forget <query>", "warning");
return;
}
const scopeCtx = getScopeCtx();
const filters = resolveSearchFilters(config.defaultScope, scopeCtx);
const result = await mem0.search(query, { filters });
const memories = result.results ?? [];
if (memories.length === 0) {
captureCommandEvent("mem0-forget", { result_count: 0 }, telemetryCtx);
ctx.ui.notify("No matching memories found.", "info");
return;
}
if (memories.length === 1) {
const target = memories[0];
const confirmed = await ctx.ui.confirm(
"Delete this memory?",
formatMemoryCompact(target),
);
if (!confirmed) {
ctx.ui.notify("Cancelled.", "info");
return;
}
await mem0.delete(target.id);
captureCommandEvent("mem0-forget", { deleted_count: 1 }, telemetryCtx);
ctx.ui.notify(`Deleted: ${formatMemoryCompact(target)}`, "info");
return;
}
const labels = memories.map((m) => formatMemoryCompact(m));
const selected = await ctx.ui.select("Which memory should I delete?", labels);
if (!selected) {
ctx.ui.notify("Cancelled.", "info");
return;
}
const idx = labels.indexOf(selected);
if (idx < 0) return;
const target = memories[idx];
await mem0.delete(target.id);
captureCommandEvent("mem0-forget", { deleted_count: 1 }, telemetryCtx);
ctx.ui.notify(`Deleted: ${formatMemoryCompact(target)}`, "info");
},
});
// ── /mem0-search ────────────────────────────────────────────────────
pi.registerCommand("mem0-search", {
description: "Semantic search across memories",
handler: async (args, ctx) => {
const query = args?.trim();
if (!query) {
ctx.ui.notify("Usage: /mem0-search <query>", "warning");
return;
}
const scopeCtx = getScopeCtx();
const filters = resolveSearchFilters(config.defaultScope, scopeCtx);
const result = await mem0.search(query, { filters });
const memories = result.results ?? [];
captureCommandEvent("mem0-search", { result_count: memories.length }, telemetryCtx);
pi.sendMessage({
customType: "mem0-search",
content: formatMemoryList(memories),
display: true,
});
},
});
// ── /mem0-tour ──────────────────────────────────────────────────────
pi.registerCommand("mem0-tour", {
description: "Browse all memories grouped by category",
handler: async (args, ctx) => {
const raw = args?.trim().toLowerCase();
const validScopes: Scope[] = ["project", "session", "global"];
if (raw && !validScopes.includes(raw as Scope)) {
ctx.ui.notify(`Invalid scope "${raw}". Must be one of: ${validScopes.join(", ")}`, "warning");
return;
}
const scope: Scope = (raw as Scope) || config.defaultScope;
const scopeCtx = getScopeCtx();
const filters = resolveSearchFilters(scope, scopeCtx);
const result = await mem0.getAll({ filters });
const memories = result.results ?? [];
if (memories.length === 0) {
captureCommandEvent("mem0-tour", { memory_count: 0, scope }, telemetryCtx);
pi.sendMessage({ customType: "mem0-tour", content: "No memories found.", display: true });
return;
}
const groups = groupByCategory(memories);
const lines: string[] = [`**Memory Tour** (${memories.length} total, scope: ${scope})`, ""];
for (const [category, items] of groups) {
lines.push(`### ${category} (${items.length})`);
for (const m of items) {
lines.push(`- ${formatMemoryCompact(m)}`);
}
lines.push("");
}
captureCommandEvent("mem0-tour", { memory_count: memories.length, scope }, telemetryCtx);
pi.sendMessage({ customType: "mem0-tour", content: lines.join("\n"), display: true });
},
});
// ── /mem0-dream ─────────────────────────────────────────────────────
pi.registerCommand("mem0-dream", {
description: "Consolidate memories — merge duplicates, prune stale entries, resolve contradictions",
handler: async (_args, ctx) => {
if (!acquireDreamLock(CONFIG_DIR)) {
ctx.ui.notify("A dream consolidation is already in progress.", "warning");
return;
}
captureCommandEvent("mem0-dream", {}, telemetryCtx);
pi.sendMessage({ customType: "mem0-dream", content: DREAM_PROTOCOL, display: true }, { triggerTurn: true });
ctx.ui.notify("Dream consolidation started.", "info");
},
});
// ── /mem0-pin ───────────────────────────────────────────────────────
pi.registerCommand("mem0-pin", {
description: "Pin a memory to protect it from dream pruning",
handler: async (args, ctx) => {
const query = args?.trim();
if (!query) {
ctx.ui.notify("Usage: /mem0-pin <query>", "warning");
return;
}
const scopeCtx = getScopeCtx();
const filters = resolveSearchFilters(config.defaultScope, scopeCtx);
const result = await mem0.search(query, { filters });
const memories = result.results ?? [];
if (memories.length === 0) {
captureCommandEvent("mem0-pin", { result_count: 0 }, telemetryCtx);
ctx.ui.notify("No matching memories found to pin.", "info");
return;
}
if (memories.length === 1) {
const target = memories[0];
const text = target.memory ?? "";
if (text.startsWith("[PINNED]")) {
ctx.ui.notify("Already pinned.", "info");
return;
}
const confirmed = await ctx.ui.confirm(
"Pin this memory?",
formatMemoryCompact(target),
);
if (!confirmed) {
ctx.ui.notify("Cancelled.", "info");
return;
}
await mem0.update(target.id, { text: `[PINNED] ${text}` });
captureCommandEvent("mem0-pin", { pinned: true }, telemetryCtx);
ctx.ui.notify(`Pinned: ${formatMemoryCompact(target)}`, "info");
return;
}
const labels = memories.map((m) => formatMemoryCompact(m));
const selected = await ctx.ui.select("Which memory should I pin?", labels);
if (!selected) {
ctx.ui.notify("Cancelled.", "info");
return;
}
const idx = labels.indexOf(selected);
if (idx < 0) return;
const target = memories[idx];
const selectedText = target.memory ?? "";
if (selectedText.startsWith("[PINNED]")) {
ctx.ui.notify("Already pinned.", "info");
return;
}
await mem0.update(target.id, { text: `[PINNED] ${selectedText}` });
captureCommandEvent("mem0-pin", { pinned: true }, telemetryCtx);
ctx.ui.notify(`Pinned: ${formatMemoryCompact(target)}`, "info");
},
});
// ── /mem0-scope ─────────────────────────────────────────────────────
pi.registerCommand("mem0-scope", {
description: "Change default memory scope for this session (project, session, global)",
handler: async (args, ctx) => {
const scope = args?.trim().toLowerCase();
const valid: Scope[] = ["project", "session", "global"];
if (!scope) {
ctx.ui.notify(`Current scope: ${config.defaultScope}. Usage: /mem0-scope <${valid.join("|")}>`, "info");
return;
}
if (!valid.includes(scope as Scope)) {
ctx.ui.notify(`Invalid scope "${scope}". Must be one of: ${valid.join(", ")}`, "warning");
return;
}
config.defaultScope = scope as Scope;
captureCommandEvent("mem0-scope", { scope }, telemetryCtx);
ctx.ui.notify(`Default scope changed to "${scope}" for this session.`, "info");
},
});
// ── /mem0-status ────────────────────────────────────────────────────
pi.registerCommand("mem0-status", {
description: "Show connection health, identity, project, and memory count",
handler: async (_args, _ctx) => {
const scopeCtx = getScopeCtx();
const filters = resolveSearchFilters("project", scopeCtx);
let count = 0;
let connected = false;
try {
const result = await mem0.getAll({ filters });
count = result.count ?? (result.results ?? []).length;
connected = true;
} catch {
connected = false;
}
const lines = [
"**Mem0 Status**",
"",
`- Connection: ${connected ? "connected" : "disconnected"}`,
`- User: ${scopeCtx.userId}`,
`- Project: ${scopeCtx.appId}`,
`- Session: ${scopeCtx.runId}`,
`- Default scope: ${config.defaultScope}`,
`- Project memories: ${count}`,
`- Auto-capture: ${config.autoCapture ? "on" : "off"}`,
`- Dream: ${config.dream.enabled ? "enabled" : "disabled"}`,
];
captureCommandEvent("mem0-status", { connected, memory_count: count }, telemetryCtx);
pi.sendMessage({ customType: "mem0-status", content: lines.join("\n"), display: true });
},
});
}
@@ -0,0 +1,58 @@
import * as fs from "node:fs";
import * as os from "node:os";
import * as path from "node:path";
import type { Mem0Config, DreamConfig } from "../types.ts";
const AGENT_ROOT = path.join(os.homedir(), ".pi", "agent");
export const CONFIG_DIR = AGENT_ROOT;
const CONFIG_PATH = path.join(AGENT_ROOT, "mem0-config.json");
const DEFAULT_DREAM: DreamConfig = {
enabled: true,
auto: true,
minHours: 24,
minSessions: 5,
minMemories: 20,
};
const DEFAULT_CONFIG: Mem0Config = {
apiKey: "",
userId: "",
autoCapture: true,
defaultScope: "project",
contextInjection: false,
dream: DEFAULT_DREAM,
};
export function loadConfig(): Mem0Config {
let fileConfig: Partial<Mem0Config> = {};
if (fs.existsSync(CONFIG_PATH)) {
try {
const raw = fs.readFileSync(CONFIG_PATH, "utf-8");
fileConfig = JSON.parse(raw);
} catch {
// Corrupted config — use defaults
}
}
const dream: DreamConfig = {
...DEFAULT_DREAM,
...(fileConfig.dream ?? {}),
};
const config: Mem0Config = {
...DEFAULT_CONFIG,
...fileConfig,
dream,
};
if (process.env.MEM0_API_KEY) {
config.apiKey = process.env.MEM0_API_KEY;
}
if (process.env.MEM0_USER_ID) {
config.userId = process.env.MEM0_USER_ID;
}
return config;
}
@@ -0,0 +1,115 @@
import * as fs from "node:fs";
import * as path from "node:path";
import type { DreamState, DreamLock, DreamConfig } from "../types.ts";
const LOCK_STALE_MS = 60 * 60 * 1000;
const DEFAULTS: DreamConfig = {
enabled: true,
auto: true,
minHours: 24,
minSessions: 5,
minMemories: 20,
};
function statePath(stateDir: string): string {
return path.join(stateDir, "mem0-dream-state.json");
}
function lockPath(stateDir: string): string {
return path.join(stateDir, "mem0-dream.lock");
}
function ensureDir(dir: string): void {
try {
fs.mkdirSync(dir, { recursive: true });
} catch { /* exists */ }
}
function readState(stateDir: string): DreamState {
try {
const raw = fs.readFileSync(statePath(stateDir), "utf-8");
return JSON.parse(raw) as DreamState;
} catch {
return { lastConsolidatedAt: 0, sessionsSince: 0, lastSessionId: null };
}
}
function writeState(stateDir: string, state: DreamState): void {
ensureDir(stateDir);
fs.writeFileSync(statePath(stateDir), JSON.stringify(state, null, 2));
}
export function incrementSessionCount(stateDir: string, sessionId: string): void {
const state = readState(stateDir);
if (state.lastSessionId !== sessionId) {
state.sessionsSince++;
state.lastSessionId = sessionId;
writeState(stateDir, state);
}
}
export function checkCheapGates(
stateDir: string,
config: Partial<DreamConfig>,
): { proceed: boolean; reason?: string } {
const minHours = config.minHours ?? DEFAULTS.minHours;
const minSessions = config.minSessions ?? DEFAULTS.minSessions;
const state = readState(stateDir);
const hoursSince = (Date.now() - state.lastConsolidatedAt) / 3_600_000;
if (hoursSince < minHours) {
return { proceed: false, reason: `time: ${hoursSince.toFixed(1)}h < ${minHours}h` };
}
if (state.sessionsSince < minSessions) {
return { proceed: false, reason: `sessions: ${state.sessionsSince} < ${minSessions}` };
}
return { proceed: true };
}
export function checkMemoryGate(
memoryCount: number,
config: Partial<DreamConfig>,
): { pass: boolean; reason?: string } {
const minMemories = config.minMemories ?? DEFAULTS.minMemories;
if (memoryCount < minMemories) {
return { pass: false, reason: `memories: ${memoryCount} < ${minMemories}` };
}
return { pass: true };
}
export function acquireDreamLock(stateDir: string): boolean {
ensureDir(stateDir);
const lp = lockPath(stateDir);
try {
const raw = fs.readFileSync(lp, "utf-8");
const lock = JSON.parse(raw) as DreamLock;
if (Date.now() - lock.startedAt < LOCK_STALE_MS) {
return false;
}
try { fs.unlinkSync(lp); } catch { /* race ok */ }
} catch { /* no lock file */ }
const lock: DreamLock = { pid: process.pid, startedAt: Date.now() };
try {
fs.writeFileSync(lp, JSON.stringify(lock), { flag: "wx" });
return true;
} catch {
return false;
}
}
export function releaseDreamLock(stateDir: string): void {
try { fs.unlinkSync(lockPath(stateDir)); } catch { /* already gone */ }
}
export function recordDreamCompletion(stateDir: string): void {
const state = readState(stateDir);
state.lastConsolidatedAt = Date.now();
state.sessionsSince = 0;
state.lastSessionId = null;
writeState(stateDir, state);
}
@@ -0,0 +1,22 @@
export const DREAM_PROTOCOL = `<mem0-dream>
You are running memory consolidation. Complete these steps using the mem0_memory tool:
1. ORIENT — Call mem0_memory with action "get_all" to list all memories. Count by category. Note oldest/newest.
2. GATHER TARGETS — Review each memory. Classify as:
- DELETE: sensitive information (API keys, passwords, tokens), expired/stale entries, noise, redundant operational details
- MERGE: near-duplicates (same fact stated differently). Keep the better-worded one, delete the other.
- REWRITE: vague, first-person, or poorly-categorized entries. Use mem0_memory "add" with improved text, then "delete" the old one.
- KEEP: everything else.
Skip any memory starting with "[PINNED]".
3. CONSOLIDATE — Execute the changes:
- Delete stale/duplicate entries
- For merges: add the merged text, delete both originals
- For rewrites: add improved version, delete original
4. REPORT — Summarize: how many reviewed, deleted, merged, rewritten, final count.
Quality targets: zero sensitive data stored, zero duplicates, all entries are atomic (one fact each), 15-50 words each.
After consolidation, respond to the user's message normally.
</mem0-dream>`;
@@ -0,0 +1,34 @@
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
import { resolveUserId } from "./entry.ts";
describe("resolveUserId", () => {
const originalEnv = { ...process.env };
afterEach(() => {
process.env = { ...originalEnv };
});
it("returns config userId when set", () => {
expect(resolveUserId("config-user")).toBe("config-user");
});
it("falls back to USER env var", () => {
process.env.USER = "env-user";
delete process.env.USERNAME;
expect(resolveUserId("")).toBe("env-user");
});
it("falls back to USERNAME env var on Windows", () => {
delete process.env.USER;
process.env.USERNAME = "win-user";
expect(resolveUserId("")).toBe("win-user");
});
it("falls back to os.userInfo() when env vars are missing", () => {
delete process.env.USER;
delete process.env.USERNAME;
const result = resolveUserId("");
expect(typeof result).toBe("string");
expect(result.length).toBeGreaterThan(0);
});
});
+146
View File
@@ -0,0 +1,146 @@
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
import MemoryClient from "mem0ai";
import { loadConfig, CONFIG_DIR } from "./config/index.ts";
import { detectAppId, detectRunId, resolveSearchFilters } from "./memory/scoping.ts";
import { registerMemoryTool } from "./memory/tools.ts";
import { registerCommands } from "./commands.ts";
import { setupAutoCapture } from "./capture/index.ts";
import { MEMORY_POLICY } from "./prompt.ts";
import { DREAM_PROTOCOL } from "./dream/prompt.ts";
import {
incrementSessionCount,
checkCheapGates,
checkMemoryGate,
acquireDreamLock,
releaseDreamLock,
recordDreamCompletion,
} from "./dream/index.ts";
import { captureEvent } from "./telemetry.ts";
import * as os from "node:os";
import type { ScopeContext } from "./types.ts";
export function resolveUserId(configUserId: string): string {
if (configUserId) return configUserId;
if (process.env.USER) return process.env.USER;
if (process.env.USERNAME) return process.env.USERNAME;
try { return os.userInfo().username; } catch { return "default"; }
}
export default function mem0Extension(pi: ExtensionAPI): void {
const config = loadConfig();
if (!config.apiKey) {
console.warn("[mem0] No API key found. Set MEM0_API_KEY or add apiKey to ~/.pi/agent/mem0-config.json. Extension disabled.");
return;
}
const mem0 = new MemoryClient({ apiKey: config.apiKey });
const scopeCtx: ScopeContext = {
userId: resolveUserId(config.userId),
appId: "",
runId: "unknown",
};
function getScopeCtx(): ScopeContext {
return scopeCtx;
}
const telemetryCtx = { apiKey: config.apiKey };
// ── Register tool + commands + auto-capture ─────────────────────────
registerMemoryTool(pi, mem0, config, getScopeCtx, telemetryCtx);
registerCommands(pi, mem0, config, getScopeCtx, telemetryCtx);
setupAutoCapture(pi, mem0, config, getScopeCtx, telemetryCtx);
captureEvent("pi.plugin.registered", {
auto_capture: config.autoCapture,
dream_enabled: config.dream.enabled,
default_scope: config.defaultScope,
}, telemetryCtx);
// ── session_start: detect project + session, reconstruct scope ──────
pi.on("session_start", async (_event, ctx) => {
scopeCtx.appId = detectAppId(ctx.cwd);
const sessionFile = ctx.sessionManager?.getSessionFile?.();
scopeCtx.runId = detectRunId(sessionFile);
if (config.userId) {
scopeCtx.userId = config.userId;
}
if (config.dream.enabled) {
incrementSessionCount(CONFIG_DIR, scopeCtx.runId);
}
captureEvent("pi.session.start", {}, telemetryCtx);
});
// ── before_agent_start: append memory policy + auto-dream trigger ───
let dreamTriggered = false;
let dreamChecked = false;
pi.on("before_agent_start", async (event, _ctx) => {
let extra = MEMORY_POLICY;
if (config.dream.enabled && config.dream.auto && !dreamTriggered && !dreamChecked) {
const gates = checkCheapGates(CONFIG_DIR, config.dream);
if (gates.proceed) {
try {
const filters = resolveSearchFilters("project", scopeCtx);
const result = await mem0.getAll({ filters });
const count = result.count ?? (result.results ?? []).length;
dreamChecked = true;
const memGate = checkMemoryGate(count, config.dream);
if (memGate.pass && acquireDreamLock(CONFIG_DIR)) {
dreamTriggered = true;
extra += "\n\n" + DREAM_PROTOCOL;
captureEvent("pi.dream.triggered", { memory_count: count }, telemetryCtx);
}
} catch {
// Transient error — retry next turn
}
}
}
return {
systemPrompt: (event.systemPrompt ?? "") + "\n\n" + extra,
};
});
// ── agent_end: dream completion check ───────────────────────────────
pi.on("agent_end", async (event) => {
if (!dreamTriggered) return;
const messages = event.messages ?? [];
const hadWriteAction = messages.some((m) => {
if (m.role !== "assistant") return false;
const content = Array.isArray(m.content) ? m.content : [];
return content.some(
(block: any) =>
block.type === "tool_use" &&
block.name === "mem0_memory" &&
["add", "delete", "delete_all"].includes(block.input?.action),
);
});
if (hadWriteAction) {
recordDreamCompletion(CONFIG_DIR);
captureEvent("pi.dream.completed", {}, telemetryCtx);
}
releaseDreamLock(CONFIG_DIR);
dreamTriggered = false;
});
// ── session_shutdown: release dream lock if still held ──────────────
pi.on("session_shutdown", async () => {
captureEvent("pi.session.stop", {}, telemetryCtx);
if (dreamTriggered) {
releaseDreamLock(CONFIG_DIR);
dreamTriggered = false;
}
});
}
+34
View File
@@ -0,0 +1,34 @@
export type {
Scope,
Mem0Config,
DreamConfig,
ScopeContext,
CustomCategory,
} from "./types.ts";
export { DEFAULT_CUSTOM_CATEGORIES } from "./types.ts";
export { loadConfig, CONFIG_DIR } from "./config/index.ts";
export { registerMemoryTool, buildToolExecute } from "./memory/tools.ts";
export { detectAppId, detectRunId, resolveSearchFilters, resolveAddParams } from "./memory/scoping.ts";
export { formatAge, formatMemoryCompact, formatMemoryList, groupByCategory } from "./memory/formatting.ts";
export { setupAutoCapture, extractConversation } from "./capture/index.ts";
export {
incrementSessionCount,
checkCheapGates,
checkMemoryGate,
acquireDreamLock,
releaseDreamLock,
recordDreamCompletion,
} from "./dream/index.ts";
export { DREAM_PROTOCOL } from "./dream/prompt.ts";
export { MEMORY_POLICY } from "./prompt.ts";
export { registerCommands } from "./commands.ts";
export { captureEvent, captureToolEvent, captureCommandEvent, _getEventQueue, _resetForTesting } from "./telemetry.ts";
export { default as mem0Extension } from "./entry.ts";
@@ -0,0 +1,43 @@
interface MemoryLike {
id: string;
memory?: string;
categories?: string[];
createdAt?: Date | string;
}
export function formatAge(date: Date | string): string {
const d = typeof date === "string" ? new Date(date) : date;
const ms = Date.now() - d.getTime();
const minutes = Math.floor(ms / 60_000);
if (minutes < 60) return `${minutes}m ago`;
const hours = Math.floor(minutes / 60);
if (hours < 24) return `${hours}h ago`;
const days = Math.floor(hours / 24);
return `${days}d ago`;
}
export function formatMemoryCompact(mem: MemoryLike): string {
const cat = mem.categories?.[0] ?? "uncategorized";
const age = mem.createdAt ? ` (${formatAge(mem.createdAt)})` : "";
return `[${cat}] ${mem.memory ?? "(empty)"}${age} [mem0:${mem.id}]`;
}
export function formatMemoryList(memories: MemoryLike[]): string {
if (memories.length === 0) return "No memories found.";
return memories
.map((m, i) => `${i + 1}. ${formatMemoryCompact(m)}`)
.join("\n");
}
export function groupByCategory(
memories: MemoryLike[],
): Map<string, MemoryLike[]> {
const groups = new Map<string, MemoryLike[]>();
for (const m of memories) {
const cat = m.categories?.[0] ?? "uncategorized";
const list = groups.get(cat) ?? [];
list.push(m);
groups.set(cat, list);
}
return groups;
}
@@ -0,0 +1,85 @@
import { describe, it, expect, vi, beforeEach } from "vitest";
import { detectRunId, resolveSearchFilters, resolveAddParams } from "./scoping.ts";
const mockExecFileSync = vi.fn();
vi.mock("node:child_process", () => ({
execFileSync: (...args: any[]) => mockExecFileSync(...args),
}));
const { detectAppId } = await import("./scoping.ts");
describe("detectAppId", () => {
beforeEach(() => {
mockExecFileSync.mockReset();
});
it("uses git root basename for a git repo", () => {
mockExecFileSync.mockReturnValue("/home/user/projects/my-app\n");
expect(detectAppId("/home/user/projects/my-app")).toBe("my-app");
});
it("returns same app_id from any subdirectory in a monorepo", () => {
mockExecFileSync.mockReturnValue("/home/user/projects/monorepo\n");
const root = detectAppId("/home/user/projects/monorepo");
const sub = detectAppId("/home/user/projects/monorepo/packages/core");
expect(root).toBe("monorepo");
expect(sub).toBe("monorepo");
});
it("falls back to basename when not in a git repo", () => {
mockExecFileSync.mockImplementation(() => {
throw new Error("fatal: not a git repository");
});
expect(detectAppId("/home/user/scratch")).toBe("scratch");
});
});
describe("detectRunId", () => {
it("returns 'unknown' when no session file", () => {
expect(detectRunId(undefined)).toBe("unknown");
});
it("returns a 12-char hex hash for a session file", () => {
const id = detectRunId("/tmp/session-abc.json");
expect(id).toMatch(/^[0-9a-f]{12}$/);
});
it("produces different IDs for different session files", () => {
const a = detectRunId("/tmp/session-a.json");
const b = detectRunId("/tmp/session-b.json");
expect(a).not.toBe(b);
});
});
describe("resolveSearchFilters", () => {
const ctx = { userId: "u1", appId: "a1", runId: "r1" };
it("includes user_id and app_id for project scope", () => {
expect(resolveSearchFilters("project", ctx)).toEqual({ user_id: "u1", app_id: "a1" });
});
it("includes run_id for session scope", () => {
expect(resolveSearchFilters("session", ctx)).toEqual({ user_id: "u1", app_id: "a1", run_id: "r1" });
});
it("uses wildcard app_id for global scope", () => {
expect(resolveSearchFilters("global", ctx)).toEqual({ user_id: "u1", app_id: "*" });
});
});
describe("resolveAddParams", () => {
const ctx = { userId: "u1", appId: "a1", runId: "r1" };
it("includes userId and appId for project scope", () => {
expect(resolveAddParams("project", ctx)).toEqual({ userId: "u1", appId: "a1" });
});
it("includes runId for session scope", () => {
expect(resolveAddParams("session", ctx)).toEqual({ userId: "u1", appId: "a1", runId: "r1" });
});
it("only includes userId for global scope", () => {
expect(resolveAddParams("global", ctx)).toEqual({ userId: "u1" });
});
});
@@ -0,0 +1,51 @@
import * as path from "node:path";
import * as crypto from "node:crypto";
import { execFileSync } from "node:child_process";
import type { Scope, ScopeContext } from "../types.ts";
export function detectAppId(cwd: string): string {
try {
const root = execFileSync("git", ["rev-parse", "--show-toplevel"], {
cwd,
encoding: "utf-8",
timeout: 3000,
stdio: ["ignore", "pipe", "ignore"],
}).trim();
return path.basename(root);
} catch {
return path.basename(cwd);
}
}
export function detectRunId(sessionFile: string | undefined): string {
if (!sessionFile) return "unknown";
return crypto.createHash("sha256").update(sessionFile).digest("hex").slice(0, 12);
}
export function resolveSearchFilters(
scope: Scope,
ctx: ScopeContext,
): Record<string, string> {
switch (scope) {
case "project":
return { user_id: ctx.userId, app_id: ctx.appId };
case "session":
return { user_id: ctx.userId, app_id: ctx.appId, run_id: ctx.runId };
case "global":
return { user_id: ctx.userId, app_id: "*" };
}
}
export function resolveAddParams(
scope: Scope,
ctx: ScopeContext,
): Record<string, string> {
switch (scope) {
case "project":
return { userId: ctx.userId, appId: ctx.appId };
case "session":
return { userId: ctx.userId, appId: ctx.appId, runId: ctx.runId };
case "global":
return { userId: ctx.userId };
}
}
@@ -0,0 +1,195 @@
import type { ExtensionAPI } from "@earendil-works/pi-coding-agent";
import { Type } from "typebox";
import { StringEnum } from "@earendil-works/pi-ai";
import type MemoryClient from "mem0ai";
import type { Scope, ScopeContext, Mem0Config } from "../types.ts";
import { DEFAULT_CUSTOM_CATEGORIES } from "../types.ts";
import { resolveSearchFilters, resolveAddParams } from "./scoping.ts";
import { formatMemoryList } from "./formatting.ts";
import { captureToolEvent } from "../telemetry.ts";
interface MemoryResult {
message?: string;
eventId?: string;
status?: string;
}
const MAX_OUTPUT_LINES = 200;
const MAX_OUTPUT_BYTES = 50_000;
function truncateOutput(text: string): string {
const lines = text.split("\n");
if (lines.length <= MAX_OUTPUT_LINES && text.length <= MAX_OUTPUT_BYTES) {
return text;
}
const kept = lines.slice(0, MAX_OUTPUT_LINES);
let result = kept.join("\n");
if (result.length > MAX_OUTPUT_BYTES) {
result = result.slice(0, MAX_OUTPUT_BYTES);
}
const dropped = lines.length - kept.length;
if (dropped > 0 || text.length > MAX_OUTPUT_BYTES) {
result += `\n\n[Output truncated: showing ${kept.length} of ${lines.length} lines]`;
}
return result;
}
interface ToolParams {
action: "search" | "add" | "get_all" | "update" | "delete" | "delete_all";
query?: string;
content?: string;
memory_id?: string;
scope?: Scope;
}
export function buildToolExecute(
mem0: MemoryClient,
scopeCtx: ScopeContext,
defaultScope: Scope,
) {
return async (params: ToolParams, signal?: AbortSignal) => {
const scope = params.scope ?? defaultScope;
switch (params.action) {
case "search": {
if (signal?.aborted) throw new Error("Cancelled");
if (!params.query) throw new Error("query is required for search");
const filters = resolveSearchFilters(scope, scopeCtx);
const result = await mem0.search(params.query, { filters });
const memories = result.results ?? [];
return {
content: [{ type: "text" as const, text: truncateOutput(formatMemoryList(memories)) }],
details: { matchCount: memories.length },
};
}
case "add": {
if (signal?.aborted) throw new Error("Cancelled");
if (!params.content) throw new Error("content is required for add");
const addParams = resolveAddParams(scope, scopeCtx);
const result = await mem0.add(
[{ role: "user", content: params.content }],
{ ...addParams, customCategories: DEFAULT_CUSTOM_CATEGORIES },
);
const res = result as MemoryResult;
const msg = res.message ?? "Memory stored.";
return {
content: [{ type: "text" as const, text: msg }],
details: { eventId: res.eventId ?? null, status: res.status ?? null },
};
}
case "get_all": {
if (signal?.aborted) throw new Error("Cancelled");
const filters = resolveSearchFilters(scope, scopeCtx);
const result = await mem0.getAll({ filters });
const memories = result.results ?? [];
return {
content: [{ type: "text" as const, text: truncateOutput(formatMemoryList(memories)) }],
details: { totalCount: result.count ?? memories.length },
};
}
case "update": {
if (signal?.aborted) throw new Error("Cancelled");
if (!params.memory_id) throw new Error("memory_id is required for update");
if (!params.content) throw new Error("content is required for update");
const updateResult = await mem0.update(params.memory_id, { text: params.content });
const res = updateResult as MemoryResult;
return {
content: [{ type: "text" as const, text: res.status ?? "Memory updated." }],
details: { memoryId: params.memory_id },
};
}
case "delete": {
if (signal?.aborted) throw new Error("Cancelled");
if (!params.memory_id) throw new Error("memory_id is required for delete");
const result = await mem0.delete(params.memory_id);
return {
content: [{ type: "text" as const, text: result.message ?? "Memory deleted." }],
details: {},
};
}
case "delete_all": {
if (signal?.aborted) throw new Error("Cancelled");
const delParams = resolveAddParams(scope, scopeCtx);
const result = await mem0.deleteAll(delParams);
return {
content: [{ type: "text" as const, text: result.message ?? "All memories deleted." }],
details: {},
};
}
}
};
}
export function registerMemoryTool(
pi: ExtensionAPI,
mem0: MemoryClient,
config: Mem0Config,
getScopeCtx: () => ScopeContext,
telemetryCtx?: { apiKey?: string },
): void {
pi.registerTool({
name: "mem0_memory",
label: "Mem0 Memory",
description:
"Search, add, update, and manage persistent semantic memories powered by Mem0. Memories persist across sessions and devices. Output is truncated to 200 lines / 50KB.",
promptSnippet: "Semantic memory search and storage via Mem0",
promptGuidelines: [
'Use mem0_memory with action "search" when the user asks about past conversations, preferences, or decisions',
'Use mem0_memory with action "add" to save important facts, preferences, goals, decisions, or lessons the user shares',
'Use mem0_memory with action "update" to modify an existing memory — requires memory_id and content. Preserves the memory ID',
"Always use the default project scope unless the user EXPLICITLY asks to search across all projects — only then use scope \"global\"",
"Do NOT pass scope at all for normal queries — omitting it uses the project default automatically",
],
parameters: Type.Object({
action: StringEnum([
"search",
"add",
"get_all",
"update",
"delete",
"delete_all",
] as const),
query: Type.Optional(
Type.String({ description: "Search query or memory text" }),
),
content: Type.Optional(
Type.String({ description: "Memory content to store or updated text" }),
),
memory_id: Type.Optional(
Type.String({ description: "Memory ID for update or delete" }),
),
scope: Type.Optional(
StringEnum(["project", "session", "global"] as const),
),
}),
async execute(toolCallId, params, signal, onUpdate, ctx) {
const scopeCtx = getScopeCtx();
const exec = buildToolExecute(mem0, scopeCtx, config.defaultScope);
const start = Date.now();
try {
const result = await exec(params as ToolParams, signal);
const details = (result as any).details ?? {};
captureToolEvent((params as ToolParams).action, {
success: true,
latency_ms: Date.now() - start,
result_count: details.matchCount ?? details.totalCount ?? undefined,
}, telemetryCtx);
return result;
} catch (err) {
captureToolEvent((params as ToolParams).action, {
success: false,
latency_ms: Date.now() - start,
error_type: err instanceof Error ? err.name : "unknown",
}, telemetryCtx);
throw err;
}
},
});
}
@@ -0,0 +1,16 @@
export const MEMORY_POLICY = `<mem0-memory-policy>
You have persistent semantic memory via the mem0_memory tool, powered by Mem0.
Memory is scoped to the current project by default. Do not change the scope unless explicitly asked.
- "project" (default): memories for this project — use this for all normal queries
- "session": memories from this session only
- "global": all memories across projects — ONLY use when the user explicitly asks for cross-project search
When to use memory:
- Search when the user references past conversations, preferences, or decisions
- Save important facts, user preferences, key decisions, and lessons learned
- Check memory before asking the user something they may have already told you
- Save identity information, goals, relationships, and routines the user shares
Memory persists across sessions and devices via Mem0's cloud.
</mem0-memory-policy>`;
@@ -0,0 +1,240 @@
/**
* Plugin telemetry — anonymous usage tracking via PostHog.
*
* Sends fire-and-forget events to PostHog using native fetch().
* Events are batched and flushed every 5 seconds or when the queue
* reaches 10 events, whichever comes first.
*
* Disable with: MEM0_TELEMETRY=false
*/
import { createHash, randomUUID } from "node:crypto";
import * as fs from "node:fs";
import * as path from "node:path";
import { CONFIG_DIR } from "./config/index.ts";
const POSTHOG_API_KEY = "phc_hgJkUVJFYtmaJqrvf6CYN67TIQ8yhXAkWzUn9AMU4yX";
const POSTHOG_HOST = "https://us.i.posthog.com/i/v0/e/";
const FLUSH_INTERVAL_MS = 5_000;
const FLUSH_THRESHOLD = 10;
let eventQueue: Record<string, unknown>[] = [];
let flushTimer: ReturnType<typeof setInterval> | undefined;
function _loadPluginVersion(): string {
try {
const pkgUrl = new URL("../package.json", import.meta.url);
const pkg = JSON.parse(fs.readFileSync(pkgUrl, "utf-8"));
return pkg.version ?? "unknown";
} catch {
return "unknown";
}
}
const PLUGIN_VERSION = _loadPluginVersion();
// ── Opt-out ──────────────────────────────────────────────────────────────
function isTelemetryEnabled(): boolean {
try {
const val = process.env.MEM0_TELEMETRY;
if (val !== undefined) {
const s = val.toLowerCase();
return s !== "false" && s !== "0" && s !== "no" && s !== "off";
}
return true;
} catch {
return true;
}
}
// ── Identity ─────────────────────────────────────────────────────────────
const TELEMETRY_ID_PATH = path.join(CONFIG_DIR, "mem0-telemetry-id.json");
let _cachedAnonymousId: string | undefined;
function getOrCreateAnonymousId(): string {
if (_cachedAnonymousId) return _cachedAnonymousId;
try {
if (fs.existsSync(TELEMETRY_ID_PATH)) {
const data = JSON.parse(fs.readFileSync(TELEMETRY_ID_PATH, "utf-8"));
if (data.anonymousId) {
_cachedAnonymousId = data.anonymousId;
return _cachedAnonymousId!;
}
}
} catch { /* ignore */ }
const newId = `pi-mem0-anon-${randomUUID().replace(/-/g, "")}`;
try {
fs.mkdirSync(CONFIG_DIR, { recursive: true });
fs.writeFileSync(TELEMETRY_ID_PATH, JSON.stringify({ anonymousId: newId }), "utf-8");
} catch { /* ignore */ }
_cachedAnonymousId = newId;
return newId;
}
function getDistinctId(apiKey?: string): string {
if (apiKey) {
return createHash("sha256").update(apiKey).digest("hex");
}
return getOrCreateAnonymousId();
}
let _identifyDone = false;
function maybeBuildIdentifyEvent(distinctId: string): Record<string, unknown> | null {
if (_identifyDone) return null;
if (!distinctId || distinctId.startsWith("pi-mem0-anon-")) return null;
try {
if (!fs.existsSync(TELEMETRY_ID_PATH)) {
_identifyDone = true;
return null;
}
const data = JSON.parse(fs.readFileSync(TELEMETRY_ID_PATH, "utf-8"));
const storedAnon = data.anonymousId;
if (!storedAnon) {
_identifyDone = true;
return null;
}
const identifyEvent = {
event: "$identify",
distinct_id: distinctId,
properties: { $anon_distinct_id: storedAnon, $lib: "posthog-node" },
};
try {
fs.unlinkSync(TELEMETRY_ID_PATH);
} catch { /* ignore */ }
_identifyDone = true;
_cachedAnonymousId = undefined;
return identifyEvent;
} catch {
return null;
}
}
// ── Flush machinery ──────────────────────────────────────────────────────
function ensureFlushTimer(): void {
if (flushTimer) return;
flushTimer = setInterval(flushEvents, FLUSH_INTERVAL_MS);
if (typeof flushTimer === "object" && "unref" in flushTimer) {
flushTimer.unref();
}
}
let _exitHandlerInstalled = false;
function ensureExitHandler(): void {
if (_exitHandlerInstalled) return;
_exitHandlerInstalled = true;
process.on("beforeExit", async () => {
if (eventQueue.length === 0) return;
const batch = eventQueue;
eventQueue = [];
const body = JSON.stringify({ api_key: POSTHOG_API_KEY, batch });
try {
await fetch(POSTHOG_HOST, {
method: "POST",
headers: {
"Content-Type": "application/json",
"Content-Length": String(Buffer.byteLength(body)),
},
body,
signal: AbortSignal.timeout(3_000),
});
} catch { /* silently swallow */ }
});
}
function flushEvents(): void {
if (eventQueue.length === 0) return;
const batch = eventQueue;
eventQueue = [];
const body = JSON.stringify({ api_key: POSTHOG_API_KEY, batch });
fetch(POSTHOG_HOST, {
method: "POST",
headers: {
"Content-Type": "application/json",
"Content-Length": String(Buffer.byteLength(body)),
},
body,
signal: AbortSignal.timeout(3_000),
}).catch(() => { /* silently swallow */ });
}
// ── Public API ───────────────────────────────────────────────────────────
export function captureEvent(
eventName: string,
properties: Record<string, unknown> = {},
ctx?: { apiKey?: string },
): void {
if (!isTelemetryEnabled()) return;
try {
const distinctId = getDistinctId(ctx?.apiKey);
const identifyEvent = maybeBuildIdentifyEvent(distinctId);
if (identifyEvent) {
eventQueue.push(identifyEvent);
}
eventQueue.push({
event: eventName,
distinct_id: distinctId,
properties: {
source: "PI_AGENT_PLUGIN",
language: "node",
plugin_version: PLUGIN_VERSION,
node_version: process.version,
os: process.platform,
$process_person_profile: false,
$lib: "posthog-node",
...properties,
},
});
ensureFlushTimer();
ensureExitHandler();
if (eventQueue.length >= FLUSH_THRESHOLD) {
flushEvents();
}
} catch { /* silently swallow */ }
}
export function captureToolEvent(
action: string,
properties: Record<string, unknown> = {},
ctx?: { apiKey?: string },
): void {
captureEvent("pi.tool.mem0_memory", { action, ...properties }, ctx);
}
export function captureCommandEvent(
command: string,
properties: Record<string, unknown> = {},
ctx?: { apiKey?: string },
): void {
captureEvent(`pi.command.${command}`, properties, ctx);
}
// ── Test helpers ─────────────────────────────────────────────────────────
export function _getEventQueue(): Record<string, unknown>[] {
return eventQueue;
}
export function _resetForTesting(): void {
eventQueue = [];
if (flushTimer) {
clearInterval(flushTimer);
flushTimer = undefined;
}
_cachedAnonymousId = undefined;
_identifyDone = false;
}
+52
View File
@@ -0,0 +1,52 @@
export type Scope = "project" | "session" | "global";
export interface DreamConfig {
enabled: boolean;
auto: boolean;
minHours: number;
minSessions: number;
minMemories: number;
}
export interface Mem0Config {
apiKey: string;
userId: string;
autoCapture: boolean;
defaultScope: Scope;
contextInjection: boolean;
dream: DreamConfig;
}
export interface DreamState {
lastConsolidatedAt: number;
sessionsSince: number;
lastSessionId: string | null;
}
export interface DreamLock {
pid: number;
startedAt: number;
}
export interface ScopeContext {
userId: string;
appId: string;
runId: string;
}
export interface CustomCategory {
[key: string]: string;
}
export const DEFAULT_CUSTOM_CATEGORIES: CustomCategory[] = [
{ identity: "Personal details, background, and self-descriptions" },
{ preferences: "Likes, dislikes, habits, and preferred ways of doing things" },
{ goals: "Objectives, aspirations, and targets the user is working toward" },
{ projects: "Ongoing work, initiatives, and areas of focus" },
{ decisions: "Choices made, rationale, and trade-offs considered" },
{ technical: "Technical knowledge, tools, configurations, and environment details" },
{ relationships: "People, teams, organizations, and their roles" },
{ routines: "Recurring patterns, workflows, schedules, and processes" },
{ lessons: "Insights learned, mistakes to avoid, and best practices discovered" },
{ work: "Professional context, role, responsibilities, and work environment" },
];