diff --git a/package.json b/package.json index aa12931..682ee25 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@ztimson/ai-utils", - "version": "1.3.0", + "version": "1.3.1", "description": "AI Utility library", "author": "Zak Timson", "license": "MIT", diff --git a/src/llm.ts b/src/llm.ts index a479155..015c707 100644 --- a/src/llm.ts +++ b/src/llm.ts @@ -340,7 +340,7 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam if(typeof resp === 'string' && !resp.trim() && lastDelegateResp !== null) resp = lastDelegateResp; // Trim memory injections from history - if(mem?.tool) history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'recall')); + if(mem?.tool) history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'memory_recall')); // Auto-memorize before compressing if(options.compress && this.estimateTokens(history) >= options.compress.max) { diff --git a/src/memory.ts b/src/memory.ts index a7d2be9..ae96d98 100644 --- a/src/memory.ts +++ b/src/memory.ts @@ -221,6 +221,8 @@ function getWeekSunday(monday: string): string { } export class MemoryManager { + private recentlyTouched = new Map(); + private pendingMemorizations = new Map m.name === args.name); if (!mem) return 'Document not found'; + this.touch(mem.name); return mem.content; }, }), @@ -292,6 +295,96 @@ ${conversation}`; }; } + private applyHeader(content: string, header: string): string { + return `${header}\n\n${this.stripHeader(content)}`; + } + + private async backgroundMemorization(conversation: string, memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string): Promise { + const mem = memories instanceof MemoryCache ? memories.memories : memories; + const monday = getWeekMonday(); + const sunday = getWeekSunday(monday); + const buckets = await this.factAgent(conversation, mem, options, monday); + if(!buckets.length) return; + const jobs = [...buckets].map(({subject, facts}) => { + let node = mem.find(m => m.name === subject); + if(!node) { + node = {name: subject, description: '', content: '', embedding: [],}; + mem.push(node); + } + const week = subject.startsWith('Journal/') ? {monday, sunday} : undefined; + return this.enqueue(node, facts, mem, options, tempName, week); + }); + await Promise.all(jobs); + } + + private buildHeader(node: Memory, week?: {monday: string, sunday: string}, links: string[] = [], backlinks: string[] = []): string { + const tags = node.name.split('/')[0]?.toLowerCase(); + const lines = [ + '---', + `name: ${node.name}`, + `description: ${node.description || ''}`, + tags ? `tags: [${tags}]` : '', + links.length ? `links: [${links.map(l => `"${l}"`).join(', ')}]` : 'links: []', + backlinks.length ? `backlinks: [${backlinks.map(l => `"${l}"`).join(', ')}]` : 'backlinks: []', + week ? `week: ${week.monday} – ${week.sunday}` : '', + `modified: ${new Date().toISOString()}`, + '---', + ].filter(Boolean); + return lines.join('\n'); + } + + private cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] { + const scored = memories + .filter(m => m.embedding?.length) + .map(m => ({ + ref: {name: m.name, description: m.description}, + distance: cosineDistance(query, m.embedding), + })) + .sort((a, b) => a.distance - b.distance) + .slice(0, limit); + return scored.map(s => s.ref); + } + + /** + * Coalescing queue: if a doc is already compiling, abort the in-flight run, merge its + * facts with the new ones and restart. Never blocks a pending update, never drops facts. + */ + private enqueue(node: Memory, facts: string[], memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string, week?: {monday: string, sunday: string}): Promise { + const key = node.name; + const existing = this.queues.get(key); + if (existing) { + existing.pending.push(...facts); + existing.request?.abort?.(); + return existing.task; + } + + const entry: {pending: string[], request: {abort?: () => void} | null, task: Promise} = {pending: [...facts], request: null, task: Promise.resolve()}; + this.queues.set(key, entry); + const m = memories instanceof MemoryCache ? memories.memories : memories; + entry.task = (async () => { + while (entry.pending.length) { + const batch = dedupeFacts(entry.pending.splice(0, entry.pending.length)); + const written = await this.docAgent(node, batch, m, options, tempName, week, entry); + if (!written) entry.pending.unshift(...batch); + } + })().finally(() => { + this.queues.delete(key); + if(!this.queues.size && memories instanceof MemoryCache) memories.rebuild(); + }); + return entry.task; + } + + private listNodes(memories: Memory[]): MemoryRef[] { + return memories.map(m => ({name: m.name, description: m.description})); + } + + decay() { + for(const [name, ttl] of this.recentlyTouched) { + if(ttl <= 1) this.recentlyTouched.delete(name); + else this.recentlyTouched.set(name, ttl - 1); + } + } + forget(name: string, memories: Memory[] | MemoryCache): boolean { const mem = memories instanceof MemoryCache ? memories.memories : memories; const idx = mem.findIndex(m => m.name === name); @@ -316,20 +409,43 @@ ${conversation}`; return true; } - private cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] { - const scored = memories - .filter(m => m.embedding?.length) - .map(m => ({ - ref: {name: m.name, description: m.description}, - distance: cosineDistance(query, m.embedding), - })) - .sort((a, b) => a.distance - b.distance) - .slice(0, limit); - return scored.map(s => s.ref); + getTouched(): string[] { + return [...this.recentlyTouched.keys()]; } - private listNodes(memories: Memory[]): MemoryRef[] { - return memories.map(m => ({name: m.name, description: m.description})); + async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise { + const conversation = history + .filter(h => h.role === 'user' || h.role === 'assistant') + .map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim(); + if(!conversation) return []; + + const trackingId = `${Date.now()}_${Math.random()}`; + const tempMemory = await this.createTempMemory(conversation); + const mem = memories instanceof MemoryCache ? memories.memories : memories; + mem.push(tempMemory); + if (memories instanceof MemoryCache) memories.rebuild(); + this.pendingMemorizations.set(trackingId, { + memories, + tempMemoryName: tempMemory.name, + timestamp: Date.now(), + }); + + try { + await this.backgroundMemorization(conversation, memories, options, tempMemory.name); + const finalMem = memories instanceof MemoryCache ? memories.memories : memories; + return finalMem.filter(m => !m.name.startsWith('_temp_')); + } finally { + const pending = this.pendingMemorizations.get(trackingId); + if (pending) { + const cleanMem = pending.memories instanceof MemoryCache + ? pending.memories.memories + : pending.memories; + const idx = cleanMem.findIndex(m => m.name === pending.tempMemoryName); + if (idx !== -1) cleanMem.splice(idx, 1); + if (pending.memories instanceof MemoryCache) pending.memories.rebuild(); + } + this.pendingMemorizations.delete(trackingId); + } } async recollect(query: string, memories: Memory[] | MemoryCache, limit = 5, graphDepth = 1): Promise { @@ -370,106 +486,8 @@ ${conversation}`; return ordered.map(n => mem.find(m => m.name === n)!).filter(Boolean); } - async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise { - const conversation = history - .filter(h => h.role === 'user' || h.role === 'assistant') - .map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim(); - if(!conversation) return []; - - const trackingId = `${Date.now()}_${Math.random()}`; - const tempMemory = await this.createTempMemory(conversation); - const mem = memories instanceof MemoryCache ? memories.memories : memories; - mem.push(tempMemory); - if (memories instanceof MemoryCache) memories.rebuild(); - this.pendingMemorizations.set(trackingId, { - memories, - tempMemoryName: tempMemory.name, - timestamp: Date.now(), - }); - - try { - await this._memorizeBackground(conversation, memories, options, tempMemory.name); - const finalMem = memories instanceof MemoryCache ? memories.memories : memories; - return finalMem.filter(m => !m.name.startsWith('_temp_')); - } finally { - const pending = this.pendingMemorizations.get(trackingId); - if (pending) { - const cleanMem = pending.memories instanceof MemoryCache - ? pending.memories.memories - : pending.memories; - const idx = cleanMem.findIndex(m => m.name === pending.tempMemoryName); - if (idx !== -1) cleanMem.splice(idx, 1); - if (pending.memories instanceof MemoryCache) pending.memories.rebuild(); - } - this.pendingMemorizations.delete(trackingId); - } - } - - private async _memorizeBackground(conversation: string, memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string): Promise { - const mem = memories instanceof MemoryCache ? memories.memories : memories; - const monday = getWeekMonday(); - const sunday = getWeekSunday(monday); - const buckets = await this.factAgent(conversation, mem, options, monday); - if(!buckets.length) return; - const jobs = [...buckets].map(({subject, facts}) => { - let node = mem.find(m => m.name === subject); - if(!node) { - node = {name: subject, description: '', content: '', embedding: [],}; - mem.push(node); - } - const week = subject.startsWith('Journal/') ? {monday, sunday} : undefined; - return this.enqueue(node, facts, mem, options, tempName, week); - }); - await Promise.all(jobs); - } - - /** - * Coalescing queue: if a doc is already compiling, abort the in-flight run, merge its - * facts with the new ones and restart. Never blocks a pending update, never drops facts. - */ - private enqueue(node: Memory, facts: string[], memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string, week?: {monday: string, sunday: string}): Promise { - const key = node.name; - const existing = this.queues.get(key); - if (existing) { - existing.pending.push(...facts); - existing.request?.abort?.(); - return existing.task; - } - - const entry: {pending: string[], request: {abort?: () => void} | null, task: Promise} = {pending: [...facts], request: null, task: Promise.resolve()}; - this.queues.set(key, entry); - const m = memories instanceof MemoryCache ? memories.memories : memories; - entry.task = (async () => { - while (entry.pending.length) { - const batch = dedupeFacts(entry.pending.splice(0, entry.pending.length)); - const written = await this.docAgent(node, batch, m, options, tempName, week, entry); - if (!written) entry.pending.unshift(...batch); - } - })().finally(() => { - this.queues.delete(key); - if(!this.queues.size && memories instanceof MemoryCache) memories.rebuild(); - }); - return entry.task; - } - - private buildHeader(node: Memory, week?: {monday: string, sunday: string}, links: string[] = [], backlinks: string[] = []): string { - const tags = node.name.split('/')[0]?.toLowerCase(); - const lines = [ - '---', - `name: ${node.name}`, - `description: ${node.description || ''}`, - tags ? `tags: [${tags}]` : '', - links.length ? `links: [${links.map(l => `"${l}"`).join(', ')}]` : 'links: []', - backlinks.length ? `backlinks: [${backlinks.map(l => `"${l}"`).join(', ')}]` : 'backlinks: []', - week ? `week: ${week.monday} – ${week.sunday}` : '', - `modified: ${new Date().toISOString()}`, - '---', - ].filter(Boolean); - return lines.join('\n'); - } - - private applyHeader(content: string, header: string): string { - return `${header}\n\n${this.stripHeader(content)}`; + touch(name: string, ttl = 2) { + this.recentlyTouched.set(name, ttl); } private updateFrontmatter(content: string, updates: {links?: string[], backlinks?: string[]}): string {