@llmnesia/mcp
Advanced tools
| #!/usr/bin/env node | ||
| import { createRequire as __llmnesiaCreateRequire } from "node:module"; | ||
| const require = __llmnesiaCreateRequire(import.meta.url); | ||
| var __create = Object.create; | ||
| var __defProp = Object.defineProperty; | ||
| var __getOwnPropDesc = Object.getOwnPropertyDescriptor; | ||
| var __getOwnPropNames = Object.getOwnPropertyNames; | ||
| var __getProtoOf = Object.getPrototypeOf; | ||
| var __hasOwnProp = Object.prototype.hasOwnProperty; | ||
| var __require = /* @__PURE__ */ ((x) => typeof require !== "undefined" ? require : typeof Proxy !== "undefined" ? new Proxy(x, { | ||
| get: (a, b) => (typeof require !== "undefined" ? require : a)[b] | ||
| }) : x)(function(x) { | ||
| if (typeof require !== "undefined") return require.apply(this, arguments); | ||
| throw Error('Dynamic require of "' + x + '" is not supported'); | ||
| }); | ||
| var __commonJS = (cb, mod) => function __require2() { | ||
| return mod || (0, cb[__getOwnPropNames(cb)[0]])((mod = { exports: {} }).exports, mod), mod.exports; | ||
| }; | ||
| var __export = (target, all) => { | ||
| for (var name in all) | ||
| __defProp(target, name, { get: all[name], enumerable: true }); | ||
| }; | ||
| var __copyProps = (to, from, except, desc) => { | ||
| if (from && typeof from === "object" || typeof from === "function") { | ||
| for (let key of __getOwnPropNames(from)) | ||
| if (!__hasOwnProp.call(to, key) && key !== except) | ||
| __defProp(to, key, { get: () => from[key], enumerable: !(desc = __getOwnPropDesc(from, key)) || desc.enumerable }); | ||
| } | ||
| return to; | ||
| }; | ||
| var __toESM = (mod, isNodeMode, target) => (target = mod != null ? __create(__getProtoOf(mod)) : {}, __copyProps( | ||
| // If the importer is in node compatibility mode or this is not an ESM | ||
| // file that has been converted to a CommonJS file using a Babel- | ||
| // compatible transform (i.e. "__esModule" has not been set), then set | ||
| // "default" to the CommonJS "module.exports" for node compatibility. | ||
| isNodeMode || !mod || !mod.__esModule ? __defProp(target, "default", { value: mod, enumerable: true }) : target, | ||
| mod | ||
| )); | ||
| // src/version.ts | ||
| var VERSION = true ? "0.2.4" : "0.0.0-dev"; | ||
| export { | ||
| __require, | ||
| __commonJS, | ||
| __export, | ||
| __toESM, | ||
| VERSION | ||
| }; |
| #!/usr/bin/env node | ||
| import { createRequire as __llmnesiaCreateRequire } from "node:module"; | ||
| const require = __llmnesiaCreateRequire(import.meta.url); | ||
| import { | ||
| VERSION | ||
| } from "./chunk-JSNSTAFD.js"; | ||
| // src/config.ts | ||
| import { readFileSync } from "fs"; | ||
| import { homedir } from "os"; | ||
| import { join, resolve } from "path"; | ||
| var CORPUS_FORMAT_VERSION = 2; | ||
| function nativeHostConfigPath() { | ||
| return join(homedir(), ".llmnesia", "config.json"); | ||
| } | ||
| function readNativeHostCorpusRoot() { | ||
| try { | ||
| const parsed = JSON.parse(readFileSync(nativeHostConfigPath(), "utf8")); | ||
| return typeof parsed.corpusRoot === "string" ? parsed.corpusRoot : void 0; | ||
| } catch { | ||
| return void 0; | ||
| } | ||
| } | ||
| function resolveCorpusDir(explicit) { | ||
| const raw = explicit ?? process.env.LLMNESIA_CORPUS_DIR ?? readNativeHostCorpusRoot() ?? "~/.llmnesia/corpus"; | ||
| return expandHome(raw); | ||
| } | ||
| function expandHome(p) { | ||
| if (p === "~") { | ||
| return homedir(); | ||
| } | ||
| if (p.startsWith("~/") || p.startsWith("~\\")) { | ||
| return resolve(homedir(), p.slice(2)); | ||
| } | ||
| return resolve(p); | ||
| } | ||
| // src/paths.ts | ||
| import { join as join2 } from "path"; | ||
| import { readdir } from "fs/promises"; | ||
| import { createHash } from "crypto"; | ||
| var INDEX_FORMAT_VERSION = 3; | ||
| function indexDbFilename(version) { | ||
| return version === 1 ? "llmnesia.db" : `llmnesia-v${version}.db`; | ||
| } | ||
| async function staleIndexFiles(dir) { | ||
| const keep = indexDbFilename(INDEX_FORMAT_VERSION); | ||
| let entries; | ||
| try { | ||
| entries = await readdir(dir); | ||
| } catch { | ||
| return []; | ||
| } | ||
| return entries.filter( | ||
| (name) => /^llmnesia(-v\d+)?\.db$/.test(name) && name !== keep | ||
| ); | ||
| } | ||
| var CorpusPaths = class { | ||
| constructor(root) { | ||
| this.root = root; | ||
| } | ||
| root; | ||
| get metaFile() { | ||
| return join2(this.root, "meta.json"); | ||
| } | ||
| get inboxDir() { | ||
| return join2(this.root, "inbox"); | ||
| } | ||
| get processedDir() { | ||
| return join2(this.inboxDir, "processed"); | ||
| } | ||
| /** Input that was readable but not clean enough to discard automatically. */ | ||
| get quarantineDir() { | ||
| return join2(this.inboxDir, "quarantine"); | ||
| } | ||
| get conversationsDir() { | ||
| return join2(this.root, "conversations"); | ||
| } | ||
| get memoriesDir() { | ||
| return join2(this.root, "memories"); | ||
| } | ||
| get indexDir() { | ||
| return join2(this.root, "index"); | ||
| } | ||
| get indexDbFile() { | ||
| return join2(this.indexDir, indexDbFilename(INDEX_FORMAT_VERSION)); | ||
| } | ||
| /** Directory holding a platform's conversation JSON files. */ | ||
| platformDir(platform) { | ||
| return join2(this.conversationsDir, safeSegment(platform)); | ||
| } | ||
| /** Absolute path to the JSON file backing a given docId. */ | ||
| conversationFile(platform, docId) { | ||
| return this.compressedConversationFile(platform, docId); | ||
| } | ||
| /** Canonical v2 path. Compression is lossless and independently readable. */ | ||
| compressedConversationFile(platform, docId) { | ||
| return join2(this.platformDir(platform), `${slugifyDocId(docId)}.json.gz`); | ||
| } | ||
| /** v1 path accepted until confirmed/resumable compaction migrates it. */ | ||
| legacyConversationFile(platform, docId) { | ||
| return join2(this.platformDir(platform), `${slugifyDocId(docId)}.json`); | ||
| } | ||
| /** | ||
| * Absolute path to the JSON file backing a given memoryId. Memories are flat | ||
| * rather than partitioned by platform: they are user-curated and few, and a | ||
| * memory is not owned by any one platform the way a conversation is. | ||
| */ | ||
| memoryFile(memoryId) { | ||
| return join2(this.memoriesDir, `${slugifyDocId(memoryId)}.json`); | ||
| } | ||
| }; | ||
| var SLUG_MAX = 120; | ||
| function slugifyDocId(docId) { | ||
| const lossless = docId.length > 0 && docId.length <= SLUG_MAX && /^[a-zA-Z0-9._-]+$/.test(docId); | ||
| if (lossless) { | ||
| return docId; | ||
| } | ||
| const hash = createHash("sha1").update(docId).digest("hex").slice(0, 12); | ||
| const readable = docId.replace(/[^a-zA-Z0-9._-]+/g, "-").replace(/^-+|-+$/g, "").slice(0, SLUG_MAX - 13); | ||
| return readable.length > 0 ? `${readable}-${hash}` : hash; | ||
| } | ||
| function safeSegment(value) { | ||
| const cleaned = value.replace(/[^a-zA-Z0-9._-]+/g, "-").replace(/^[.-]+|[.-]+$/g, ""); | ||
| return cleaned.length > 0 ? cleaned : "unknown"; | ||
| } | ||
| // src/corpus.ts | ||
| import { promises as fs, readFileSync as readFileSync2 } from "fs"; | ||
| import { dirname, join as join3 } from "path"; | ||
| import { promisify } from "util"; | ||
| import { gzip as gzipCallback, gunzip as gunzipCallback, gunzipSync } from "zlib"; | ||
| // src/serialize.ts | ||
| function canonicalizeConversation(record) { | ||
| return { | ||
| conversation: canonicalConversationDoc(record.conversation), | ||
| messages: [...record.messages].sort((a, b) => a.msgIndex - b.msgIndex).map(canonicalMessageDoc) | ||
| }; | ||
| } | ||
| function canonicalizeMemory(memory) { | ||
| const out = { | ||
| memoryId: memory.memoryId, | ||
| kind: memory.kind, | ||
| label: memory.label, | ||
| text: memory.text, | ||
| source: memory.source | ||
| }; | ||
| assignIfDefined(out, "sourceDocId", memory.sourceDocId); | ||
| assignIfDefined(out, "sourcePlatform", memory.sourcePlatform); | ||
| assignIfDefined(out, "sourceUrl", memory.sourceUrl); | ||
| assignIfDefined(out, "sourceTitle", memory.sourceTitle); | ||
| assignIfDefined(out, "sourceMsgStart", memory.sourceMsgStart); | ||
| assignIfDefined(out, "sourceMsgEnd", memory.sourceMsgEnd); | ||
| out.pinned = memory.pinned; | ||
| out.enabled = memory.enabled; | ||
| if (memory.tags && memory.tags.length > 0) { | ||
| out.tags = [...memory.tags].sort(); | ||
| } | ||
| out.createdAt = memory.createdAt; | ||
| out.updatedAt = memory.updatedAt; | ||
| assignIfDefined(out, "lastUsedAt", memory.lastUsedAt); | ||
| assignIfDefined(out, "useCount", memory.useCount); | ||
| return out; | ||
| } | ||
| function canonicalConversationDoc(c) { | ||
| const out = { | ||
| docId: c.docId, | ||
| platform: c.platform, | ||
| title: c.title, | ||
| url: c.url, | ||
| createdAt: c.createdAt, | ||
| contentUpdatedAt: c.contentUpdatedAt, | ||
| indexedAt: c.indexedAt, | ||
| corpusChangedAt: c.corpusChangedAt, | ||
| updatedAt: c.updatedAt, | ||
| lastSeenAt: c.lastSeenAt, | ||
| messageCount: c.messageCount, | ||
| preview: c.preview, | ||
| hasImages: c.hasImages, | ||
| pinned: c.pinned, | ||
| indexLevel: c.indexLevel, | ||
| estimatedBytes: c.estimatedBytes | ||
| }; | ||
| assignIfDefined(out, "searchDocLength", c.searchDocLength); | ||
| assignIfDefined(out, "contentHash", c.contentHash); | ||
| assignIfDefined(out, "truncated", c.truncated); | ||
| assignIfDefined(out, "kind", c.kind); | ||
| assignIfDefined(out, "sourceId", c.sourceId); | ||
| assignIfDefined(out, "baseUrl", c.baseUrl); | ||
| assignIfDefined(out, "isDeeplinkable", c.isDeeplinkable); | ||
| assignIfDefined(out, "origin", c.origin); | ||
| assignIfDefined(out, "accountId", c.accountId); | ||
| assignIfDefined(out, "accountLabel", c.accountLabel); | ||
| if (c.summary) { | ||
| out.summary = { | ||
| text: c.summary.text, | ||
| recipe: c.summary.recipe, | ||
| method: c.summary.method, | ||
| createdAt: c.summary.createdAt | ||
| }; | ||
| } | ||
| return out; | ||
| } | ||
| function canonicalMessageDoc(m) { | ||
| const out = { | ||
| msgKey: m.msgKey, | ||
| docId: m.docId, | ||
| msgIndex: m.msgIndex, | ||
| role: m.role, | ||
| text: m.text | ||
| }; | ||
| assignIfDefined(out, "formattedText", m.formattedText); | ||
| out.createdAt = m.createdAt; | ||
| out.updatedAt = m.updatedAt; | ||
| out.contentHash = m.contentHash; | ||
| assignIfDefined(out, "platformMessageId", m.platformMessageId); | ||
| return out; | ||
| } | ||
| function assignIfDefined(target, key, value) { | ||
| if (value !== void 0) { | ||
| target[key] = value; | ||
| } | ||
| } | ||
| // src/corpus.ts | ||
| var DIR_MODE = 448; | ||
| var FILE_MODE = 384; | ||
| var gzip = promisify(gzipCallback); | ||
| var gunzip = promisify(gunzipCallback); | ||
| var Corpus = class { | ||
| paths; | ||
| // Serialize writes within this process so concurrent upserts (inbox drain + | ||
| // MCP save_conversation) don't interleave. Cross-process safety relies on the | ||
| // attended, single-companion v1 model (see PLAN-MCP.md open question 3). | ||
| writeChain = Promise.resolve(); | ||
| constructor(root) { | ||
| this.paths = new CorpusPaths(root); | ||
| } | ||
| /** Create the directory skeleton and meta.json if absent. Idempotent. */ | ||
| async init() { | ||
| await fs.mkdir(this.paths.conversationsDir, { recursive: true, mode: DIR_MODE }); | ||
| await fs.mkdir(this.paths.memoriesDir, { recursive: true, mode: DIR_MODE }); | ||
| await fs.mkdir(this.paths.processedDir, { recursive: true, mode: DIR_MODE }); | ||
| await fs.mkdir(this.paths.quarantineDir, { recursive: true, mode: DIR_MODE }); | ||
| await fs.mkdir(this.paths.indexDir, { recursive: true, mode: DIR_MODE }); | ||
| await hardenDir(this.paths.root); | ||
| await hardenDir(this.paths.conversationsDir); | ||
| await hardenDir(this.paths.memoriesDir); | ||
| await hardenDir(this.paths.inboxDir); | ||
| await hardenDir(this.paths.processedDir); | ||
| await hardenDir(this.paths.quarantineDir); | ||
| await hardenDir(this.paths.indexDir); | ||
| try { | ||
| await fs.access(this.paths.metaFile); | ||
| } catch { | ||
| const meta = { | ||
| formatVersion: CORPUS_FORMAT_VERSION, | ||
| createdAt: Date.now(), | ||
| lastIngestAt: null, | ||
| conversationCount: 0 | ||
| }; | ||
| await writeFileAtomic(this.paths.metaFile, `${JSON.stringify(meta, null, 2)} | ||
| `); | ||
| } | ||
| } | ||
| async readMeta() { | ||
| try { | ||
| const raw = await fs.readFile(this.paths.metaFile, "utf8"); | ||
| return JSON.parse(raw); | ||
| } catch { | ||
| return null; | ||
| } | ||
| } | ||
| async writeMeta(patch) { | ||
| const current = await this.readMeta() ?? { | ||
| formatVersion: CORPUS_FORMAT_VERSION, | ||
| createdAt: Date.now(), | ||
| lastIngestAt: null, | ||
| conversationCount: 0 | ||
| }; | ||
| const next = { ...current, ...patch }; | ||
| await writeFileAtomic(this.paths.metaFile, `${JSON.stringify(next, null, 2)} | ||
| `); | ||
| } | ||
| /** | ||
| * Insert or replace a conversation, keyed on docId. New bridge records use | ||
| * the local `corpusChangedAt` mutation clock; older records fall back to the | ||
| * source `updatedAt`. If mutation clocks tie, source freshness breaks the | ||
| * tie. Re-writing identical content remains a byte-identical no-op. | ||
| */ | ||
| async upsert(record) { | ||
| return this.enqueueWrite(async () => { | ||
| const canonical = canonicalizeConversation(record); | ||
| const { conversation } = canonical; | ||
| const file = this.paths.compressedConversationFile(conversation.platform, conversation.docId); | ||
| const legacyFile = this.paths.legacyConversationFile(conversation.platform, conversation.docId); | ||
| const existing = await this.readBestConversation(conversation.platform, conversation.docId); | ||
| if (existing) { | ||
| const existingChangedAt = existing.conversation.corpusChangedAt ?? existing.conversation.updatedAt ?? 0; | ||
| const incomingChangedAt = conversation.corpusChangedAt ?? conversation.updatedAt ?? 0; | ||
| const existingUpdatedAt = existing.conversation.updatedAt ?? 0; | ||
| const incomingUpdatedAt = conversation.updatedAt ?? 0; | ||
| if (incomingChangedAt < existingChangedAt || incomingChangedAt === existingChangedAt && incomingUpdatedAt < existingUpdatedAt) { | ||
| return { docId: conversation.docId, outcome: "skipped" }; | ||
| } | ||
| } | ||
| const serialized = serializeConversation(canonical); | ||
| if (existing) { | ||
| const existingSerialized = `${JSON.stringify(canonicalizeConversation(existing), null, 2)} | ||
| `; | ||
| if (existingSerialized === serialized) { | ||
| return { docId: conversation.docId, outcome: "skipped" }; | ||
| } | ||
| } | ||
| await fs.mkdir(dirname(file), { recursive: true, mode: DIR_MODE }); | ||
| await writeCompressedConversationAtomic(file, serialized); | ||
| await fs.rm(legacyFile, { force: true }); | ||
| return { docId: conversation.docId, outcome: existing ? "updated" : "created" }; | ||
| }); | ||
| } | ||
| /** Read a conversation by docId, scanning platform dirs for the slug file. */ | ||
| async get(docId) { | ||
| for (const platform of await this.listPlatforms()) { | ||
| const record = await this.readBestConversation(platform, docId); | ||
| if (record && record.conversation.docId === docId) { | ||
| return record; | ||
| } | ||
| } | ||
| return null; | ||
| } | ||
| /** Direct form used by search snippet hydration, avoiding a platform scan. */ | ||
| async getFromPlatform(platform, docId) { | ||
| const record = await this.readBestConversation(platform, docId); | ||
| return record?.conversation.docId === docId ? record : null; | ||
| } | ||
| /** | ||
| * Synchronous twin used after SQLite has ranked a small result set. The | ||
| * search connection is synchronous already (`node:sqlite` DatabaseSync), and | ||
| * reading at most `limit` local gzip files keeps the public retriever API | ||
| * stable while avoiding a second raw-text store inside FTS5. | ||
| */ | ||
| getFromPlatformSync(platform, docId) { | ||
| const record = readBestConversationFilesSync( | ||
| this.paths.compressedConversationFile(platform, docId), | ||
| this.paths.legacyConversationFile(platform, docId) | ||
| ); | ||
| return record?.conversation.docId === docId ? record : null; | ||
| } | ||
| /** | ||
| * Write an agent-authored memory, keyed on memoryId. | ||
| * | ||
| * The semantics differ from {@link upsert} on purpose. A memoryId is a | ||
| * content hash, so an existing file at that id already holds this exact | ||
| * memory: re-writing it would only churn the timestamps, and it would stamp | ||
| * `enabled: false` back over a memory the user had reviewed and turned on. | ||
| * So an existing record wins and the write is skipped. | ||
| * | ||
| * A record whose `source` is not "agent" is refused outright rather than | ||
| * skipped: a captured conversation or a memory the user wrote by hand must be | ||
| * untouchable from the agent write path, and the caller should hear about it | ||
| * rather than believe its write landed. | ||
| */ | ||
| async upsertMemory(memory) { | ||
| return this.enqueueWrite(async () => { | ||
| const canonical = canonicalizeMemory(memory); | ||
| const file = this.paths.memoryFile(canonical.memoryId); | ||
| const existing = await readJsonIfExists(file); | ||
| if (existing) { | ||
| if (existing.source !== "agent") { | ||
| return { memoryId: canonical.memoryId, outcome: "refused" }; | ||
| } | ||
| return { memoryId: canonical.memoryId, outcome: "skipped" }; | ||
| } | ||
| await fs.mkdir(dirname(file), { recursive: true, mode: DIR_MODE }); | ||
| await writeFileAtomic(file, `${JSON.stringify(canonical, null, 2)} | ||
| `); | ||
| return { memoryId: canonical.memoryId, outcome: "created" }; | ||
| }); | ||
| } | ||
| /** Read a memory by memoryId, or null when it has never been written. */ | ||
| async getMemory(memoryId) { | ||
| const record = await readJsonIfExists(this.paths.memoryFile(memoryId)); | ||
| return record && record.memoryId === memoryId ? record : null; | ||
| } | ||
| /** Async-iterate every stored memory (used by the index rebuild). */ | ||
| async *iterateMemories() { | ||
| let files; | ||
| try { | ||
| files = await fs.readdir(this.paths.memoriesDir); | ||
| } catch { | ||
| return; | ||
| } | ||
| for (const name of files.sort()) { | ||
| if (!name.endsWith(".json")) { | ||
| continue; | ||
| } | ||
| const record = await readJsonIfExists(join3(this.paths.memoriesDir, name)); | ||
| if (record?.memoryId) { | ||
| yield record; | ||
| } | ||
| } | ||
| } | ||
| /** Platform subdirectory names under conversations/. */ | ||
| async listPlatforms() { | ||
| try { | ||
| const entries = await fs.readdir(this.paths.conversationsDir, { withFileTypes: true }); | ||
| return entries.filter((e) => e.isDirectory()).map((e) => e.name); | ||
| } catch { | ||
| return []; | ||
| } | ||
| } | ||
| /** Async-iterate every stored conversation (used by reindex + stats). */ | ||
| async *iterate() { | ||
| for (const platform of await this.listPlatforms()) { | ||
| const dir = this.paths.platformDir(platform); | ||
| let files; | ||
| try { | ||
| files = await fs.readdir(dir); | ||
| } catch { | ||
| continue; | ||
| } | ||
| const stems = /* @__PURE__ */ new Set(); | ||
| for (const name of files) { | ||
| if (name.endsWith(".json.gz")) stems.add(name.slice(0, -".json.gz".length)); | ||
| else if (name.endsWith(".json")) stems.add(name.slice(0, -".json".length)); | ||
| } | ||
| for (const stem of [...stems].sort()) { | ||
| const record = await readBestConversationFiles( | ||
| join3(dir, `${stem}.json.gz`), | ||
| join3(dir, `${stem}.json`) | ||
| ); | ||
| if (record) { | ||
| yield record; | ||
| } | ||
| } | ||
| } | ||
| } | ||
| /** Byte/count breakdown used by `stats` and `compact --dry-run`. */ | ||
| async conversationStorageStats() { | ||
| const stats = { | ||
| legacyFiles: 0, | ||
| legacyBytes: 0, | ||
| compressedFiles: 0, | ||
| compressedBytes: 0 | ||
| }; | ||
| for (const platform of await this.listPlatforms()) { | ||
| const dir = this.paths.platformDir(platform); | ||
| let entries; | ||
| try { | ||
| entries = await fs.readdir(dir, { withFileTypes: true }); | ||
| } catch { | ||
| continue; | ||
| } | ||
| for (const entry of entries) { | ||
| if (!entry.isFile()) continue; | ||
| const file = join3(dir, entry.name); | ||
| if (entry.name.endsWith(".json.gz")) { | ||
| stats.compressedFiles += 1; | ||
| stats.compressedBytes += (await fs.stat(file)).size; | ||
| } else if (entry.name.endsWith(".json")) { | ||
| stats.legacyFiles += 1; | ||
| stats.legacyBytes += (await fs.stat(file)).size; | ||
| } | ||
| } | ||
| } | ||
| return stats; | ||
| } | ||
| /** | ||
| * Convert every legacy conversation independently. Each candidate is written, | ||
| * read back, and byte-compared in canonical JSON form before its plaintext | ||
| * source is removed, making the operation resumable after any interruption. | ||
| */ | ||
| async migrateLegacyConversations(options = {}) { | ||
| const dryRun = options.dryRun === true; | ||
| const before = await this.conversationStorageStats(); | ||
| const report = { | ||
| ...before, | ||
| dryRun, | ||
| candidates: before.legacyFiles, | ||
| migrated: 0, | ||
| legacyRemoved: 0, | ||
| projectedCompressedBytes: before.compressedBytes, | ||
| failures: [] | ||
| }; | ||
| for (const platform of await this.listPlatforms()) { | ||
| const dir = this.paths.platformDir(platform); | ||
| let names; | ||
| try { | ||
| names = (await fs.readdir(dir)).filter((name) => name.endsWith(".json")).sort(); | ||
| } catch { | ||
| continue; | ||
| } | ||
| for (const name of names) { | ||
| const legacyFile = join3(dir, name); | ||
| const compressedFile = join3(dir, `${name}.gz`); | ||
| try { | ||
| const legacy = await readPlainConversationIfExists(legacyFile); | ||
| if (!legacy) throw new Error("legacy JSON could not be parsed"); | ||
| const compressed = await readCompressedConversationIfExists(compressedFile); | ||
| const winner = pickNewerConversation(compressed, legacy) ?? legacy; | ||
| const serialized = serializeConversation(canonicalizeConversation(winner)); | ||
| const compressedBytes = await gzip(Buffer.from(serialized, "utf8"), { level: 9 }); | ||
| report.projectedCompressedBytes -= await fileSizeIfExists(compressedFile); | ||
| report.projectedCompressedBytes += compressedBytes.byteLength; | ||
| if (!dryRun) { | ||
| await writeVerifiedCompressedConversationAtomic(compressedFile, compressedBytes, serialized); | ||
| await fs.rm(legacyFile); | ||
| } | ||
| report.migrated += 1; | ||
| report.legacyRemoved += 1; | ||
| } catch (error) { | ||
| report.failures.push({ | ||
| file: legacyFile, | ||
| error: error instanceof Error ? error.message : String(error) | ||
| }); | ||
| } | ||
| } | ||
| } | ||
| if (!dryRun && report.failures.length === 0) { | ||
| await this.writeMeta({ formatVersion: CORPUS_FORMAT_VERSION }); | ||
| } | ||
| return report; | ||
| } | ||
| async stats() { | ||
| const byPlatform = {}; | ||
| let conversationCount = 0; | ||
| let messageCount = 0; | ||
| for await (const record of this.iterate()) { | ||
| conversationCount += 1; | ||
| messageCount += record.messages.length; | ||
| const platform = safeSegment(record.conversation.platform); | ||
| byPlatform[platform] = (byPlatform[platform] ?? 0) + 1; | ||
| } | ||
| return { conversationCount, messageCount, byPlatform }; | ||
| } | ||
| /** Refresh meta.json's conversationCount / lastIngestAt after an ingest run. */ | ||
| async touchIngestMeta(receipt) { | ||
| const { conversationCount } = await this.stats(); | ||
| await this.writeMeta({ | ||
| lastIngestAt: Date.now(), | ||
| conversationCount, | ||
| ...receipt ? { lastIngestReceipt: receipt } : {} | ||
| }); | ||
| } | ||
| enqueueWrite(fn) { | ||
| const run = this.writeChain.then(fn, fn); | ||
| this.writeChain = run.then( | ||
| () => void 0, | ||
| () => void 0 | ||
| ); | ||
| return run; | ||
| } | ||
| readBestConversation(platform, docId) { | ||
| return readBestConversationFiles( | ||
| this.paths.compressedConversationFile(platform, docId), | ||
| this.paths.legacyConversationFile(platform, docId) | ||
| ); | ||
| } | ||
| }; | ||
| async function readJsonIfExists(file) { | ||
| try { | ||
| const raw = await fs.readFile(file, "utf8"); | ||
| return JSON.parse(raw); | ||
| } catch { | ||
| return null; | ||
| } | ||
| } | ||
| function serializeConversation(record) { | ||
| return `${JSON.stringify(canonicalizeConversation(record), null, 2)} | ||
| `; | ||
| } | ||
| async function readPlainConversationIfExists(file) { | ||
| return readJsonIfExists(file); | ||
| } | ||
| async function readCompressedConversationIfExists(file) { | ||
| try { | ||
| const raw = await fs.readFile(file); | ||
| return JSON.parse((await gunzip(raw)).toString("utf8")); | ||
| } catch { | ||
| return null; | ||
| } | ||
| } | ||
| async function readBestConversationFiles(compressedFile, legacyFile) { | ||
| const [compressed, legacy] = await Promise.all([ | ||
| readCompressedConversationIfExists(compressedFile), | ||
| readPlainConversationIfExists(legacyFile) | ||
| ]); | ||
| return pickNewerConversation(compressed, legacy); | ||
| } | ||
| function readBestConversationFilesSync(compressedFile, legacyFile) { | ||
| let compressed = null; | ||
| let legacy = null; | ||
| try { | ||
| compressed = JSON.parse(gunzipSync(readFileSync2(compressedFile)).toString("utf8")); | ||
| } catch { | ||
| } | ||
| try { | ||
| legacy = JSON.parse(readFileSync2(legacyFile, "utf8")); | ||
| } catch { | ||
| } | ||
| return pickNewerConversation(compressed, legacy); | ||
| } | ||
| function pickNewerConversation(a, b) { | ||
| if (!a) return b; | ||
| if (!b) return a; | ||
| const aChanged = a.conversation.corpusChangedAt ?? a.conversation.updatedAt ?? 0; | ||
| const bChanged = b.conversation.corpusChangedAt ?? b.conversation.updatedAt ?? 0; | ||
| if (aChanged !== bChanged) return aChanged > bChanged ? a : b; | ||
| const aUpdated = a.conversation.updatedAt ?? 0; | ||
| const bUpdated = b.conversation.updatedAt ?? 0; | ||
| return aUpdated >= bUpdated ? a : b; | ||
| } | ||
| async function writeCompressedConversationAtomic(file, serialized) { | ||
| const compressed = await gzip(Buffer.from(serialized, "utf8"), { level: 9 }); | ||
| await writeFileAtomic(file, compressed); | ||
| } | ||
| async function writeVerifiedCompressedConversationAtomic(file, compressed, serialized) { | ||
| const tmp = `${file}.migration-${process.pid}-${Date.now()}`; | ||
| try { | ||
| await fs.writeFile(tmp, compressed, { mode: FILE_MODE }); | ||
| const verified = await readCompressedConversationIfExists(tmp); | ||
| if (!verified || serializeConversation(canonicalizeConversation(verified)) !== serialized) { | ||
| throw new Error("compressed read-back verification failed"); | ||
| } | ||
| await fs.rename(tmp, file); | ||
| } catch (error) { | ||
| await fs.rm(tmp, { force: true }); | ||
| throw error; | ||
| } | ||
| } | ||
| async function fileSizeIfExists(file) { | ||
| try { | ||
| return (await fs.stat(file)).size; | ||
| } catch { | ||
| return 0; | ||
| } | ||
| } | ||
| async function writeFileAtomic(file, contents) { | ||
| const tmp = `${file}.tmp-${process.pid}-${Date.now()}`; | ||
| if (typeof contents === "string") { | ||
| await fs.writeFile(tmp, contents, { encoding: "utf8", mode: FILE_MODE }); | ||
| } else { | ||
| await fs.writeFile(tmp, contents, { mode: FILE_MODE }); | ||
| } | ||
| await fs.rename(tmp, file); | ||
| } | ||
| async function hardenDir(dir) { | ||
| try { | ||
| await fs.chmod(dir, DIR_MODE); | ||
| } catch { | ||
| } | ||
| } | ||
| // src/ingest.ts | ||
| import { promises as fs2, createReadStream } from "fs"; | ||
| import { join as join4 } from "path"; | ||
| import { createInterface } from "readline"; | ||
| import { createGunzip } from "zlib"; | ||
| var EMPTY_REPORT = { | ||
| conversations: 0, | ||
| created: 0, | ||
| updated: 0, | ||
| skipped: 0, | ||
| messages: 0, | ||
| malformedLines: 0 | ||
| }; | ||
| async function ingestFile(corpus, file, index) { | ||
| const report = { ...EMPTY_REPORT }; | ||
| let current = null; | ||
| let buffered = []; | ||
| const flush = async () => { | ||
| if (!current) { | ||
| return; | ||
| } | ||
| const record = { conversation: current, messages: buffered }; | ||
| const result = await corpus.upsert(record); | ||
| tally(report, result); | ||
| if (index && result.outcome !== "skipped") { | ||
| index.upsertConversation(record); | ||
| } | ||
| report.conversations += 1; | ||
| report.messages += buffered.length; | ||
| current = null; | ||
| buffered = []; | ||
| }; | ||
| const rl = createInterface({ input: openLineSource(file), crlfDelay: Infinity }); | ||
| for await (const raw of rl) { | ||
| const line = raw.trim(); | ||
| if (!line) { | ||
| continue; | ||
| } | ||
| let parsed; | ||
| try { | ||
| parsed = JSON.parse(line); | ||
| } catch { | ||
| report.malformedLines += 1; | ||
| continue; | ||
| } | ||
| if (typeof parsed !== "object" || parsed === null) { | ||
| report.malformedLines += 1; | ||
| continue; | ||
| } | ||
| const record = parsed; | ||
| const type = record.type; | ||
| if (type === "meta") { | ||
| continue; | ||
| } | ||
| if (type === "conversation") { | ||
| await flush(); | ||
| current = coerceConversation(record.conversation); | ||
| buffered = []; | ||
| if (!current) { | ||
| report.malformedLines += 1; | ||
| } | ||
| continue; | ||
| } | ||
| if (type === "message" && current) { | ||
| const docId = typeof record.docId === "string" ? record.docId : ""; | ||
| if (docId !== current.docId) { | ||
| report.malformedLines += 1; | ||
| continue; | ||
| } | ||
| const message = coerceMessage(current.docId, record.message); | ||
| if (message) { | ||
| buffered.push(message); | ||
| } else { | ||
| report.malformedLines += 1; | ||
| } | ||
| } | ||
| } | ||
| await flush(); | ||
| return report; | ||
| } | ||
| async function pendingInboxFiles(corpus) { | ||
| try { | ||
| return (await fs2.readdir(corpus.paths.inboxDir, { withFileTypes: true })).filter((e) => e.isFile() && isNdjson(e.name)).map((e) => e.name).sort(); | ||
| } catch { | ||
| return []; | ||
| } | ||
| } | ||
| async function drainInbox(corpus, index, options = {}) { | ||
| await corpus.init(); | ||
| const inbox = corpus.paths.inboxDir; | ||
| const entries = await pendingInboxFiles(corpus); | ||
| const total = { | ||
| ...EMPTY_REPORT, | ||
| files: 0, | ||
| failed: 0, | ||
| discarded: 0, | ||
| quarantined: 0 | ||
| }; | ||
| for (const name of entries) { | ||
| const src = join4(inbox, name); | ||
| try { | ||
| const report = await ingestFile(corpus, src, index); | ||
| accumulate(total, report); | ||
| total.files += 1; | ||
| if (report.malformedLines > 0) { | ||
| await moveToQuarantine(corpus, src, name); | ||
| total.quarantined += 1; | ||
| } else { | ||
| await fs2.unlink(src); | ||
| total.discarded += 1; | ||
| } | ||
| } catch (error) { | ||
| if (await stillPending(src)) { | ||
| total.failed += 1; | ||
| if (isCorruptDeltaError(error)) { | ||
| await moveToQuarantine(corpus, src, name); | ||
| total.quarantined += 1; | ||
| } | ||
| options.onFileError?.(name, error); | ||
| } | ||
| } | ||
| } | ||
| if (total.files > 0) { | ||
| await corpus.touchIngestMeta({ | ||
| files: total.files, | ||
| conversations: total.conversations, | ||
| messages: total.messages, | ||
| malformedLines: total.malformedLines, | ||
| discarded: total.discarded, | ||
| quarantined: total.quarantined | ||
| }); | ||
| } | ||
| return total; | ||
| } | ||
| async function stillPending(src) { | ||
| try { | ||
| await fs2.access(src); | ||
| return true; | ||
| } catch { | ||
| return false; | ||
| } | ||
| } | ||
| async function moveToQuarantine(corpus, src, name) { | ||
| const dest = join4(corpus.paths.quarantineDir, name); | ||
| try { | ||
| await fs2.rename(src, dest); | ||
| } catch { | ||
| const alt = join4(corpus.paths.quarantineDir, `${Date.now()}-${name}`); | ||
| await fs2.copyFile(src, alt); | ||
| await fs2.unlink(src); | ||
| } | ||
| } | ||
| function isCorruptDeltaError(error) { | ||
| if (!error || typeof error !== "object" || !("code" in error)) { | ||
| return false; | ||
| } | ||
| const code = String(error.code ?? ""); | ||
| return code === "Z_DATA_ERROR" || code === "Z_BUF_ERROR" || code.startsWith("ERR_ZLIB_"); | ||
| } | ||
| function isNdjson(name) { | ||
| const lower = name.toLowerCase(); | ||
| return lower.endsWith(".jsonl") || lower.endsWith(".ndjson") || lower.endsWith(".jsonl.gz") || lower.endsWith(".ndjson.gz"); | ||
| } | ||
| function openLineSource(file) { | ||
| const stream = createReadStream(file); | ||
| if (file.toLowerCase().endsWith(".gz")) { | ||
| return stream.pipe(createGunzip()); | ||
| } | ||
| return stream; | ||
| } | ||
| function coerceConversation(value) { | ||
| if (!value || typeof value !== "object") { | ||
| return null; | ||
| } | ||
| const c = value; | ||
| if (typeof c.docId !== "string" || !c.docId) { | ||
| return null; | ||
| } | ||
| const conversation = { ...c }; | ||
| if (typeof conversation.updatedAt !== "number") { | ||
| conversation.updatedAt = typeof conversation.contentUpdatedAt === "number" ? conversation.contentUpdatedAt : 0; | ||
| } | ||
| return conversation; | ||
| } | ||
| function coerceMessage(docId, value) { | ||
| if (!value || typeof value !== "object") { | ||
| return null; | ||
| } | ||
| const m = value; | ||
| const text = typeof m.text === "string" ? m.text : ""; | ||
| if (!text.trim()) { | ||
| return null; | ||
| } | ||
| const msgIndex = typeof m.msgIndex === "number" ? m.msgIndex : 0; | ||
| const message = { | ||
| // The backup line omits msgKey; reconstruct the extension's convention. | ||
| msgKey: `${docId}:${msgIndex}`, | ||
| docId, | ||
| msgIndex, | ||
| role: coerceRole(m.role), | ||
| text, | ||
| createdAt: typeof m.createdAt === "number" ? m.createdAt : 0, | ||
| updatedAt: typeof m.updatedAt === "number" ? m.updatedAt : 0, | ||
| contentHash: typeof m.contentHash === "string" ? m.contentHash : "" | ||
| }; | ||
| if (typeof m.formattedText === "string") { | ||
| message.formattedText = m.formattedText; | ||
| } | ||
| if (typeof m.platformMessageId === "string") { | ||
| message.platformMessageId = m.platformMessageId; | ||
| } else if (m.platformMessageId === null) { | ||
| message.platformMessageId = null; | ||
| } | ||
| return message; | ||
| } | ||
| function coerceRole(value) { | ||
| return value === "assistant" || value === "user" || value === "system" || value === "unknown" ? value : "unknown"; | ||
| } | ||
| function tally(report, result) { | ||
| if (result.outcome === "created") { | ||
| report.created += 1; | ||
| } else if (result.outcome === "updated") { | ||
| report.updated += 1; | ||
| } else { | ||
| report.skipped += 1; | ||
| } | ||
| } | ||
| function accumulate(total, report) { | ||
| total.conversations += report.conversations; | ||
| total.created += report.created; | ||
| total.updated += report.updated; | ||
| total.skipped += report.skipped; | ||
| total.messages += report.messages; | ||
| total.malformedLines += report.malformedLines; | ||
| } | ||
| // src/searchIndex.ts | ||
| import { promises as fs3, chmodSync, closeSync, openSync } from "fs"; | ||
| import { createRequire } from "module"; | ||
| // src/sqliteWarning.ts | ||
| var INSTALLED = /* @__PURE__ */ Symbol.for("llmnesia.sqliteWarningFilterInstalled"); | ||
| function isSqliteExperimentalWarning(warning) { | ||
| return warning.name === "ExperimentalWarning" && /\bSQLite\b/i.test(warning.message); | ||
| } | ||
| function install() { | ||
| const flagged = globalThis; | ||
| if (flagged[INSTALLED]) { | ||
| return; | ||
| } | ||
| flagged[INSTALLED] = true; | ||
| const previous = process.listeners("warning"); | ||
| process.removeAllListeners("warning"); | ||
| process.on("warning", (warning) => { | ||
| if (isSqliteExperimentalWarning(warning)) { | ||
| return; | ||
| } | ||
| for (const listener of previous) { | ||
| listener.call(process, warning); | ||
| } | ||
| }); | ||
| } | ||
| install(); | ||
| // src/searchIndex.ts | ||
| var nodeRequire = createRequire(import.meta.url); | ||
| var sqlite; | ||
| function loadSqlite() { | ||
| return sqlite ??= nodeRequire("node:sqlite"); | ||
| } | ||
| var DB_FILE_MODE = 384; | ||
| function hardenDbFile(dbPath) { | ||
| try { | ||
| closeSync(openSync(dbPath, "a", DB_FILE_MODE)); | ||
| chmodSync(dbPath, DB_FILE_MODE); | ||
| } catch { | ||
| } | ||
| } | ||
| var TITLE_ROW_INDEX = -1; | ||
| var TITLE_ROLE = "title"; | ||
| var DEFAULT_LIMIT = 10; | ||
| var SNIPPET_TOKENS = 12; | ||
| var META_SQL = ` | ||
| CREATE TABLE IF NOT EXISTS index_meta ( | ||
| key TEXT PRIMARY KEY, | ||
| value TEXT NOT NULL | ||
| ); | ||
| `; | ||
| var DROP_DERIVED_SQL = ` | ||
| DROP TABLE IF EXISTS conversations; | ||
| DROP TABLE IF EXISTS messages_fts; | ||
| DROP TABLE IF EXISTS message_rows; | ||
| DROP TABLE IF EXISTS memories; | ||
| DROP TABLE IF EXISTS memories_fts; | ||
| `; | ||
| var SCHEMA_SQL = ` | ||
| CREATE TABLE IF NOT EXISTS index_meta ( | ||
| key TEXT PRIMARY KEY, | ||
| value TEXT NOT NULL | ||
| ); | ||
| CREATE TABLE IF NOT EXISTS conversations ( | ||
| docId TEXT PRIMARY KEY, | ||
| platform TEXT NOT NULL, | ||
| title TEXT NOT NULL, | ||
| url TEXT NOT NULL, | ||
| createdAt INTEGER NOT NULL, | ||
| updatedAt INTEGER NOT NULL, | ||
| messageCount INTEGER NOT NULL, | ||
| preview TEXT NOT NULL, | ||
| pinned INTEGER NOT NULL, | ||
| accountId TEXT, | ||
| accountLabel TEXT | ||
| ); | ||
| CREATE INDEX IF NOT EXISTS idx_conversations_platform ON conversations(platform); | ||
| CREATE INDEX IF NOT EXISTS idx_conversations_createdAt ON conversations(createdAt); | ||
| CREATE TABLE IF NOT EXISTS message_rows ( | ||
| rowid INTEGER PRIMARY KEY, | ||
| docId TEXT NOT NULL, | ||
| msgIndex INTEGER NOT NULL, | ||
| role TEXT NOT NULL | ||
| ); | ||
| CREATE INDEX IF NOT EXISTS idx_message_rows_docId ON message_rows(docId); | ||
| CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts USING fts5( | ||
| text, | ||
| content = '', | ||
| contentless_delete = 1, | ||
| tokenize = 'porter unicode61' | ||
| ); | ||
| CREATE TABLE IF NOT EXISTS memories ( | ||
| memoryId TEXT PRIMARY KEY, | ||
| kind TEXT NOT NULL, | ||
| label TEXT NOT NULL, | ||
| text TEXT NOT NULL, | ||
| source TEXT NOT NULL, | ||
| pinned INTEGER NOT NULL, | ||
| enabled INTEGER NOT NULL, | ||
| createdAt INTEGER NOT NULL, | ||
| updatedAt INTEGER NOT NULL | ||
| ); | ||
| CREATE VIRTUAL TABLE IF NOT EXISTS memories_fts USING fts5( | ||
| memoryId UNINDEXED, | ||
| label, | ||
| text, | ||
| tokenize = 'porter unicode61' | ||
| ); | ||
| `; | ||
| var SearchIndex = class { | ||
| db; | ||
| stmts; | ||
| /** Whether the on-disk index was a stale format version when it was opened. */ | ||
| staleOnOpen; | ||
| /** Canonical source used only to hydrate the few snippets returned by rank. */ | ||
| corpus; | ||
| constructor(dbPath, corpus) { | ||
| this.corpus = corpus; | ||
| hardenDbFile(dbPath); | ||
| this.db = new (loadSqlite()).DatabaseSync(dbPath); | ||
| this.db.exec("PRAGMA journal_mode = WAL"); | ||
| this.db.exec("PRAGMA foreign_keys = OFF"); | ||
| this.db.exec(META_SQL); | ||
| this.staleOnOpen = this.getMeta("format_version") !== String(INDEX_FORMAT_VERSION); | ||
| if (this.staleOnOpen) { | ||
| this.db.exec(DROP_DERIVED_SQL); | ||
| } | ||
| this.db.exec(SCHEMA_SQL); | ||
| this.db.exec("PRAGMA temp_store = MEMORY"); | ||
| this.db.exec(`CREATE VIRTUAL TABLE IF NOT EXISTS temp.snippet_fts USING fts5( | ||
| text, tokenize = 'porter unicode61' | ||
| )`); | ||
| this.setMeta("format_version", String(INDEX_FORMAT_VERSION)); | ||
| this.stmts = { | ||
| selectFtsRows: this.db.prepare("SELECT rowid FROM message_rows WHERE docId = ?"), | ||
| deleteFts: this.db.prepare("DELETE FROM messages_fts WHERE rowid = ?"), | ||
| deleteMessageRows: this.db.prepare("DELETE FROM message_rows WHERE docId = ?"), | ||
| insertMessageRow: this.db.prepare( | ||
| "INSERT INTO message_rows (docId, msgIndex, role) VALUES (?, ?, ?)" | ||
| ), | ||
| insertFts: this.db.prepare("INSERT INTO messages_fts (rowid, text) VALUES (?, ?)"), | ||
| deleteConversation: this.db.prepare("DELETE FROM conversations WHERE docId = ?"), | ||
| insertConversation: this.db.prepare( | ||
| `INSERT INTO conversations | ||
| (docId, platform, title, url, createdAt, updatedAt, messageCount, preview, pinned, accountId, accountLabel) | ||
| VALUES (@docId, @platform, @title, @url, @createdAt, @updatedAt, @messageCount, @preview, @pinned, @accountId, @accountLabel)` | ||
| ), | ||
| deleteMemoryFts: this.db.prepare("DELETE FROM memories_fts WHERE memoryId = ?"), | ||
| insertMemoryFts: this.db.prepare( | ||
| "INSERT INTO memories_fts (memoryId, label, text) VALUES (?, ?, ?)" | ||
| ), | ||
| deleteMemory: this.db.prepare("DELETE FROM memories WHERE memoryId = ?"), | ||
| insertMemory: this.db.prepare( | ||
| `INSERT INTO memories | ||
| (memoryId, kind, label, text, source, pinned, enabled, createdAt, updatedAt) | ||
| VALUES (@memoryId, @kind, @label, @text, @source, @pinned, @enabled, @createdAt, @updatedAt)` | ||
| ), | ||
| clearSnippetFts: this.db.prepare("DELETE FROM temp.snippet_fts"), | ||
| insertSnippetFts: this.db.prepare("INSERT INTO temp.snippet_fts (text) VALUES (?)"), | ||
| renderSnippet: this.db.prepare( | ||
| `SELECT snippet(snippet_fts, 0, '[', ']', '\u2026', ${SNIPPET_TOKENS}) AS snippet | ||
| FROM temp.snippet_fts WHERE snippet_fts MATCH ?` | ||
| ) | ||
| }; | ||
| } | ||
| /** True when the stored format version was missing or stale — caller rebuilds. */ | ||
| needsRebuild() { | ||
| return this.staleOnOpen; | ||
| } | ||
| /** Number of indexed conversations. */ | ||
| count() { | ||
| const row = this.db.prepare("SELECT COUNT(*) AS n FROM conversations").get(); | ||
| return row.n; | ||
| } | ||
| /** Number of indexed memories. */ | ||
| countMemories() { | ||
| const row = this.db.prepare("SELECT COUNT(*) AS n FROM memories").get(); | ||
| return row.n; | ||
| } | ||
| /** Most-recently-updated conversations (metadata only), newest first. */ | ||
| listRecent(n) { | ||
| return this.listConversations(n, "recent"); | ||
| } | ||
| /** Conversation metadata in a caller-selected, deterministic chronology. */ | ||
| listConversations(n, sort) { | ||
| const limit = Math.max(1, n); | ||
| const orderBy = sort === "oldest" ? "createdAt ASC, updatedAt ASC" : sort === "newest" ? "createdAt DESC, updatedAt DESC" : "updatedAt DESC, createdAt DESC"; | ||
| const rows = this.db.prepare( | ||
| `SELECT docId, platform, title, url, createdAt, updatedAt, messageCount, preview, pinned, | ||
| accountId, accountLabel | ||
| FROM conversations ORDER BY ${orderBy} LIMIT ?` | ||
| ).all(limit); | ||
| return rows.map((r) => { | ||
| const { accountId, accountLabel, ...rest } = r; | ||
| const entry = { ...rest, pinned: r.pinned === 1 }; | ||
| if (accountId) entry.accountId = accountId; | ||
| if (accountLabel) entry.accountLabel = accountLabel; | ||
| return entry; | ||
| }); | ||
| } | ||
| /** | ||
| * Run `fn` inside a transaction, rolling back if it throws. node:sqlite has | ||
| * no transaction() wrapper, so this is the one place BEGIN/COMMIT/ROLLBACK | ||
| * lives for synchronous writes; {@link rebuildFrom} spells it out separately | ||
| * because its body is async and cannot be expressed as a sync callback. | ||
| */ | ||
| transaction(fn) { | ||
| this.db.exec("BEGIN"); | ||
| try { | ||
| const result = fn(); | ||
| this.db.exec("COMMIT"); | ||
| return result; | ||
| } catch (err) { | ||
| this.db.exec("ROLLBACK"); | ||
| throw err; | ||
| } | ||
| } | ||
| upsertConversation(record) { | ||
| this.transaction(() => this.writeRecord(record)); | ||
| } | ||
| upsertMemory(memory) { | ||
| this.transaction(() => this.writeMemory(memory)); | ||
| } | ||
| removeConversation(docId) { | ||
| this.transaction(() => { | ||
| this.deleteFtsRows(docId); | ||
| this.stmts.deleteConversation.run(docId); | ||
| }); | ||
| } | ||
| /** | ||
| * Clear the index and repopulate it from the corpus — proving the corpus is | ||
| * the source of truth. Uses one manual transaction spanning the async file | ||
| * reads, which the sync {@link transaction} helper cannot hold. | ||
| */ | ||
| async rebuildFrom(corpus) { | ||
| this.corpus = corpus; | ||
| this.db.exec("DELETE FROM messages_fts; DELETE FROM message_rows; DELETE FROM conversations;"); | ||
| this.db.exec("DELETE FROM memories_fts; DELETE FROM memories;"); | ||
| let conversations = 0; | ||
| let memories = 0; | ||
| this.db.exec("BEGIN"); | ||
| try { | ||
| for await (const record of corpus.iterate()) { | ||
| this.writeRecord(record); | ||
| conversations += 1; | ||
| } | ||
| for await (const memory of corpus.iterateMemories()) { | ||
| this.writeMemory(memory); | ||
| memories += 1; | ||
| } | ||
| this.db.exec("COMMIT"); | ||
| } catch (err) { | ||
| this.db.exec("ROLLBACK"); | ||
| throw err; | ||
| } | ||
| return { conversations, memories }; | ||
| } | ||
| // Replace a conversation's index rows. Caller supplies the transaction (the | ||
| // transaction() helper for single upserts, or rebuildFrom's manual one). | ||
| writeRecord(rec) { | ||
| const { conversation, messages } = rec; | ||
| this.deleteFtsRows(conversation.docId); | ||
| this.stmts.deleteConversation.run(conversation.docId); | ||
| const titleText = [conversation.title, conversation.preview].filter(Boolean).join("\n"); | ||
| if (titleText.trim()) { | ||
| this.insertFtsRow(conversation.docId, TITLE_ROW_INDEX, TITLE_ROLE, titleText); | ||
| } | ||
| for (const message of messages) { | ||
| if (message.text.trim()) { | ||
| this.insertFtsRow(conversation.docId, message.msgIndex, message.role, message.text); | ||
| } | ||
| } | ||
| this.stmts.insertConversation.run({ | ||
| docId: conversation.docId, | ||
| platform: conversation.platform, | ||
| title: conversation.title, | ||
| url: conversation.url, | ||
| createdAt: conversation.createdAt, | ||
| updatedAt: conversation.updatedAt ?? conversation.contentUpdatedAt ?? 0, | ||
| messageCount: conversation.messageCount ?? messages.length, | ||
| preview: conversation.preview ?? "", | ||
| pinned: conversation.pinned ? 1 : 0, | ||
| accountId: conversation.accountId ?? null, | ||
| accountLabel: conversation.accountLabel ?? null | ||
| }); | ||
| } | ||
| insertFtsRow(docId, msgIndex, role, text) { | ||
| const inserted = this.stmts.insertMessageRow.run(docId, msgIndex, role); | ||
| this.stmts.insertFts.run(inserted.lastInsertRowid, text); | ||
| } | ||
| deleteFtsRows(docId) { | ||
| const rows = this.stmts.selectFtsRows.all(docId); | ||
| for (const row of rows) { | ||
| this.stmts.deleteFts.run(row.rowid); | ||
| } | ||
| this.stmts.deleteMessageRows.run(docId); | ||
| } | ||
| // Replace a memory's index rows. Caller supplies the transaction, as with | ||
| // writeRecord. Memories live in their own tables rather than in messages_fts: | ||
| // a memory is not a turn in a conversation, and mixing them would let a | ||
| // stored note answer a question about what the user actually said. | ||
| writeMemory(memory) { | ||
| this.stmts.deleteMemoryFts.run(memory.memoryId); | ||
| this.stmts.deleteMemory.run(memory.memoryId); | ||
| const searchable = [memory.label, memory.text].filter(Boolean).join("\n"); | ||
| if (searchable.trim()) { | ||
| this.stmts.insertMemoryFts.run(memory.memoryId, memory.label, memory.text); | ||
| } | ||
| this.stmts.insertMemory.run({ | ||
| memoryId: memory.memoryId, | ||
| kind: memory.kind, | ||
| label: memory.label, | ||
| text: memory.text, | ||
| source: memory.source, | ||
| pinned: memory.pinned ? 1 : 0, | ||
| enabled: memory.enabled ? 1 : 0, | ||
| createdAt: memory.createdAt, | ||
| updatedAt: memory.updatedAt | ||
| }); | ||
| } | ||
| search(query, filters = {}) { | ||
| const match = toMatchQuery(query, filters.matchMode); | ||
| if (!match) { | ||
| return []; | ||
| } | ||
| const limit = Math.max(1, filters.limit ?? DEFAULT_LIMIT); | ||
| const clauses = ["messages_fts MATCH @match"]; | ||
| const params = { match }; | ||
| if (filters.platform) { | ||
| clauses.push("c.platform = @platform"); | ||
| params.platform = filters.platform; | ||
| } | ||
| if (filters.title) { | ||
| clauses.push("c.title LIKE @title ESCAPE '\\'"); | ||
| params.title = `%${escapeLike(filters.title)}%`; | ||
| } | ||
| if (typeof filters.dateFrom === "number") { | ||
| clauses.push("c.createdAt >= @dateFrom"); | ||
| params.dateFrom = filters.dateFrom; | ||
| } | ||
| if (typeof filters.dateTo === "number") { | ||
| clauses.push("c.createdAt <= @dateTo"); | ||
| params.dateTo = filters.dateTo; | ||
| } | ||
| const scanCap = Math.min(2e3, Math.max(200, limit * 40)); | ||
| const sql = ` | ||
| SELECT | ||
| c.docId, c.platform, c.title, c.url, c.createdAt, c.updatedAt, | ||
| c.messageCount, c.preview, c.pinned, c.accountId, c.accountLabel, | ||
| m.msgIndex AS msgIndex, m.role AS role, | ||
| bm25(messages_fts) AS bm25 | ||
| FROM messages_fts | ||
| JOIN message_rows m ON m.rowid = messages_fts.rowid | ||
| JOIN conversations c ON c.docId = m.docId | ||
| WHERE ${clauses.join(" AND ")} | ||
| ORDER BY bm25 | ||
| LIMIT @scanCap | ||
| `; | ||
| const rows = this.db.prepare(sql).all({ ...params, scanCap }); | ||
| const seen = /* @__PURE__ */ new Set(); | ||
| const hits = []; | ||
| for (const row of rows) { | ||
| if (seen.has(row.docId)) { | ||
| continue; | ||
| } | ||
| seen.add(row.docId); | ||
| const hit = { | ||
| docId: row.docId, | ||
| platform: row.platform, | ||
| title: row.title, | ||
| url: row.url, | ||
| createdAt: row.createdAt, | ||
| updatedAt: row.updatedAt, | ||
| messageCount: row.messageCount, | ||
| preview: row.preview, | ||
| pinned: row.pinned === 1, | ||
| // BM25 is negative (more negative = better); expose a positive relevance. | ||
| score: -row.bm25, | ||
| snippet: this.hydrateSnippet(row, match), | ||
| matchRole: normalizeMatchRole(row.role), | ||
| matchMsgIndex: row.msgIndex | ||
| }; | ||
| if (row.accountId) hit.accountId = row.accountId; | ||
| if (row.accountLabel) hit.accountLabel = row.accountLabel; | ||
| hits.push(hit); | ||
| if (hits.length >= limit) { | ||
| break; | ||
| } | ||
| } | ||
| return hits; | ||
| } | ||
| /** Reclaim free pages and merge FTS segments after a confirmed compaction. */ | ||
| compact() { | ||
| this.db.exec("INSERT INTO messages_fts(messages_fts) VALUES('optimize')"); | ||
| this.db.exec("INSERT INTO memories_fts(memories_fts) VALUES('optimize')"); | ||
| this.db.exec("PRAGMA wal_checkpoint(TRUNCATE)"); | ||
| this.db.exec("VACUUM"); | ||
| } | ||
| hydrateSnippet(row, match) { | ||
| const record = this.corpus?.getFromPlatformSync(row.platform, row.docId); | ||
| const source = row.msgIndex === TITLE_ROW_INDEX ? [record?.conversation.title ?? row.title, record?.conversation.preview ?? row.preview].filter(Boolean).join("\n") : record?.messages.find((message) => message.msgIndex === row.msgIndex)?.text ?? row.preview; | ||
| this.stmts.clearSnippetFts.run(); | ||
| try { | ||
| this.stmts.insertSnippetFts.run(source); | ||
| const rendered = this.stmts.renderSnippet.get(match); | ||
| return rendered?.snippet ?? fallbackSnippet(source, SNIPPET_TOKENS); | ||
| } finally { | ||
| this.stmts.clearSnippetFts.run(); | ||
| } | ||
| } | ||
| close() { | ||
| this.db.close(); | ||
| } | ||
| getMeta(key) { | ||
| const row = this.db.prepare("SELECT value FROM index_meta WHERE key = ?").get(key); | ||
| return row ? row.value : null; | ||
| } | ||
| setMeta(key, value) { | ||
| this.db.prepare("INSERT INTO index_meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value").run(key, value); | ||
| } | ||
| }; | ||
| async function openSearchIndex(corpus, options = {}) { | ||
| await corpus.init(); | ||
| const index = new SearchIndex(corpus.paths.indexDbFile, corpus); | ||
| if (options.autoBuild !== false && (index.needsRebuild() || index.count() === 0 && index.countMemories() === 0)) { | ||
| await index.rebuildFrom(corpus); | ||
| } | ||
| return index; | ||
| } | ||
| async function deleteIndexDb(dbFile) { | ||
| for (const suffix of ["", "-wal", "-shm"]) { | ||
| await fs3.rm(`${dbFile}${suffix}`, { force: true }); | ||
| } | ||
| } | ||
| function toMatchQuery(query, mode = "any") { | ||
| const terms = queryTerms(query); | ||
| if (terms.length === 0) { | ||
| return null; | ||
| } | ||
| if (mode === "phrase") { | ||
| return `"${terms.join(" ")}"`; | ||
| } | ||
| return terms.map((t) => `"${t}"`).join(mode === "all" ? " AND " : " OR "); | ||
| } | ||
| function queryTerms(query) { | ||
| return query.split(/[^\p{L}\p{N}]+/u).filter((term) => term.length > 0); | ||
| } | ||
| function fallbackSnippet(source, maxTokens) { | ||
| const collapsed = source.replace(/\s+/g, " ").trim(); | ||
| if (!collapsed) return ""; | ||
| const words = collapsed.split(" "); | ||
| const width = Math.max(1, maxTokens); | ||
| const end = Math.min(words.length, width); | ||
| return `${words.slice(0, end).join(" ")}${end < words.length ? "\u2026" : ""}`; | ||
| } | ||
| function escapeLike(value) { | ||
| return value.replace(/[\\%_]/g, "\\$&"); | ||
| } | ||
| function normalizeMatchRole(role) { | ||
| if (role === "title") { | ||
| return "title"; | ||
| } | ||
| return role === "assistant" || role === "user" || role === "system" ? role : "unknown"; | ||
| } | ||
| // src/runtimeInstall.ts | ||
| import { cp, mkdir, readdir as readdir2, rename, rm, writeFile } from "fs/promises"; | ||
| import { existsSync } from "fs"; | ||
| import { homedir as homedir2 } from "os"; | ||
| import { dirname as dirname2, join as join5 } from "path"; | ||
| import { fileURLToPath } from "url"; | ||
| function defaultSourceDir() { | ||
| return dirname2(fileURLToPath(import.meta.url)); | ||
| } | ||
| function runtimeBinDir(home = homedir2()) { | ||
| return join5(home, ".llmnesia", "bin"); | ||
| } | ||
| async function installRuntime(options = {}) { | ||
| const home = options.home ?? homedir2(); | ||
| const sourceDir = options.sourceDir ?? defaultSourceDir(); | ||
| const nodePath = options.nodePath ?? process.execPath; | ||
| const binDir = runtimeBinDir(home); | ||
| const sourceEntry = join5(sourceDir, "cli.js"); | ||
| if (!existsSync(sourceEntry)) { | ||
| throw new Error( | ||
| `Cannot locate the built runtime: no cli.js in ${sourceDir}. Run this from an installed package (npx -y @llmnesia/mcp@latest install) or build first (npm run build).` | ||
| ); | ||
| } | ||
| const parentDir = dirname2(binDir); | ||
| const nonce = `${process.pid}-${Date.now()}`; | ||
| const stagingDir = join5(parentDir, `.bin-install-${nonce}`); | ||
| const backupDir = join5(parentDir, `.bin-backup-${nonce}`); | ||
| await mkdir(parentDir, { recursive: true }); | ||
| await rm(stagingDir, { recursive: true, force: true }); | ||
| await rm(backupDir, { recursive: true, force: true }); | ||
| await mkdir(stagingDir, { recursive: true }); | ||
| try { | ||
| for (const name of await readdir2(sourceDir)) { | ||
| await cp(join5(sourceDir, name), join5(stagingDir, name), { recursive: true }); | ||
| } | ||
| await writeFile( | ||
| join5(stagingDir, "package.json"), | ||
| `${JSON.stringify( | ||
| { | ||
| name: "llmnesia-mcp-runtime", | ||
| private: true, | ||
| type: "module", | ||
| // Stamped so the installed build is identifiable without spawning it. | ||
| // Its absence is how a machine sat on 0.1.3 while npm was ten | ||
| // releases ahead and nothing, doctor included, could say so. | ||
| version: options.version ?? VERSION, | ||
| comment: "Installed copy of @llmnesia/mcp's dist. Managed by the LLMnesia installer; edits are lost on reinstall." | ||
| }, | ||
| null, | ||
| 2 | ||
| )} | ||
| `, | ||
| "utf8" | ||
| ); | ||
| const hadPreviousRuntime = existsSync(binDir); | ||
| if (hadPreviousRuntime) { | ||
| await rename(binDir, backupDir); | ||
| } | ||
| try { | ||
| await rename(stagingDir, binDir); | ||
| } catch (err) { | ||
| if (hadPreviousRuntime && !existsSync(binDir) && existsSync(backupDir)) { | ||
| await rename(backupDir, binDir); | ||
| } | ||
| throw err; | ||
| } | ||
| await rm(backupDir, { recursive: true, force: true }); | ||
| } finally { | ||
| await rm(stagingDir, { recursive: true, force: true }); | ||
| } | ||
| const entry = join5(binDir, "cli.js"); | ||
| return { | ||
| binDir, | ||
| entry, | ||
| serve: { command: nodePath, args: [entry, "serve"] }, | ||
| sourceDir, | ||
| nodePath | ||
| }; | ||
| } | ||
| async function uninstallRuntime(home = homedir2()) { | ||
| const binDir = runtimeBinDir(home); | ||
| if (!existsSync(binDir)) { | ||
| return null; | ||
| } | ||
| await rm(binDir, { recursive: true, force: true }); | ||
| return binDir; | ||
| } | ||
| // src/selfUpdate.ts | ||
| import { spawn } from "child_process"; | ||
| import { mkdtemp, readFile, rm as rm2, writeFile as writeFile2 } from "fs/promises"; | ||
| import { dirname as dirname3, join as join6 } from "path"; | ||
| import { homedir as homedir3, tmpdir } from "os"; | ||
| import { fileURLToPath as fileURLToPath2 } from "url"; | ||
| var PACKAGE_NAME = "@llmnesia/mcp"; | ||
| var REGISTRY_LATEST_URL = "https://registry.npmjs.org/@llmnesia%2Fmcp/latest"; | ||
| var UPDATE_CHECK_INTERVAL_MS = 24 * 60 * 60 * 1e3; | ||
| var REGISTRY_TIMEOUT_MS = 5e3; | ||
| var INSTALL_TIMEOUT_MS = 12e4; | ||
| var UPDATE_CHECK_DELAY_MS = 3e4; | ||
| function stateFile(home) { | ||
| return join6(home, ".llmnesia", "update-state.json"); | ||
| } | ||
| async function readState(home) { | ||
| try { | ||
| const parsed = JSON.parse(await readFile(stateFile(home), "utf8")); | ||
| return typeof parsed === "object" && parsed !== null ? parsed : {}; | ||
| } catch { | ||
| return {}; | ||
| } | ||
| } | ||
| async function writeState(home, state) { | ||
| try { | ||
| await writeFile2(stateFile(home), `${JSON.stringify(state, null, 2)} | ||
| `, "utf8"); | ||
| } catch { | ||
| } | ||
| } | ||
| function compareVersions(a, b) { | ||
| const parse = (v) => { | ||
| const match = /^(\d+)\.(\d+)\.(\d+)/.exec(v.trim()); | ||
| return match ? [Number(match[1]), Number(match[2]), Number(match[3])] : []; | ||
| }; | ||
| const left = parse(a); | ||
| const right = parse(b); | ||
| if (left.length === 0 || right.length === 0) { | ||
| return 0; | ||
| } | ||
| for (let i = 0; i < 3; i += 1) { | ||
| if (left[i] !== right[i]) { | ||
| return left[i] - right[i]; | ||
| } | ||
| } | ||
| return 0; | ||
| } | ||
| function isStableRelease(version) { | ||
| return /^\d+\.\d+\.\d+$/.test(version.trim()); | ||
| } | ||
| function isSelfUpdateDisabled(env = process.env) { | ||
| const raw = env.LLMNESIA_NO_AUTO_UPDATE; | ||
| return raw === "1" || raw === "true"; | ||
| } | ||
| async function fetchLatestFromRegistry() { | ||
| const controller = new AbortController(); | ||
| const timer = setTimeout(() => controller.abort(), REGISTRY_TIMEOUT_MS); | ||
| try { | ||
| const response = await fetch(REGISTRY_LATEST_URL, { | ||
| signal: controller.signal, | ||
| headers: { accept: "application/json" } | ||
| }); | ||
| if (!response.ok) { | ||
| throw new Error(`registry responded ${response.status}`); | ||
| } | ||
| const body = await response.json(); | ||
| if (typeof body.version !== "string") { | ||
| throw new Error("registry response had no version"); | ||
| } | ||
| return body.version; | ||
| } finally { | ||
| clearTimeout(timer); | ||
| } | ||
| } | ||
| function npmCommand() { | ||
| const binDir = dirname3(process.execPath); | ||
| return process.platform === "win32" ? join6(binDir, "npm.cmd") : join6(binDir, "npm"); | ||
| } | ||
| async function npmDownload(version) { | ||
| const prefix = await mkdtemp(join6(tmpdir(), "llmnesia-update-")); | ||
| await new Promise((resolve2, reject) => { | ||
| const child = spawn( | ||
| npmCommand(), | ||
| [ | ||
| "install", | ||
| `${PACKAGE_NAME}@${version}`, | ||
| "--prefix", | ||
| prefix, | ||
| "--no-audit", | ||
| "--no-fund", | ||
| "--no-save", | ||
| "--loglevel", | ||
| "error" | ||
| ], | ||
| { stdio: ["ignore", "ignore", "pipe"], windowsHide: true } | ||
| ); | ||
| let stderr = ""; | ||
| child.stderr?.on("data", (chunk) => { | ||
| stderr += String(chunk); | ||
| }); | ||
| const timer = setTimeout(() => child.kill(), INSTALL_TIMEOUT_MS); | ||
| child.once("error", (err) => { | ||
| clearTimeout(timer); | ||
| reject(err); | ||
| }); | ||
| child.once("close", (code) => { | ||
| clearTimeout(timer); | ||
| if (code === 0) { | ||
| resolve2(); | ||
| } else { | ||
| reject(new Error(`npm install exited ${code}${stderr ? `: ${stderr.trim()}` : ""}`)); | ||
| } | ||
| }); | ||
| }); | ||
| return join6(prefix, "node_modules", PACKAGE_NAME, "dist"); | ||
| } | ||
| async function maybeSelfUpdate(options = {}) { | ||
| const env = options.env ?? process.env; | ||
| if (isSelfUpdateDisabled(env)) { | ||
| return { action: "disabled" }; | ||
| } | ||
| const home = options.home ?? homedir3(); | ||
| const currentVersion = options.currentVersion ?? VERSION; | ||
| const runningDir = options.runningDir ?? dirname3(fileURLToPath2(import.meta.url)); | ||
| if (runningDir !== runtimeBinDir(home)) { | ||
| return { action: "not-installed-runtime" }; | ||
| } | ||
| if (!isStableRelease(currentVersion)) { | ||
| return { action: "not-a-release-build", currentVersion }; | ||
| } | ||
| const now = options.now ?? Date.now; | ||
| const state = await readState(home); | ||
| if (state.lastCheckedAt !== void 0 && now() - state.lastCheckedAt < UPDATE_CHECK_INTERVAL_MS) { | ||
| return { action: "throttled" }; | ||
| } | ||
| let latest; | ||
| try { | ||
| latest = await (options.fetchLatestVersion ?? fetchLatestFromRegistry)(); | ||
| } catch (err) { | ||
| await writeState(home, { ...state, lastCheckedAt: now() }); | ||
| return { action: "check-failed", error: err instanceof Error ? err.message : String(err) }; | ||
| } | ||
| const checked = { ...state, lastCheckedAt: now(), lastSeenVersion: latest }; | ||
| if (!isStableRelease(latest) || compareVersions(latest, currentVersion) <= 0) { | ||
| await writeState(home, checked); | ||
| return { action: "up-to-date", latest }; | ||
| } | ||
| if (state.lastFailedVersion === latest) { | ||
| await writeState(home, checked); | ||
| return { action: "skipped-failed-before", latest }; | ||
| } | ||
| let sourceDir; | ||
| try { | ||
| sourceDir = await (options.downloadPackage ?? npmDownload)(latest); | ||
| await installRuntime({ home, sourceDir, version: latest }); | ||
| await writeState(home, { ...checked, lastUpdatedTo: latest, lastFailedVersion: void 0 }); | ||
| return { action: "updated", from: currentVersion, to: latest }; | ||
| } catch (err) { | ||
| await writeState(home, { ...checked, lastFailedVersion: latest }); | ||
| return { action: "update-failed", latest, error: err instanceof Error ? err.message : String(err) }; | ||
| } finally { | ||
| if (sourceDir) { | ||
| await rm2(join6(sourceDir, "..", "..", ".."), { recursive: true, force: true }).catch(() => { | ||
| }); | ||
| } | ||
| } | ||
| } | ||
| function scheduleSelfUpdate(options = {}) { | ||
| const { delayMs = UPDATE_CHECK_DELAY_MS, onOutcome, ...rest } = options; | ||
| const timer = setTimeout(() => { | ||
| void maybeSelfUpdate(rest).then( | ||
| (outcome) => onOutcome?.(outcome), | ||
| () => { | ||
| } | ||
| ); | ||
| }, delayMs); | ||
| timer.unref?.(); | ||
| } | ||
| function describeOutcome(outcome) { | ||
| switch (outcome.action) { | ||
| case "updated": | ||
| return `updated ${outcome.from} to ${outcome.to}; it takes effect next launch`; | ||
| case "update-failed": | ||
| return `could not install ${outcome.latest}: ${outcome.error}`; | ||
| case "check-failed": | ||
| return `could not reach the npm registry: ${outcome.error}`; | ||
| case "up-to-date": | ||
| return `already current (latest published is ${outcome.latest})`; | ||
| case "skipped-failed-before": | ||
| return `${outcome.latest} failed to install previously; not retrying automatically`; | ||
| case "throttled": | ||
| return "checked recently"; | ||
| case "disabled": | ||
| return "automatic updates are turned off (LLMNESIA_NO_AUTO_UPDATE)"; | ||
| case "not-installed-runtime": | ||
| return "not running from ~/.llmnesia/bin, so nothing to update"; | ||
| case "not-a-release-build": | ||
| return `running an unreleased build (${outcome.currentVersion}); leaving the installed runtime alone`; | ||
| } | ||
| } | ||
| export { | ||
| resolveCorpusDir, | ||
| staleIndexFiles, | ||
| Corpus, | ||
| ingestFile, | ||
| pendingInboxFiles, | ||
| drainInbox, | ||
| SearchIndex, | ||
| openSearchIndex, | ||
| deleteIndexDb, | ||
| runtimeBinDir, | ||
| installRuntime, | ||
| uninstallRuntime, | ||
| isSelfUpdateDisabled, | ||
| scheduleSelfUpdate, | ||
| describeOutcome | ||
| }; |
Sorry, the diff of this file is too big to display
+2
-2
@@ -6,3 +6,3 @@ #!/usr/bin/env node | ||
| VERSION | ||
| } from "./chunk-6H2TI5IW.js"; | ||
| } from "./chunk-JSNSTAFD.js"; | ||
@@ -145,3 +145,3 @@ // src/nodeDownload.ts | ||
| } | ||
| const { main } = await import("./cli-77JECNWA.js"); | ||
| const { main } = await import("./cli-ZEA7WTVW.js"); | ||
| return main(); | ||
@@ -148,0 +148,0 @@ } |
+1
-1
| { | ||
| "name": "@llmnesia/mcp", | ||
| "version": "0.2.3", | ||
| "version": "0.2.4", | ||
| "private": false, | ||
@@ -5,0 +5,0 @@ "type": "module", |
+26
-0
@@ -143,2 +143,26 @@ # @llmnesia/mcp | ||
| ### Local storage and compaction | ||
| Conversation files are stored as losslessly compressed JSON. Search uses a | ||
| contentless FTS5 token index: the canonical text is not duplicated inside the | ||
| database, and result snippets are loaded from the canonical conversation only | ||
| after SQLite has ranked the matches. No conversation or searchable token is | ||
| dropped to save space. | ||
| Current builds discard a clean browser-to-corpus transport payload after both | ||
| the canonical conversation and active index have accepted it. Malformed input | ||
| is kept under `inbox/quarantine/` for inspection. Older releases archived every | ||
| successful payload and may also have left versioned indexes behind. Preview the | ||
| exact recoverable space without writing anything: | ||
| ```bash | ||
| npx -y @llmnesia/mcp@latest compact --dry-run | ||
| ``` | ||
| To apply it, fully quit connected desktop AI/MCP clients, then run `compact` | ||
| without `--dry-run`. It asks for confirmation, verifies every compressed | ||
| conversation by reading it back before removing the legacy JSON, builds the | ||
| current index before deleting older derived versions, and never removes | ||
| quarantined input. | ||
| ## Tools | ||
@@ -221,2 +245,4 @@ | ||
| llmnesia-mcp stats # corpus summary | ||
| llmnesia-mcp compact --dry-run # exact read-only storage/savings report | ||
| llmnesia-mcp compact # confirmed lossless migration + redundant cleanup | ||
| llmnesia-mcp reindex # drop + rebuild the search index from the corpus | ||
@@ -223,0 +249,0 @@ llmnesia-mcp search <q> # one-shot search from the terminal |
+2
-2
@@ -6,3 +6,3 @@ { | ||
| "description": "Search your AI chat history from Claude, Cursor, Codex, and other MCP clients.", | ||
| "version": "0.2.3", | ||
| "version": "0.2.4", | ||
| "packages": [ | ||
@@ -13,3 +13,3 @@ { | ||
| "identifier": "@llmnesia/mcp", | ||
| "version": "0.2.3", | ||
| "version": "0.2.4", | ||
| "transport": { | ||
@@ -16,0 +16,0 @@ "type": "stdio" |
| #!/usr/bin/env node | ||
| import { createRequire as __llmnesiaCreateRequire } from "node:module"; | ||
| const require = __llmnesiaCreateRequire(import.meta.url); | ||
| var __create = Object.create; | ||
| var __defProp = Object.defineProperty; | ||
| var __getOwnPropDesc = Object.getOwnPropertyDescriptor; | ||
| var __getOwnPropNames = Object.getOwnPropertyNames; | ||
| var __getProtoOf = Object.getPrototypeOf; | ||
| var __hasOwnProp = Object.prototype.hasOwnProperty; | ||
| var __require = /* @__PURE__ */ ((x) => typeof require !== "undefined" ? require : typeof Proxy !== "undefined" ? new Proxy(x, { | ||
| get: (a, b) => (typeof require !== "undefined" ? require : a)[b] | ||
| }) : x)(function(x) { | ||
| if (typeof require !== "undefined") return require.apply(this, arguments); | ||
| throw Error('Dynamic require of "' + x + '" is not supported'); | ||
| }); | ||
| var __commonJS = (cb, mod) => function __require2() { | ||
| return mod || (0, cb[__getOwnPropNames(cb)[0]])((mod = { exports: {} }).exports, mod), mod.exports; | ||
| }; | ||
| var __export = (target, all) => { | ||
| for (var name in all) | ||
| __defProp(target, name, { get: all[name], enumerable: true }); | ||
| }; | ||
| var __copyProps = (to, from, except, desc) => { | ||
| if (from && typeof from === "object" || typeof from === "function") { | ||
| for (let key of __getOwnPropNames(from)) | ||
| if (!__hasOwnProp.call(to, key) && key !== except) | ||
| __defProp(to, key, { get: () => from[key], enumerable: !(desc = __getOwnPropDesc(from, key)) || desc.enumerable }); | ||
| } | ||
| return to; | ||
| }; | ||
| var __toESM = (mod, isNodeMode, target) => (target = mod != null ? __create(__getProtoOf(mod)) : {}, __copyProps( | ||
| // If the importer is in node compatibility mode or this is not an ESM | ||
| // file that has been converted to a CommonJS file using a Babel- | ||
| // compatible transform (i.e. "__esModule" has not been set), then set | ||
| // "default" to the CommonJS "module.exports" for node compatibility. | ||
| isNodeMode || !mod || !mod.__esModule ? __defProp(target, "default", { value: mod, enumerable: true }) : target, | ||
| mod | ||
| )); | ||
| // src/version.ts | ||
| var VERSION = true ? "0.2.3" : "0.0.0-dev"; | ||
| export { | ||
| __require, | ||
| __commonJS, | ||
| __export, | ||
| __toESM, | ||
| VERSION | ||
| }; |
| #!/usr/bin/env node | ||
| import { createRequire as __llmnesiaCreateRequire } from "node:module"; | ||
| const require = __llmnesiaCreateRequire(import.meta.url); | ||
| import { | ||
| VERSION | ||
| } from "./chunk-6H2TI5IW.js"; | ||
| // src/config.ts | ||
| import { readFileSync } from "fs"; | ||
| import { homedir } from "os"; | ||
| import { join, resolve } from "path"; | ||
| var CORPUS_FORMAT_VERSION = 1; | ||
| function nativeHostConfigPath() { | ||
| return join(homedir(), ".llmnesia", "config.json"); | ||
| } | ||
| function readNativeHostCorpusRoot() { | ||
| try { | ||
| const parsed = JSON.parse(readFileSync(nativeHostConfigPath(), "utf8")); | ||
| return typeof parsed.corpusRoot === "string" ? parsed.corpusRoot : void 0; | ||
| } catch { | ||
| return void 0; | ||
| } | ||
| } | ||
| function resolveCorpusDir(explicit) { | ||
| const raw = explicit ?? process.env.LLMNESIA_CORPUS_DIR ?? readNativeHostCorpusRoot() ?? "~/.llmnesia/corpus"; | ||
| return expandHome(raw); | ||
| } | ||
| function expandHome(p) { | ||
| if (p === "~") { | ||
| return homedir(); | ||
| } | ||
| if (p.startsWith("~/") || p.startsWith("~\\")) { | ||
| return resolve(homedir(), p.slice(2)); | ||
| } | ||
| return resolve(p); | ||
| } | ||
| // src/paths.ts | ||
| import { join as join2 } from "path"; | ||
| import { readdir } from "fs/promises"; | ||
| import { createHash } from "crypto"; | ||
| var INDEX_FORMAT_VERSION = 2; | ||
| function indexDbFilename(version) { | ||
| return version === 1 ? "llmnesia.db" : `llmnesia-v${version}.db`; | ||
| } | ||
| async function staleIndexFiles(dir) { | ||
| const keep = indexDbFilename(INDEX_FORMAT_VERSION); | ||
| let entries; | ||
| try { | ||
| entries = await readdir(dir); | ||
| } catch { | ||
| return []; | ||
| } | ||
| return entries.filter( | ||
| (name) => /^llmnesia(-v\d+)?\.db$/.test(name) && name !== keep | ||
| ); | ||
| } | ||
| var CorpusPaths = class { | ||
| constructor(root) { | ||
| this.root = root; | ||
| } | ||
| root; | ||
| get metaFile() { | ||
| return join2(this.root, "meta.json"); | ||
| } | ||
| get inboxDir() { | ||
| return join2(this.root, "inbox"); | ||
| } | ||
| get processedDir() { | ||
| return join2(this.inboxDir, "processed"); | ||
| } | ||
| get conversationsDir() { | ||
| return join2(this.root, "conversations"); | ||
| } | ||
| get memoriesDir() { | ||
| return join2(this.root, "memories"); | ||
| } | ||
| get indexDir() { | ||
| return join2(this.root, "index"); | ||
| } | ||
| get indexDbFile() { | ||
| return join2(this.indexDir, indexDbFilename(INDEX_FORMAT_VERSION)); | ||
| } | ||
| /** Directory holding a platform's conversation JSON files. */ | ||
| platformDir(platform) { | ||
| return join2(this.conversationsDir, safeSegment(platform)); | ||
| } | ||
| /** Absolute path to the JSON file backing a given docId. */ | ||
| conversationFile(platform, docId) { | ||
| return join2(this.platformDir(platform), `${slugifyDocId(docId)}.json`); | ||
| } | ||
| /** | ||
| * Absolute path to the JSON file backing a given memoryId. Memories are flat | ||
| * rather than partitioned by platform: they are user-curated and few, and a | ||
| * memory is not owned by any one platform the way a conversation is. | ||
| */ | ||
| memoryFile(memoryId) { | ||
| return join2(this.memoriesDir, `${slugifyDocId(memoryId)}.json`); | ||
| } | ||
| }; | ||
| var SLUG_MAX = 120; | ||
| function slugifyDocId(docId) { | ||
| const lossless = docId.length > 0 && docId.length <= SLUG_MAX && /^[a-zA-Z0-9._-]+$/.test(docId); | ||
| if (lossless) { | ||
| return docId; | ||
| } | ||
| const hash = createHash("sha1").update(docId).digest("hex").slice(0, 12); | ||
| const readable = docId.replace(/[^a-zA-Z0-9._-]+/g, "-").replace(/^-+|-+$/g, "").slice(0, SLUG_MAX - 13); | ||
| return readable.length > 0 ? `${readable}-${hash}` : hash; | ||
| } | ||
| function safeSegment(value) { | ||
| const cleaned = value.replace(/[^a-zA-Z0-9._-]+/g, "-").replace(/^[.-]+|[.-]+$/g, ""); | ||
| return cleaned.length > 0 ? cleaned : "unknown"; | ||
| } | ||
| // src/corpus.ts | ||
| import { promises as fs } from "fs"; | ||
| import { dirname, join as join3 } from "path"; | ||
| // src/serialize.ts | ||
| function canonicalizeConversation(record) { | ||
| return { | ||
| conversation: canonicalConversationDoc(record.conversation), | ||
| messages: [...record.messages].sort((a, b) => a.msgIndex - b.msgIndex).map(canonicalMessageDoc) | ||
| }; | ||
| } | ||
| function canonicalizeMemory(memory) { | ||
| const out = { | ||
| memoryId: memory.memoryId, | ||
| kind: memory.kind, | ||
| label: memory.label, | ||
| text: memory.text, | ||
| source: memory.source | ||
| }; | ||
| assignIfDefined(out, "sourceDocId", memory.sourceDocId); | ||
| assignIfDefined(out, "sourcePlatform", memory.sourcePlatform); | ||
| assignIfDefined(out, "sourceUrl", memory.sourceUrl); | ||
| assignIfDefined(out, "sourceTitle", memory.sourceTitle); | ||
| assignIfDefined(out, "sourceMsgStart", memory.sourceMsgStart); | ||
| assignIfDefined(out, "sourceMsgEnd", memory.sourceMsgEnd); | ||
| out.pinned = memory.pinned; | ||
| out.enabled = memory.enabled; | ||
| if (memory.tags && memory.tags.length > 0) { | ||
| out.tags = [...memory.tags].sort(); | ||
| } | ||
| out.createdAt = memory.createdAt; | ||
| out.updatedAt = memory.updatedAt; | ||
| assignIfDefined(out, "lastUsedAt", memory.lastUsedAt); | ||
| assignIfDefined(out, "useCount", memory.useCount); | ||
| return out; | ||
| } | ||
| function canonicalConversationDoc(c) { | ||
| const out = { | ||
| docId: c.docId, | ||
| platform: c.platform, | ||
| title: c.title, | ||
| url: c.url, | ||
| createdAt: c.createdAt, | ||
| contentUpdatedAt: c.contentUpdatedAt, | ||
| indexedAt: c.indexedAt, | ||
| corpusChangedAt: c.corpusChangedAt, | ||
| updatedAt: c.updatedAt, | ||
| lastSeenAt: c.lastSeenAt, | ||
| messageCount: c.messageCount, | ||
| preview: c.preview, | ||
| hasImages: c.hasImages, | ||
| pinned: c.pinned, | ||
| indexLevel: c.indexLevel, | ||
| estimatedBytes: c.estimatedBytes | ||
| }; | ||
| assignIfDefined(out, "searchDocLength", c.searchDocLength); | ||
| assignIfDefined(out, "contentHash", c.contentHash); | ||
| assignIfDefined(out, "truncated", c.truncated); | ||
| assignIfDefined(out, "kind", c.kind); | ||
| assignIfDefined(out, "sourceId", c.sourceId); | ||
| assignIfDefined(out, "baseUrl", c.baseUrl); | ||
| assignIfDefined(out, "isDeeplinkable", c.isDeeplinkable); | ||
| assignIfDefined(out, "origin", c.origin); | ||
| assignIfDefined(out, "accountId", c.accountId); | ||
| assignIfDefined(out, "accountLabel", c.accountLabel); | ||
| if (c.summary) { | ||
| out.summary = { | ||
| text: c.summary.text, | ||
| recipe: c.summary.recipe, | ||
| method: c.summary.method, | ||
| createdAt: c.summary.createdAt | ||
| }; | ||
| } | ||
| return out; | ||
| } | ||
| function canonicalMessageDoc(m) { | ||
| const out = { | ||
| msgKey: m.msgKey, | ||
| docId: m.docId, | ||
| msgIndex: m.msgIndex, | ||
| role: m.role, | ||
| text: m.text | ||
| }; | ||
| assignIfDefined(out, "formattedText", m.formattedText); | ||
| out.createdAt = m.createdAt; | ||
| out.updatedAt = m.updatedAt; | ||
| out.contentHash = m.contentHash; | ||
| assignIfDefined(out, "platformMessageId", m.platformMessageId); | ||
| return out; | ||
| } | ||
| function assignIfDefined(target, key, value) { | ||
| if (value !== void 0) { | ||
| target[key] = value; | ||
| } | ||
| } | ||
| // src/corpus.ts | ||
| var DIR_MODE = 448; | ||
| var FILE_MODE = 384; | ||
| var Corpus = class { | ||
| paths; | ||
| // Serialize writes within this process so concurrent upserts (inbox drain + | ||
| // MCP save_conversation) don't interleave. Cross-process safety relies on the | ||
| // attended, single-companion v1 model (see PLAN-MCP.md open question 3). | ||
| writeChain = Promise.resolve(); | ||
| constructor(root) { | ||
| this.paths = new CorpusPaths(root); | ||
| } | ||
| /** Create the directory skeleton and meta.json if absent. Idempotent. */ | ||
| async init() { | ||
| await fs.mkdir(this.paths.conversationsDir, { recursive: true, mode: DIR_MODE }); | ||
| await fs.mkdir(this.paths.memoriesDir, { recursive: true, mode: DIR_MODE }); | ||
| await fs.mkdir(this.paths.processedDir, { recursive: true, mode: DIR_MODE }); | ||
| await fs.mkdir(this.paths.indexDir, { recursive: true, mode: DIR_MODE }); | ||
| await hardenDir(this.paths.root); | ||
| await hardenDir(this.paths.conversationsDir); | ||
| await hardenDir(this.paths.memoriesDir); | ||
| await hardenDir(this.paths.inboxDir); | ||
| await hardenDir(this.paths.processedDir); | ||
| await hardenDir(this.paths.indexDir); | ||
| try { | ||
| await fs.access(this.paths.metaFile); | ||
| } catch { | ||
| const meta = { | ||
| formatVersion: CORPUS_FORMAT_VERSION, | ||
| createdAt: Date.now(), | ||
| lastIngestAt: null, | ||
| conversationCount: 0 | ||
| }; | ||
| await writeFileAtomic(this.paths.metaFile, `${JSON.stringify(meta, null, 2)} | ||
| `); | ||
| } | ||
| } | ||
| async readMeta() { | ||
| try { | ||
| const raw = await fs.readFile(this.paths.metaFile, "utf8"); | ||
| return JSON.parse(raw); | ||
| } catch { | ||
| return null; | ||
| } | ||
| } | ||
| async writeMeta(patch) { | ||
| const current = await this.readMeta() ?? { | ||
| formatVersion: CORPUS_FORMAT_VERSION, | ||
| createdAt: Date.now(), | ||
| lastIngestAt: null, | ||
| conversationCount: 0 | ||
| }; | ||
| const next = { ...current, ...patch }; | ||
| await writeFileAtomic(this.paths.metaFile, `${JSON.stringify(next, null, 2)} | ||
| `); | ||
| } | ||
| /** | ||
| * Insert or replace a conversation, keyed on docId. New bridge records use | ||
| * the local `corpusChangedAt` mutation clock; older records fall back to the | ||
| * source `updatedAt`. If mutation clocks tie, source freshness breaks the | ||
| * tie. Re-writing identical content remains a byte-identical no-op. | ||
| */ | ||
| async upsert(record) { | ||
| return this.enqueueWrite(async () => { | ||
| const canonical = canonicalizeConversation(record); | ||
| const { conversation } = canonical; | ||
| const file = this.paths.conversationFile(conversation.platform, conversation.docId); | ||
| const existing = await readJsonIfExists(file); | ||
| if (existing) { | ||
| const existingChangedAt = existing.conversation.corpusChangedAt ?? existing.conversation.updatedAt ?? 0; | ||
| const incomingChangedAt = conversation.corpusChangedAt ?? conversation.updatedAt ?? 0; | ||
| const existingUpdatedAt = existing.conversation.updatedAt ?? 0; | ||
| const incomingUpdatedAt = conversation.updatedAt ?? 0; | ||
| if (incomingChangedAt < existingChangedAt || incomingChangedAt === existingChangedAt && incomingUpdatedAt < existingUpdatedAt) { | ||
| return { docId: conversation.docId, outcome: "skipped" }; | ||
| } | ||
| } | ||
| const serialized = `${JSON.stringify(canonical, null, 2)} | ||
| `; | ||
| if (existing) { | ||
| const existingSerialized = `${JSON.stringify(canonicalizeConversation(existing), null, 2)} | ||
| `; | ||
| if (existingSerialized === serialized) { | ||
| return { docId: conversation.docId, outcome: "skipped" }; | ||
| } | ||
| } | ||
| await fs.mkdir(dirname(file), { recursive: true, mode: DIR_MODE }); | ||
| await writeFileAtomic(file, serialized); | ||
| return { docId: conversation.docId, outcome: existing ? "updated" : "created" }; | ||
| }); | ||
| } | ||
| /** Read a conversation by docId, scanning platform dirs for the slug file. */ | ||
| async get(docId) { | ||
| const filename = `${slugifyDocId(docId)}.json`; | ||
| for (const platform of await this.listPlatforms()) { | ||
| const candidate = join3(this.paths.platformDir(platform), filename); | ||
| const record = await readJsonIfExists(candidate); | ||
| if (record && record.conversation.docId === docId) { | ||
| return record; | ||
| } | ||
| } | ||
| return null; | ||
| } | ||
| /** | ||
| * Write an agent-authored memory, keyed on memoryId. | ||
| * | ||
| * The semantics differ from {@link upsert} on purpose. A memoryId is a | ||
| * content hash, so an existing file at that id already holds this exact | ||
| * memory: re-writing it would only churn the timestamps, and it would stamp | ||
| * `enabled: false` back over a memory the user had reviewed and turned on. | ||
| * So an existing record wins and the write is skipped. | ||
| * | ||
| * A record whose `source` is not "agent" is refused outright rather than | ||
| * skipped: a captured conversation or a memory the user wrote by hand must be | ||
| * untouchable from the agent write path, and the caller should hear about it | ||
| * rather than believe its write landed. | ||
| */ | ||
| async upsertMemory(memory) { | ||
| return this.enqueueWrite(async () => { | ||
| const canonical = canonicalizeMemory(memory); | ||
| const file = this.paths.memoryFile(canonical.memoryId); | ||
| const existing = await readJsonIfExists(file); | ||
| if (existing) { | ||
| if (existing.source !== "agent") { | ||
| return { memoryId: canonical.memoryId, outcome: "refused" }; | ||
| } | ||
| return { memoryId: canonical.memoryId, outcome: "skipped" }; | ||
| } | ||
| await fs.mkdir(dirname(file), { recursive: true, mode: DIR_MODE }); | ||
| await writeFileAtomic(file, `${JSON.stringify(canonical, null, 2)} | ||
| `); | ||
| return { memoryId: canonical.memoryId, outcome: "created" }; | ||
| }); | ||
| } | ||
| /** Read a memory by memoryId, or null when it has never been written. */ | ||
| async getMemory(memoryId) { | ||
| const record = await readJsonIfExists(this.paths.memoryFile(memoryId)); | ||
| return record && record.memoryId === memoryId ? record : null; | ||
| } | ||
| /** Async-iterate every stored memory (used by the index rebuild). */ | ||
| async *iterateMemories() { | ||
| let files; | ||
| try { | ||
| files = await fs.readdir(this.paths.memoriesDir); | ||
| } catch { | ||
| return; | ||
| } | ||
| for (const name of files.sort()) { | ||
| if (!name.endsWith(".json")) { | ||
| continue; | ||
| } | ||
| const record = await readJsonIfExists(join3(this.paths.memoriesDir, name)); | ||
| if (record?.memoryId) { | ||
| yield record; | ||
| } | ||
| } | ||
| } | ||
| /** Platform subdirectory names under conversations/. */ | ||
| async listPlatforms() { | ||
| try { | ||
| const entries = await fs.readdir(this.paths.conversationsDir, { withFileTypes: true }); | ||
| return entries.filter((e) => e.isDirectory()).map((e) => e.name); | ||
| } catch { | ||
| return []; | ||
| } | ||
| } | ||
| /** Async-iterate every stored conversation (used by reindex + stats). */ | ||
| async *iterate() { | ||
| for (const platform of await this.listPlatforms()) { | ||
| const dir = this.paths.platformDir(platform); | ||
| let files; | ||
| try { | ||
| files = await fs.readdir(dir); | ||
| } catch { | ||
| continue; | ||
| } | ||
| for (const name of files) { | ||
| if (!name.endsWith(".json")) { | ||
| continue; | ||
| } | ||
| const record = await readJsonIfExists(join3(dir, name)); | ||
| if (record) { | ||
| yield record; | ||
| } | ||
| } | ||
| } | ||
| } | ||
| async stats() { | ||
| const byPlatform = {}; | ||
| let conversationCount = 0; | ||
| let messageCount = 0; | ||
| for await (const record of this.iterate()) { | ||
| conversationCount += 1; | ||
| messageCount += record.messages.length; | ||
| const platform = safeSegment(record.conversation.platform); | ||
| byPlatform[platform] = (byPlatform[platform] ?? 0) + 1; | ||
| } | ||
| return { conversationCount, messageCount, byPlatform }; | ||
| } | ||
| /** Refresh meta.json's conversationCount / lastIngestAt after an ingest run. */ | ||
| async touchIngestMeta() { | ||
| const { conversationCount } = await this.stats(); | ||
| await this.writeMeta({ lastIngestAt: Date.now(), conversationCount }); | ||
| } | ||
| enqueueWrite(fn) { | ||
| const run = this.writeChain.then(fn, fn); | ||
| this.writeChain = run.then( | ||
| () => void 0, | ||
| () => void 0 | ||
| ); | ||
| return run; | ||
| } | ||
| }; | ||
| async function readJsonIfExists(file) { | ||
| try { | ||
| const raw = await fs.readFile(file, "utf8"); | ||
| return JSON.parse(raw); | ||
| } catch { | ||
| return null; | ||
| } | ||
| } | ||
| async function writeFileAtomic(file, contents) { | ||
| const tmp = `${file}.tmp-${process.pid}-${Date.now()}`; | ||
| await fs.writeFile(tmp, contents, { encoding: "utf8", mode: FILE_MODE }); | ||
| await fs.rename(tmp, file); | ||
| } | ||
| async function hardenDir(dir) { | ||
| try { | ||
| await fs.chmod(dir, DIR_MODE); | ||
| } catch { | ||
| } | ||
| } | ||
| // src/ingest.ts | ||
| import { promises as fs2, createReadStream } from "fs"; | ||
| import { join as join4 } from "path"; | ||
| import { createInterface } from "readline"; | ||
| import { createGunzip } from "zlib"; | ||
| var EMPTY_REPORT = { | ||
| conversations: 0, | ||
| created: 0, | ||
| updated: 0, | ||
| skipped: 0, | ||
| messages: 0, | ||
| malformedLines: 0 | ||
| }; | ||
| async function ingestFile(corpus, file, index) { | ||
| const report = { ...EMPTY_REPORT }; | ||
| let current = null; | ||
| let buffered = []; | ||
| const flush = async () => { | ||
| if (!current) { | ||
| return; | ||
| } | ||
| const record = { conversation: current, messages: buffered }; | ||
| const result = await corpus.upsert(record); | ||
| tally(report, result); | ||
| if (index && result.outcome !== "skipped") { | ||
| index.upsertConversation(record); | ||
| } | ||
| report.conversations += 1; | ||
| report.messages += buffered.length; | ||
| current = null; | ||
| buffered = []; | ||
| }; | ||
| const rl = createInterface({ input: openLineSource(file), crlfDelay: Infinity }); | ||
| for await (const raw of rl) { | ||
| const line = raw.trim(); | ||
| if (!line) { | ||
| continue; | ||
| } | ||
| let parsed; | ||
| try { | ||
| parsed = JSON.parse(line); | ||
| } catch { | ||
| report.malformedLines += 1; | ||
| continue; | ||
| } | ||
| if (typeof parsed !== "object" || parsed === null) { | ||
| report.malformedLines += 1; | ||
| continue; | ||
| } | ||
| const record = parsed; | ||
| const type = record.type; | ||
| if (type === "meta") { | ||
| continue; | ||
| } | ||
| if (type === "conversation") { | ||
| await flush(); | ||
| current = coerceConversation(record.conversation); | ||
| buffered = []; | ||
| if (!current) { | ||
| report.malformedLines += 1; | ||
| } | ||
| continue; | ||
| } | ||
| if (type === "message" && current) { | ||
| const docId = typeof record.docId === "string" ? record.docId : ""; | ||
| if (docId !== current.docId) { | ||
| report.malformedLines += 1; | ||
| continue; | ||
| } | ||
| const message = coerceMessage(current.docId, record.message); | ||
| if (message) { | ||
| buffered.push(message); | ||
| } else { | ||
| report.malformedLines += 1; | ||
| } | ||
| } | ||
| } | ||
| await flush(); | ||
| return report; | ||
| } | ||
| async function pendingInboxFiles(corpus) { | ||
| try { | ||
| return (await fs2.readdir(corpus.paths.inboxDir, { withFileTypes: true })).filter((e) => e.isFile() && isNdjson(e.name)).map((e) => e.name).sort(); | ||
| } catch { | ||
| return []; | ||
| } | ||
| } | ||
| async function drainInbox(corpus, index, options = {}) { | ||
| await corpus.init(); | ||
| const inbox = corpus.paths.inboxDir; | ||
| const entries = await pendingInboxFiles(corpus); | ||
| const total = { ...EMPTY_REPORT, files: 0, failed: 0 }; | ||
| for (const name of entries) { | ||
| const src = join4(inbox, name); | ||
| try { | ||
| const report = await ingestFile(corpus, src, index); | ||
| accumulate(total, report); | ||
| total.files += 1; | ||
| await moveToProcessed(corpus, src, name); | ||
| } catch (error) { | ||
| if (await stillPending(src)) { | ||
| total.failed += 1; | ||
| options.onFileError?.(name, error); | ||
| } | ||
| } | ||
| } | ||
| if (total.files > 0) { | ||
| await corpus.touchIngestMeta(); | ||
| } | ||
| return total; | ||
| } | ||
| async function stillPending(src) { | ||
| try { | ||
| await fs2.access(src); | ||
| return true; | ||
| } catch { | ||
| return false; | ||
| } | ||
| } | ||
| async function moveToProcessed(corpus, src, name) { | ||
| const dest = join4(corpus.paths.processedDir, name); | ||
| try { | ||
| await fs2.rename(src, dest); | ||
| } catch { | ||
| const alt = join4(corpus.paths.processedDir, `${Date.now()}-${name}`); | ||
| await fs2.copyFile(src, alt); | ||
| await fs2.unlink(src); | ||
| } | ||
| } | ||
| function isNdjson(name) { | ||
| const lower = name.toLowerCase(); | ||
| return lower.endsWith(".jsonl") || lower.endsWith(".ndjson") || lower.endsWith(".jsonl.gz") || lower.endsWith(".ndjson.gz"); | ||
| } | ||
| function openLineSource(file) { | ||
| const stream = createReadStream(file); | ||
| if (file.toLowerCase().endsWith(".gz")) { | ||
| return stream.pipe(createGunzip()); | ||
| } | ||
| return stream; | ||
| } | ||
| function coerceConversation(value) { | ||
| if (!value || typeof value !== "object") { | ||
| return null; | ||
| } | ||
| const c = value; | ||
| if (typeof c.docId !== "string" || !c.docId) { | ||
| return null; | ||
| } | ||
| const conversation = { ...c }; | ||
| if (typeof conversation.updatedAt !== "number") { | ||
| conversation.updatedAt = typeof conversation.contentUpdatedAt === "number" ? conversation.contentUpdatedAt : 0; | ||
| } | ||
| return conversation; | ||
| } | ||
| function coerceMessage(docId, value) { | ||
| if (!value || typeof value !== "object") { | ||
| return null; | ||
| } | ||
| const m = value; | ||
| const text = typeof m.text === "string" ? m.text : ""; | ||
| if (!text.trim()) { | ||
| return null; | ||
| } | ||
| const msgIndex = typeof m.msgIndex === "number" ? m.msgIndex : 0; | ||
| const message = { | ||
| // The backup line omits msgKey; reconstruct the extension's convention. | ||
| msgKey: `${docId}:${msgIndex}`, | ||
| docId, | ||
| msgIndex, | ||
| role: coerceRole(m.role), | ||
| text, | ||
| createdAt: typeof m.createdAt === "number" ? m.createdAt : 0, | ||
| updatedAt: typeof m.updatedAt === "number" ? m.updatedAt : 0, | ||
| contentHash: typeof m.contentHash === "string" ? m.contentHash : "" | ||
| }; | ||
| if (typeof m.formattedText === "string") { | ||
| message.formattedText = m.formattedText; | ||
| } | ||
| if (typeof m.platformMessageId === "string") { | ||
| message.platformMessageId = m.platformMessageId; | ||
| } else if (m.platformMessageId === null) { | ||
| message.platformMessageId = null; | ||
| } | ||
| return message; | ||
| } | ||
| function coerceRole(value) { | ||
| return value === "assistant" || value === "user" || value === "system" || value === "unknown" ? value : "unknown"; | ||
| } | ||
| function tally(report, result) { | ||
| if (result.outcome === "created") { | ||
| report.created += 1; | ||
| } else if (result.outcome === "updated") { | ||
| report.updated += 1; | ||
| } else { | ||
| report.skipped += 1; | ||
| } | ||
| } | ||
| function accumulate(total, report) { | ||
| total.conversations += report.conversations; | ||
| total.created += report.created; | ||
| total.updated += report.updated; | ||
| total.skipped += report.skipped; | ||
| total.messages += report.messages; | ||
| total.malformedLines += report.malformedLines; | ||
| } | ||
| // src/searchIndex.ts | ||
| import { promises as fs3, chmodSync, closeSync, openSync } from "fs"; | ||
| import { createRequire } from "module"; | ||
| // src/sqliteWarning.ts | ||
| var INSTALLED = /* @__PURE__ */ Symbol.for("llmnesia.sqliteWarningFilterInstalled"); | ||
| function isSqliteExperimentalWarning(warning) { | ||
| return warning.name === "ExperimentalWarning" && /\bSQLite\b/i.test(warning.message); | ||
| } | ||
| function install() { | ||
| const flagged = globalThis; | ||
| if (flagged[INSTALLED]) { | ||
| return; | ||
| } | ||
| flagged[INSTALLED] = true; | ||
| const previous = process.listeners("warning"); | ||
| process.removeAllListeners("warning"); | ||
| process.on("warning", (warning) => { | ||
| if (isSqliteExperimentalWarning(warning)) { | ||
| return; | ||
| } | ||
| for (const listener of previous) { | ||
| listener.call(process, warning); | ||
| } | ||
| }); | ||
| } | ||
| install(); | ||
| // src/searchIndex.ts | ||
| var nodeRequire = createRequire(import.meta.url); | ||
| var sqlite; | ||
| function loadSqlite() { | ||
| return sqlite ??= nodeRequire("node:sqlite"); | ||
| } | ||
| var DB_FILE_MODE = 384; | ||
| function hardenDbFile(dbPath) { | ||
| try { | ||
| closeSync(openSync(dbPath, "a", DB_FILE_MODE)); | ||
| chmodSync(dbPath, DB_FILE_MODE); | ||
| } catch { | ||
| } | ||
| } | ||
| var TITLE_ROW_INDEX = -1; | ||
| var TITLE_ROLE = "title"; | ||
| var DEFAULT_LIMIT = 10; | ||
| var SNIPPET_TOKENS = 12; | ||
| var META_SQL = ` | ||
| CREATE TABLE IF NOT EXISTS index_meta ( | ||
| key TEXT PRIMARY KEY, | ||
| value TEXT NOT NULL | ||
| ); | ||
| `; | ||
| var DROP_DERIVED_SQL = ` | ||
| DROP TABLE IF EXISTS conversations; | ||
| DROP TABLE IF EXISTS messages_fts; | ||
| DROP TABLE IF EXISTS memories; | ||
| DROP TABLE IF EXISTS memories_fts; | ||
| `; | ||
| var SCHEMA_SQL = ` | ||
| CREATE TABLE IF NOT EXISTS index_meta ( | ||
| key TEXT PRIMARY KEY, | ||
| value TEXT NOT NULL | ||
| ); | ||
| CREATE TABLE IF NOT EXISTS conversations ( | ||
| docId TEXT PRIMARY KEY, | ||
| platform TEXT NOT NULL, | ||
| title TEXT NOT NULL, | ||
| url TEXT NOT NULL, | ||
| createdAt INTEGER NOT NULL, | ||
| updatedAt INTEGER NOT NULL, | ||
| messageCount INTEGER NOT NULL, | ||
| preview TEXT NOT NULL, | ||
| pinned INTEGER NOT NULL, | ||
| accountId TEXT, | ||
| accountLabel TEXT | ||
| ); | ||
| CREATE INDEX IF NOT EXISTS idx_conversations_platform ON conversations(platform); | ||
| CREATE INDEX IF NOT EXISTS idx_conversations_createdAt ON conversations(createdAt); | ||
| CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts USING fts5( | ||
| docId UNINDEXED, | ||
| msgIndex UNINDEXED, | ||
| role UNINDEXED, | ||
| text, | ||
| tokenize = 'porter unicode61' | ||
| ); | ||
| CREATE TABLE IF NOT EXISTS memories ( | ||
| memoryId TEXT PRIMARY KEY, | ||
| kind TEXT NOT NULL, | ||
| label TEXT NOT NULL, | ||
| text TEXT NOT NULL, | ||
| source TEXT NOT NULL, | ||
| pinned INTEGER NOT NULL, | ||
| enabled INTEGER NOT NULL, | ||
| createdAt INTEGER NOT NULL, | ||
| updatedAt INTEGER NOT NULL | ||
| ); | ||
| CREATE VIRTUAL TABLE IF NOT EXISTS memories_fts USING fts5( | ||
| memoryId UNINDEXED, | ||
| label, | ||
| text, | ||
| tokenize = 'porter unicode61' | ||
| ); | ||
| `; | ||
| var SearchIndex = class { | ||
| db; | ||
| stmts; | ||
| /** Whether the on-disk index was a stale format version when it was opened. */ | ||
| staleOnOpen; | ||
| constructor(dbPath) { | ||
| hardenDbFile(dbPath); | ||
| this.db = new (loadSqlite()).DatabaseSync(dbPath); | ||
| this.db.exec("PRAGMA journal_mode = WAL"); | ||
| this.db.exec("PRAGMA foreign_keys = OFF"); | ||
| this.db.exec(META_SQL); | ||
| this.staleOnOpen = this.getMeta("format_version") !== String(INDEX_FORMAT_VERSION); | ||
| if (this.staleOnOpen) { | ||
| this.db.exec(DROP_DERIVED_SQL); | ||
| } | ||
| this.db.exec(SCHEMA_SQL); | ||
| this.setMeta("format_version", String(INDEX_FORMAT_VERSION)); | ||
| this.stmts = { | ||
| deleteFts: this.db.prepare("DELETE FROM messages_fts WHERE docId = ?"), | ||
| insertFts: this.db.prepare( | ||
| "INSERT INTO messages_fts (docId, msgIndex, role, text) VALUES (?, ?, ?, ?)" | ||
| ), | ||
| deleteConversation: this.db.prepare("DELETE FROM conversations WHERE docId = ?"), | ||
| insertConversation: this.db.prepare( | ||
| `INSERT INTO conversations | ||
| (docId, platform, title, url, createdAt, updatedAt, messageCount, preview, pinned, accountId, accountLabel) | ||
| VALUES (@docId, @platform, @title, @url, @createdAt, @updatedAt, @messageCount, @preview, @pinned, @accountId, @accountLabel)` | ||
| ), | ||
| deleteMemoryFts: this.db.prepare("DELETE FROM memories_fts WHERE memoryId = ?"), | ||
| insertMemoryFts: this.db.prepare( | ||
| "INSERT INTO memories_fts (memoryId, label, text) VALUES (?, ?, ?)" | ||
| ), | ||
| deleteMemory: this.db.prepare("DELETE FROM memories WHERE memoryId = ?"), | ||
| insertMemory: this.db.prepare( | ||
| `INSERT INTO memories | ||
| (memoryId, kind, label, text, source, pinned, enabled, createdAt, updatedAt) | ||
| VALUES (@memoryId, @kind, @label, @text, @source, @pinned, @enabled, @createdAt, @updatedAt)` | ||
| ) | ||
| }; | ||
| } | ||
| /** True when the stored format version was missing or stale — caller rebuilds. */ | ||
| needsRebuild() { | ||
| return this.staleOnOpen; | ||
| } | ||
| /** Number of indexed conversations. */ | ||
| count() { | ||
| const row = this.db.prepare("SELECT COUNT(*) AS n FROM conversations").get(); | ||
| return row.n; | ||
| } | ||
| /** Number of indexed memories. */ | ||
| countMemories() { | ||
| const row = this.db.prepare("SELECT COUNT(*) AS n FROM memories").get(); | ||
| return row.n; | ||
| } | ||
| /** Most-recently-updated conversations (metadata only), newest first. */ | ||
| listRecent(n) { | ||
| return this.listConversations(n, "recent"); | ||
| } | ||
| /** Conversation metadata in a caller-selected, deterministic chronology. */ | ||
| listConversations(n, sort) { | ||
| const limit = Math.max(1, n); | ||
| const orderBy = sort === "oldest" ? "createdAt ASC, updatedAt ASC" : sort === "newest" ? "createdAt DESC, updatedAt DESC" : "updatedAt DESC, createdAt DESC"; | ||
| const rows = this.db.prepare( | ||
| `SELECT docId, platform, title, url, createdAt, updatedAt, messageCount, preview, pinned, | ||
| accountId, accountLabel | ||
| FROM conversations ORDER BY ${orderBy} LIMIT ?` | ||
| ).all(limit); | ||
| return rows.map((r) => { | ||
| const { accountId, accountLabel, ...rest } = r; | ||
| const entry = { ...rest, pinned: r.pinned === 1 }; | ||
| if (accountId) entry.accountId = accountId; | ||
| if (accountLabel) entry.accountLabel = accountLabel; | ||
| return entry; | ||
| }); | ||
| } | ||
| /** | ||
| * Run `fn` inside a transaction, rolling back if it throws. node:sqlite has | ||
| * no transaction() wrapper, so this is the one place BEGIN/COMMIT/ROLLBACK | ||
| * lives for synchronous writes; {@link rebuildFrom} spells it out separately | ||
| * because its body is async and cannot be expressed as a sync callback. | ||
| */ | ||
| transaction(fn) { | ||
| this.db.exec("BEGIN"); | ||
| try { | ||
| const result = fn(); | ||
| this.db.exec("COMMIT"); | ||
| return result; | ||
| } catch (err) { | ||
| this.db.exec("ROLLBACK"); | ||
| throw err; | ||
| } | ||
| } | ||
| upsertConversation(record) { | ||
| this.transaction(() => this.writeRecord(record)); | ||
| } | ||
| upsertMemory(memory) { | ||
| this.transaction(() => this.writeMemory(memory)); | ||
| } | ||
| removeConversation(docId) { | ||
| this.transaction(() => { | ||
| this.stmts.deleteFts.run(docId); | ||
| this.stmts.deleteConversation.run(docId); | ||
| }); | ||
| } | ||
| /** | ||
| * Clear the index and repopulate it from the corpus — proving the corpus is | ||
| * the source of truth. Uses one manual transaction spanning the async file | ||
| * reads, which the sync {@link transaction} helper cannot hold. | ||
| */ | ||
| async rebuildFrom(corpus) { | ||
| this.db.exec("DELETE FROM messages_fts; DELETE FROM conversations;"); | ||
| this.db.exec("DELETE FROM memories_fts; DELETE FROM memories;"); | ||
| let conversations = 0; | ||
| let memories = 0; | ||
| this.db.exec("BEGIN"); | ||
| try { | ||
| for await (const record of corpus.iterate()) { | ||
| this.writeRecord(record); | ||
| conversations += 1; | ||
| } | ||
| for await (const memory of corpus.iterateMemories()) { | ||
| this.writeMemory(memory); | ||
| memories += 1; | ||
| } | ||
| this.db.exec("COMMIT"); | ||
| } catch (err) { | ||
| this.db.exec("ROLLBACK"); | ||
| throw err; | ||
| } | ||
| return { conversations, memories }; | ||
| } | ||
| // Replace a conversation's index rows. Caller supplies the transaction (the | ||
| // transaction() helper for single upserts, or rebuildFrom's manual one). | ||
| writeRecord(rec) { | ||
| const { conversation, messages } = rec; | ||
| this.stmts.deleteFts.run(conversation.docId); | ||
| this.stmts.deleteConversation.run(conversation.docId); | ||
| const titleText = [conversation.title, conversation.preview].filter(Boolean).join("\n"); | ||
| if (titleText.trim()) { | ||
| this.stmts.insertFts.run(conversation.docId, TITLE_ROW_INDEX, TITLE_ROLE, titleText); | ||
| } | ||
| for (const message of messages) { | ||
| if (message.text.trim()) { | ||
| this.stmts.insertFts.run(conversation.docId, message.msgIndex, message.role, message.text); | ||
| } | ||
| } | ||
| this.stmts.insertConversation.run({ | ||
| docId: conversation.docId, | ||
| platform: conversation.platform, | ||
| title: conversation.title, | ||
| url: conversation.url, | ||
| createdAt: conversation.createdAt, | ||
| updatedAt: conversation.updatedAt ?? conversation.contentUpdatedAt ?? 0, | ||
| messageCount: conversation.messageCount ?? messages.length, | ||
| preview: conversation.preview ?? "", | ||
| pinned: conversation.pinned ? 1 : 0, | ||
| accountId: conversation.accountId ?? null, | ||
| accountLabel: conversation.accountLabel ?? null | ||
| }); | ||
| } | ||
| // Replace a memory's index rows. Caller supplies the transaction, as with | ||
| // writeRecord. Memories live in their own tables rather than in messages_fts: | ||
| // a memory is not a turn in a conversation, and mixing them would let a | ||
| // stored note answer a question about what the user actually said. | ||
| writeMemory(memory) { | ||
| this.stmts.deleteMemoryFts.run(memory.memoryId); | ||
| this.stmts.deleteMemory.run(memory.memoryId); | ||
| const searchable = [memory.label, memory.text].filter(Boolean).join("\n"); | ||
| if (searchable.trim()) { | ||
| this.stmts.insertMemoryFts.run(memory.memoryId, memory.label, memory.text); | ||
| } | ||
| this.stmts.insertMemory.run({ | ||
| memoryId: memory.memoryId, | ||
| kind: memory.kind, | ||
| label: memory.label, | ||
| text: memory.text, | ||
| source: memory.source, | ||
| pinned: memory.pinned ? 1 : 0, | ||
| enabled: memory.enabled ? 1 : 0, | ||
| createdAt: memory.createdAt, | ||
| updatedAt: memory.updatedAt | ||
| }); | ||
| } | ||
| search(query, filters = {}) { | ||
| const match = toMatchQuery(query, filters.matchMode); | ||
| if (!match) { | ||
| return []; | ||
| } | ||
| const limit = Math.max(1, filters.limit ?? DEFAULT_LIMIT); | ||
| const clauses = ["messages_fts MATCH @match"]; | ||
| const params = { match }; | ||
| if (filters.platform) { | ||
| clauses.push("c.platform = @platform"); | ||
| params.platform = filters.platform; | ||
| } | ||
| if (filters.title) { | ||
| clauses.push("c.title LIKE @title ESCAPE '\\'"); | ||
| params.title = `%${escapeLike(filters.title)}%`; | ||
| } | ||
| if (typeof filters.dateFrom === "number") { | ||
| clauses.push("c.createdAt >= @dateFrom"); | ||
| params.dateFrom = filters.dateFrom; | ||
| } | ||
| if (typeof filters.dateTo === "number") { | ||
| clauses.push("c.createdAt <= @dateTo"); | ||
| params.dateTo = filters.dateTo; | ||
| } | ||
| const scanCap = Math.min(2e3, Math.max(200, limit * 40)); | ||
| const sql = ` | ||
| SELECT | ||
| c.docId, c.platform, c.title, c.url, c.createdAt, c.updatedAt, | ||
| c.messageCount, c.preview, c.pinned, c.accountId, c.accountLabel, | ||
| m.msgIndex AS msgIndex, m.role AS role, | ||
| snippet(messages_fts, 3, '[', ']', '\u2026', ${SNIPPET_TOKENS}) AS snippet, | ||
| bm25(messages_fts) AS bm25 | ||
| FROM messages_fts m | ||
| JOIN conversations c ON c.docId = m.docId | ||
| WHERE ${clauses.join(" AND ")} | ||
| ORDER BY bm25 | ||
| LIMIT @scanCap | ||
| `; | ||
| const rows = this.db.prepare(sql).all({ ...params, scanCap }); | ||
| const seen = /* @__PURE__ */ new Set(); | ||
| const hits = []; | ||
| for (const row of rows) { | ||
| if (seen.has(row.docId)) { | ||
| continue; | ||
| } | ||
| seen.add(row.docId); | ||
| const hit = { | ||
| docId: row.docId, | ||
| platform: row.platform, | ||
| title: row.title, | ||
| url: row.url, | ||
| createdAt: row.createdAt, | ||
| updatedAt: row.updatedAt, | ||
| messageCount: row.messageCount, | ||
| preview: row.preview, | ||
| pinned: row.pinned === 1, | ||
| // BM25 is negative (more negative = better); expose a positive relevance. | ||
| score: -row.bm25, | ||
| snippet: row.snippet, | ||
| matchRole: normalizeMatchRole(row.role), | ||
| matchMsgIndex: row.msgIndex | ||
| }; | ||
| if (row.accountId) hit.accountId = row.accountId; | ||
| if (row.accountLabel) hit.accountLabel = row.accountLabel; | ||
| hits.push(hit); | ||
| if (hits.length >= limit) { | ||
| break; | ||
| } | ||
| } | ||
| return hits; | ||
| } | ||
| close() { | ||
| this.db.close(); | ||
| } | ||
| getMeta(key) { | ||
| const row = this.db.prepare("SELECT value FROM index_meta WHERE key = ?").get(key); | ||
| return row ? row.value : null; | ||
| } | ||
| setMeta(key, value) { | ||
| this.db.prepare("INSERT INTO index_meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value").run(key, value); | ||
| } | ||
| }; | ||
| async function openSearchIndex(corpus, options = {}) { | ||
| await corpus.init(); | ||
| const index = new SearchIndex(corpus.paths.indexDbFile); | ||
| if (options.autoBuild !== false && (index.needsRebuild() || index.count() === 0 && index.countMemories() === 0)) { | ||
| await index.rebuildFrom(corpus); | ||
| } | ||
| return index; | ||
| } | ||
| async function deleteIndexDb(dbFile) { | ||
| for (const suffix of ["", "-wal", "-shm"]) { | ||
| await fs3.rm(`${dbFile}${suffix}`, { force: true }); | ||
| } | ||
| } | ||
| function toMatchQuery(query, mode = "any") { | ||
| const terms = query.split(/[^\p{L}\p{N}]+/u).filter((t) => t.length > 0); | ||
| if (terms.length === 0) { | ||
| return null; | ||
| } | ||
| if (mode === "phrase") { | ||
| return `"${terms.join(" ")}"`; | ||
| } | ||
| return terms.map((t) => `"${t}"`).join(mode === "all" ? " AND " : " OR "); | ||
| } | ||
| function escapeLike(value) { | ||
| return value.replace(/[\\%_]/g, "\\$&"); | ||
| } | ||
| function normalizeMatchRole(role) { | ||
| if (role === "title") { | ||
| return "title"; | ||
| } | ||
| return role === "assistant" || role === "user" || role === "system" ? role : "unknown"; | ||
| } | ||
| // src/runtimeInstall.ts | ||
| import { cp, mkdir, readdir as readdir2, rename, rm, writeFile } from "fs/promises"; | ||
| import { existsSync } from "fs"; | ||
| import { homedir as homedir2 } from "os"; | ||
| import { dirname as dirname2, join as join5 } from "path"; | ||
| import { fileURLToPath } from "url"; | ||
| function defaultSourceDir() { | ||
| return dirname2(fileURLToPath(import.meta.url)); | ||
| } | ||
| function runtimeBinDir(home = homedir2()) { | ||
| return join5(home, ".llmnesia", "bin"); | ||
| } | ||
| async function installRuntime(options = {}) { | ||
| const home = options.home ?? homedir2(); | ||
| const sourceDir = options.sourceDir ?? defaultSourceDir(); | ||
| const nodePath = options.nodePath ?? process.execPath; | ||
| const binDir = runtimeBinDir(home); | ||
| const sourceEntry = join5(sourceDir, "cli.js"); | ||
| if (!existsSync(sourceEntry)) { | ||
| throw new Error( | ||
| `Cannot locate the built runtime: no cli.js in ${sourceDir}. Run this from an installed package (npx -y @llmnesia/mcp@latest install) or build first (npm run build).` | ||
| ); | ||
| } | ||
| const parentDir = dirname2(binDir); | ||
| const nonce = `${process.pid}-${Date.now()}`; | ||
| const stagingDir = join5(parentDir, `.bin-install-${nonce}`); | ||
| const backupDir = join5(parentDir, `.bin-backup-${nonce}`); | ||
| await mkdir(parentDir, { recursive: true }); | ||
| await rm(stagingDir, { recursive: true, force: true }); | ||
| await rm(backupDir, { recursive: true, force: true }); | ||
| await mkdir(stagingDir, { recursive: true }); | ||
| try { | ||
| for (const name of await readdir2(sourceDir)) { | ||
| await cp(join5(sourceDir, name), join5(stagingDir, name), { recursive: true }); | ||
| } | ||
| await writeFile( | ||
| join5(stagingDir, "package.json"), | ||
| `${JSON.stringify( | ||
| { | ||
| name: "llmnesia-mcp-runtime", | ||
| private: true, | ||
| type: "module", | ||
| // Stamped so the installed build is identifiable without spawning it. | ||
| // Its absence is how a machine sat on 0.1.3 while npm was ten | ||
| // releases ahead and nothing, doctor included, could say so. | ||
| version: options.version ?? VERSION, | ||
| comment: "Installed copy of @llmnesia/mcp's dist. Managed by the LLMnesia installer; edits are lost on reinstall." | ||
| }, | ||
| null, | ||
| 2 | ||
| )} | ||
| `, | ||
| "utf8" | ||
| ); | ||
| const hadPreviousRuntime = existsSync(binDir); | ||
| if (hadPreviousRuntime) { | ||
| await rename(binDir, backupDir); | ||
| } | ||
| try { | ||
| await rename(stagingDir, binDir); | ||
| } catch (err) { | ||
| if (hadPreviousRuntime && !existsSync(binDir) && existsSync(backupDir)) { | ||
| await rename(backupDir, binDir); | ||
| } | ||
| throw err; | ||
| } | ||
| await rm(backupDir, { recursive: true, force: true }); | ||
| } finally { | ||
| await rm(stagingDir, { recursive: true, force: true }); | ||
| } | ||
| const entry = join5(binDir, "cli.js"); | ||
| return { | ||
| binDir, | ||
| entry, | ||
| serve: { command: nodePath, args: [entry, "serve"] }, | ||
| sourceDir, | ||
| nodePath | ||
| }; | ||
| } | ||
| async function uninstallRuntime(home = homedir2()) { | ||
| const binDir = runtimeBinDir(home); | ||
| if (!existsSync(binDir)) { | ||
| return null; | ||
| } | ||
| await rm(binDir, { recursive: true, force: true }); | ||
| return binDir; | ||
| } | ||
| // src/selfUpdate.ts | ||
| import { spawn } from "child_process"; | ||
| import { mkdtemp, readFile, rm as rm2, writeFile as writeFile2 } from "fs/promises"; | ||
| import { dirname as dirname3, join as join6 } from "path"; | ||
| import { homedir as homedir3, tmpdir } from "os"; | ||
| import { fileURLToPath as fileURLToPath2 } from "url"; | ||
| var PACKAGE_NAME = "@llmnesia/mcp"; | ||
| var REGISTRY_LATEST_URL = "https://registry.npmjs.org/@llmnesia%2Fmcp/latest"; | ||
| var UPDATE_CHECK_INTERVAL_MS = 24 * 60 * 60 * 1e3; | ||
| var REGISTRY_TIMEOUT_MS = 5e3; | ||
| var INSTALL_TIMEOUT_MS = 12e4; | ||
| var UPDATE_CHECK_DELAY_MS = 3e4; | ||
| function stateFile(home) { | ||
| return join6(home, ".llmnesia", "update-state.json"); | ||
| } | ||
| async function readState(home) { | ||
| try { | ||
| const parsed = JSON.parse(await readFile(stateFile(home), "utf8")); | ||
| return typeof parsed === "object" && parsed !== null ? parsed : {}; | ||
| } catch { | ||
| return {}; | ||
| } | ||
| } | ||
| async function writeState(home, state) { | ||
| try { | ||
| await writeFile2(stateFile(home), `${JSON.stringify(state, null, 2)} | ||
| `, "utf8"); | ||
| } catch { | ||
| } | ||
| } | ||
| function compareVersions(a, b) { | ||
| const parse = (v) => { | ||
| const match = /^(\d+)\.(\d+)\.(\d+)/.exec(v.trim()); | ||
| return match ? [Number(match[1]), Number(match[2]), Number(match[3])] : []; | ||
| }; | ||
| const left = parse(a); | ||
| const right = parse(b); | ||
| if (left.length === 0 || right.length === 0) { | ||
| return 0; | ||
| } | ||
| for (let i = 0; i < 3; i += 1) { | ||
| if (left[i] !== right[i]) { | ||
| return left[i] - right[i]; | ||
| } | ||
| } | ||
| return 0; | ||
| } | ||
| function isStableRelease(version) { | ||
| return /^\d+\.\d+\.\d+$/.test(version.trim()); | ||
| } | ||
| function isSelfUpdateDisabled(env = process.env) { | ||
| const raw = env.LLMNESIA_NO_AUTO_UPDATE; | ||
| return raw === "1" || raw === "true"; | ||
| } | ||
| async function fetchLatestFromRegistry() { | ||
| const controller = new AbortController(); | ||
| const timer = setTimeout(() => controller.abort(), REGISTRY_TIMEOUT_MS); | ||
| try { | ||
| const response = await fetch(REGISTRY_LATEST_URL, { | ||
| signal: controller.signal, | ||
| headers: { accept: "application/json" } | ||
| }); | ||
| if (!response.ok) { | ||
| throw new Error(`registry responded ${response.status}`); | ||
| } | ||
| const body = await response.json(); | ||
| if (typeof body.version !== "string") { | ||
| throw new Error("registry response had no version"); | ||
| } | ||
| return body.version; | ||
| } finally { | ||
| clearTimeout(timer); | ||
| } | ||
| } | ||
| function npmCommand() { | ||
| const binDir = dirname3(process.execPath); | ||
| return process.platform === "win32" ? join6(binDir, "npm.cmd") : join6(binDir, "npm"); | ||
| } | ||
| async function npmDownload(version) { | ||
| const prefix = await mkdtemp(join6(tmpdir(), "llmnesia-update-")); | ||
| await new Promise((resolve2, reject) => { | ||
| const child = spawn( | ||
| npmCommand(), | ||
| [ | ||
| "install", | ||
| `${PACKAGE_NAME}@${version}`, | ||
| "--prefix", | ||
| prefix, | ||
| "--no-audit", | ||
| "--no-fund", | ||
| "--no-save", | ||
| "--loglevel", | ||
| "error" | ||
| ], | ||
| { stdio: ["ignore", "ignore", "pipe"], windowsHide: true } | ||
| ); | ||
| let stderr = ""; | ||
| child.stderr?.on("data", (chunk) => { | ||
| stderr += String(chunk); | ||
| }); | ||
| const timer = setTimeout(() => child.kill(), INSTALL_TIMEOUT_MS); | ||
| child.once("error", (err) => { | ||
| clearTimeout(timer); | ||
| reject(err); | ||
| }); | ||
| child.once("close", (code) => { | ||
| clearTimeout(timer); | ||
| if (code === 0) { | ||
| resolve2(); | ||
| } else { | ||
| reject(new Error(`npm install exited ${code}${stderr ? `: ${stderr.trim()}` : ""}`)); | ||
| } | ||
| }); | ||
| }); | ||
| return join6(prefix, "node_modules", PACKAGE_NAME, "dist"); | ||
| } | ||
| async function maybeSelfUpdate(options = {}) { | ||
| const env = options.env ?? process.env; | ||
| if (isSelfUpdateDisabled(env)) { | ||
| return { action: "disabled" }; | ||
| } | ||
| const home = options.home ?? homedir3(); | ||
| const currentVersion = options.currentVersion ?? VERSION; | ||
| const runningDir = options.runningDir ?? dirname3(fileURLToPath2(import.meta.url)); | ||
| if (runningDir !== runtimeBinDir(home)) { | ||
| return { action: "not-installed-runtime" }; | ||
| } | ||
| if (!isStableRelease(currentVersion)) { | ||
| return { action: "not-a-release-build", currentVersion }; | ||
| } | ||
| const now = options.now ?? Date.now; | ||
| const state = await readState(home); | ||
| if (state.lastCheckedAt !== void 0 && now() - state.lastCheckedAt < UPDATE_CHECK_INTERVAL_MS) { | ||
| return { action: "throttled" }; | ||
| } | ||
| let latest; | ||
| try { | ||
| latest = await (options.fetchLatestVersion ?? fetchLatestFromRegistry)(); | ||
| } catch (err) { | ||
| await writeState(home, { ...state, lastCheckedAt: now() }); | ||
| return { action: "check-failed", error: err instanceof Error ? err.message : String(err) }; | ||
| } | ||
| const checked = { ...state, lastCheckedAt: now(), lastSeenVersion: latest }; | ||
| if (!isStableRelease(latest) || compareVersions(latest, currentVersion) <= 0) { | ||
| await writeState(home, checked); | ||
| return { action: "up-to-date", latest }; | ||
| } | ||
| if (state.lastFailedVersion === latest) { | ||
| await writeState(home, checked); | ||
| return { action: "skipped-failed-before", latest }; | ||
| } | ||
| let sourceDir; | ||
| try { | ||
| sourceDir = await (options.downloadPackage ?? npmDownload)(latest); | ||
| await installRuntime({ home, sourceDir, version: latest }); | ||
| await writeState(home, { ...checked, lastUpdatedTo: latest, lastFailedVersion: void 0 }); | ||
| return { action: "updated", from: currentVersion, to: latest }; | ||
| } catch (err) { | ||
| await writeState(home, { ...checked, lastFailedVersion: latest }); | ||
| return { action: "update-failed", latest, error: err instanceof Error ? err.message : String(err) }; | ||
| } finally { | ||
| if (sourceDir) { | ||
| await rm2(join6(sourceDir, "..", "..", ".."), { recursive: true, force: true }).catch(() => { | ||
| }); | ||
| } | ||
| } | ||
| } | ||
| function scheduleSelfUpdate(options = {}) { | ||
| const { delayMs = UPDATE_CHECK_DELAY_MS, onOutcome, ...rest } = options; | ||
| const timer = setTimeout(() => { | ||
| void maybeSelfUpdate(rest).then( | ||
| (outcome) => onOutcome?.(outcome), | ||
| () => { | ||
| } | ||
| ); | ||
| }, delayMs); | ||
| timer.unref?.(); | ||
| } | ||
| function describeOutcome(outcome) { | ||
| switch (outcome.action) { | ||
| case "updated": | ||
| return `updated ${outcome.from} to ${outcome.to}; it takes effect next launch`; | ||
| case "update-failed": | ||
| return `could not install ${outcome.latest}: ${outcome.error}`; | ||
| case "check-failed": | ||
| return `could not reach the npm registry: ${outcome.error}`; | ||
| case "up-to-date": | ||
| return `already current (latest published is ${outcome.latest})`; | ||
| case "skipped-failed-before": | ||
| return `${outcome.latest} failed to install previously; not retrying automatically`; | ||
| case "throttled": | ||
| return "checked recently"; | ||
| case "disabled": | ||
| return "automatic updates are turned off (LLMNESIA_NO_AUTO_UPDATE)"; | ||
| case "not-installed-runtime": | ||
| return "not running from ~/.llmnesia/bin, so nothing to update"; | ||
| case "not-a-release-build": | ||
| return `running an unreleased build (${outcome.currentVersion}); leaving the installed runtime alone`; | ||
| } | ||
| } | ||
| export { | ||
| resolveCorpusDir, | ||
| staleIndexFiles, | ||
| Corpus, | ||
| ingestFile, | ||
| pendingInboxFiles, | ||
| drainInbox, | ||
| SearchIndex, | ||
| openSearchIndex, | ||
| deleteIndexDb, | ||
| runtimeBinDir, | ||
| installRuntime, | ||
| uninstallRuntime, | ||
| isSelfUpdateDisabled, | ||
| scheduleSelfUpdate, | ||
| describeOutcome | ||
| }; |
Sorry, the diff of this file is too big to display
Sorry, the diff of this file is not supported yet
Sorry, the diff of this file is too big to display
Filesystem access
Supply chain riskAccesses the file system, and could potentially read sensitive data.
Long strings
Supply chain riskContains long string literals, which may be a sign of obfuscated or packed code.
URL strings
Supply chain riskPackage contains fragments of external URLs or IP addresses, which the package may be accessing at runtime.
Long strings
Supply chain riskContains long string literals, which may be a sign of obfuscated or packed code.
URL strings
Supply chain riskPackage contains fragments of external URLs or IP addresses, which the package may be accessing at runtime.
1426923
1.71%28456
1.84%275
10.44%60
1.69%