Compare commits

...

7 Commits
1.3.4 ... 1.4.1

Author SHA1 Message Date
566d84fd7a Added memory graph traversal helpers
All checks were successful
Publish Library / Build NPM Project (push) Successful in 43s
Publish Library / Tag Version (push) Successful in 14s
2026-08-04 12:58:39 -04:00
4230b534fc bump 1.4.0
All checks were successful
Publish Library / Build NPM Project (push) Successful in 1m17s
Publish Library / Tag Version (push) Successful in 14s
2026-08-04 12:44:41 -04:00
119f8472f2 token pools
Some checks failed
Publish Library / Tag Version (push) Has been cancelled
Publish Library / Build NPM Project (push) Has been cancelled
2026-08-04 12:44:21 -04:00
9c04e58c63 Pass deligate subagents full history, improved memory managment 2026-08-04 12:24:23 -04:00
7fbb42c26a improved subagent instructions 2026-08-04 12:03:31 -04:00
be08db8e2c Attach tps to response promise
All checks were successful
Publish Library / Build NPM Project (push) Successful in 42s
Publish Library / Tag Version (push) Successful in 9s
2026-08-04 09:48:20 -04:00
497f051c62 bump 1.3.5
All checks were successful
Publish Library / Build NPM Project (push) Successful in 52s
Publish Library / Tag Version (push) Successful in 7s
2026-08-04 09:30:45 -04:00
10 changed files with 822 additions and 329 deletions

View File

@@ -1,6 +1,6 @@
{ {
"name": "@ztimson/ai-utils", "name": "@ztimson/ai-utils",
"version": "1.3.4", "version": "1.4.1",
"description": "AI Utility library", "description": "AI Utility library",
"author": "Zak Timson", "author": "Zak Timson",
"license": "MIT", "license": "MIT",

View File

@@ -1,16 +1,27 @@
import {Anthropic as anthropic} from '@anthropic-ai/sdk'; import {Anthropic as anthropic} from '@anthropic-ai/sdk';
import {findByProp, objectMap, JSONSanitize, JSONAttemptParse} from '@ztimson/utils'; import {findByProp, objectMap, JSONSanitize, JSONAttemptParse, makeArray} from '@ztimson/utils';
import {AbortablePromise, Ai} from './ai.ts'; import {AbortablePromise, Ai} from './ai.ts';
import {LLMMessage, LLMRequest} from './llm.ts'; import {LLMMessage, LLMRequest} from './llm.ts';
import {LLMProvider} from './provider.ts'; import {LLMProvider} from './provider.ts';
import {TokenPool} from './token-pool.ts';
import {convertSchema} from './tools.ts'; import {convertSchema} from './tools.ts';
export class Anthropic extends LLMProvider { export class Anthropic extends LLMProvider {
client!: anthropic; private clients = new Map<string, anthropic>();
tokenPool!: TokenPool;
constructor(public readonly ai: Ai, public readonly apiToken: string, public model: string) { constructor(public readonly ai: Ai, public readonly apiToken: string | string[], public model: string) {
super(); super();
this.client = new anthropic({apiKey: apiToken}); this.tokenPool = new TokenPool(...makeArray(apiToken).filter(Boolean));
}
private getClient(token: string): anthropic {
let client = this.clients.get(token);
if(!client) {
client = new anthropic({apiKey: token});
this.clients.set(token, client);
}
return client;
} }
private toStandard(history: any[]): LLMMessage[] { private toStandard(history: any[]): LLMMessage[] {
@@ -90,7 +101,7 @@ export class Anthropic extends LLMProvider {
do { do {
requestParams.messages = history.map(({timestamp, ...m}) => m); requestParams.messages = history.map(({timestamp, ...m}) => m);
const callStart = Date.now(); const callStart = Date.now();
resp = await this.client.messages.create(requestParams).catch(err => { resp = await this.tokenPool.run(token => this.getClient(token).messages.create(requestParams)).catch(err => {
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`; err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
throw err; throw err;
}); });

72
src/helpers.ts Normal file
View File

@@ -0,0 +1,72 @@
import {Memory, MemoryCache} from './memory.ts';
export type MemoryNode = {
name: string;
missing: boolean;
links: string[];
backlinks: string[];
}
export function buildMemoryGraph(memories: Memory[] | MemoryCache): MemoryNode[] {
const mems = memories instanceof MemoryCache ? memories.memories : memories;
const nameSet = new Set(mems.map(m => m.name));
const ghosts = new Set<string>();
const nodes: MemoryNode[] = mems.map(m => ({
name: m.name,
missing: false,
links: m.links,
backlinks: m.backlinks,
}));
for (const node of nodes) {
for (const link of node.links) {
if (!nameSet.has(link)) ghosts.add(link);
}
}
return [
...nodes,
...[...ghosts].map(name => ({
name,
missing: true,
links: [],
backlinks: nodes
.filter(n => n.links.includes(name))
.map(n => n.name),
}))
];
}
export function renderMemoryGraph(nodes) {
if (!nodes.length) return 'No memories yet.';
const groups = new Map();
for (const node of nodes) {
const [prefix, ...rest] = node.name.split('/');
const group = rest.length ? prefix : 'Root';
const label = rest.length ? rest.join('/') : node.name;
if (!groups.has(group)) groups.set(group, []);
groups.get(group).push({...node, label});
}
const ghostCount = nodes.filter(n => n.missing).length;
const lines = [`Memory Graph (${nodes.length} nodes, ${ghostCount} ghost${ghostCount === 1 ? '' : 's'})`, ''];
for (const group of [...groups.keys()].sort()) {
const items = groups.get(group).sort((a, b) => a.label.localeCompare(b.label));
lines.push(`${group}/`);
items.forEach((n, i) => {
const last = i === items.length - 1;
const branch = last ? '└─' : '├─';
const pad = last ? ' ' : '│ ';
const tag = n.missing ? ' (ghost)' : '';
lines.push(` ${branch} ${n.label}${tag}`);
if (n.links.length) lines.push(` ${pad}${n.links.join(', ')}`);
if (n.backlinks.length) lines.push(` ${pad}${n.backlinks.join(', ')}`);
});
lines.push('');
}
return lines.join('\n').trimEnd();
}

View File

@@ -1,9 +1,11 @@
export * from './ai'; export * from './ai';
export * from './antrhopic'; export * from './antrhopic';
export * from './audio'; export * from './audio';
export * from './helpers';
export * from './llm'; export * from './llm';
export * from './memory'; export * from './memory';
export * from './open-ai'; export * from './open-ai';
export * from './provider'; export * from './provider';
export * from './token-pool'
export * from './tools'; export * from './tools';
export * from './vision'; export * from './vision';

View File

@@ -9,8 +9,10 @@ import {dirname, join} from 'path';
import {spawn} from 'node:child_process'; import {spawn} from 'node:child_process';
import {Memory, MemoryCache, MemoryManager, MemoryOptions} from './memory.ts'; import {Memory, MemoryCache, MemoryManager, MemoryOptions} from './memory.ts';
export type AnthropicConfig = {proto: 'anthropic', token: string}; const MAX_AGENT_DEPTH = 5;
export type OpenAiConfig = {proto: 'openai', host?: string, token: string};
export type AnthropicConfig = {proto: 'anthropic', token: string | string[]};
export type OpenAiConfig = {proto: 'openai', host?: string, token: string | string[]};
export type Agent = { export type Agent = {
name: string; name: string;
@@ -104,8 +106,6 @@ export type Skill = {
content: string; content: string;
} }
const MAX_AGENT_DEPTH = 5;
class LLM { class LLM {
private memoryManager!: MemoryManager; private memoryManager!: MemoryManager;
@@ -122,35 +122,33 @@ class LLM {
this.memoryManager = new MemoryManager(this); this.memoryManager = new MemoryManager(this);
} }
private setupAgent(agents: Agent[] = [], allAgents: Agent[], pending: Map<string, any>, aborts: (() => void)[], depth = 0): AiTool[] { private setupAgent(agents: Agent[] = [], allAgents: Agent[], history: LLMMessage[], aborts: (() => void)[], depth = 0, delegateState: {resp: string | null}): AiTool[] {
return agents.map(a => { return agents.map(a => {
const toolName = `${a.delegate ? '' : 'sub'}agent_${snakeCase(a.name)}`; const toolName = `${a.delegate ? '' : 'sub'}agent_${snakeCase(a.name)}`;
return { return {
name: toolName, name: toolName,
description: `${a.delegate ? 'Delegate to ' : ''}Subagent: ${a.description || a.name}`, description: `${a.delegate ? 'Delegate to ' : ''}Subagent: ${a.description || a.name}`,
args: { args: <any>(a.delegate ? {} : {
context: {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true}, context: {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true},
instructions: {type: 'string', description: 'Detailed instructions for subagent to complete', required: true}, instructions: {type: 'string', description: 'Detailed instructions for subagent to complete', required: true},
}, }),
fn: async (args: any, stream: any, ai: any, id?: string) => { fn: async (args: any, stream: any, ai: any, id?: string) => {
if(depth >= MAX_AGENT_DEPTH) return 'Max agent delegation depth exceeded'; if(depth >= MAX_AGENT_DEPTH) return 'Max agent delegation depth exceeded';
const subHistory: LLMMessage[] = [];
// Opt-in only, self always excluded regardless of whitelist // Opt-in only, self always excluded regardless of whitelist
const nested = (a.agents || []) const nested = (a.agents || [])
.map(name => allAgents.find(x => x.name === name)) .map(name => allAgents.find(x => x.name === name))
.filter((x): x is Agent => !!x && x.name !== a.name); .filter((x): x is Agent => !!x && x.name !== a.name);
const start = Date.now(); const request = this.ask(a.delegate ? '' : `${args.instructions}${args.context ? `\n\n<context>${args.context}</context>` : ''}`, {
const request = this.ask(`${args.instructions}${args.context ? `\n\n<context>${args.context}</context>` : ''}`, { system: `You are a specialized subagent. ${a.delegate ? 'Your output streams directly to the user for the remainder of this turn. You are mid conversation - dispense with greetings.' : 'You are wrapped in a tool call that will be analysis by an LLM - dispense with conversation'}
system: `You are a specialized subagent. ${a.delegate ? 'Your output streams directly to the user for the remainder of this turn.' : 'You are wrapped in a tool call that will be analysis by an LLM'}
As a subagent, focus on executing your task completely using available tools and returning only the final result - no commentary, questions, or dialogue. As a subagent, focus on executing your task completely using available tools and returning only the final result - no commentary, questions, or dialogue.
${a.system}`, ${a.system}`,
model: a.model || undefined, model: a.model || undefined,
temperature: a.temperature, temperature: a.temperature,
stream: a.delegate ? stream : undefined, stream: a.delegate ? stream : undefined,
history: subHistory, history: a.delegate ? history : [],
mcp: a.mcp || undefined, mcp: a.mcp || undefined,
skills: a.skills || undefined, skills: a.skills || undefined,
tools: a.tools || undefined, tools: a.tools || undefined,
@@ -159,14 +157,9 @@ ${a.system}`,
} as any); } as any);
aborts.push(request.abort); aborts.push(request.abort);
const resp = await request; const resp = await request;
const duration = Date.now() - start;
const assistantTurns = subHistory.filter((h: any) => h.role === 'assistant' && h.duration);
const genTime = assistantTurns.reduce((s, h: any) => s + h.duration, 0);
const genTokens = assistantTurns.reduce((s, h: any) => s + (h.tps || 0) * (h.duration / 1000), 0);
const tps = genTime > 0 ? genTokens / (genTime / 1000) : 0;
if(a.delegate) { if(a.delegate) {
pending.set(<string>id, {resp, subHistory, duration, tps}); delegateState.resp = resp;
return ''; return '';
} }
return resp; return resp;
@@ -266,7 +259,10 @@ ${a.system}`,
nestedAborts.forEach(a => a()); nestedAborts.forEach(a => a());
}; };
const promise = (async () => { let promise: any;
const requestStart = Date.now();
promise = (async () => {
let tools: AiTool[] = options.tools || this.ai.options.llm?.tools || []; let tools: AiTool[] = options.tools || this.ai.options.llm?.tools || [];
const prompts: string[] = []; const prompts: string[] = [];
let history = options.history || []; let history = options.history || [];
@@ -289,8 +285,8 @@ ${a.system}`,
// Agents // Agents
const agents = options.agents || this.ai.options?.llm?.agents; const agents = options.agents || this.ai.options?.llm?.agents;
const pendingDelegates = new Map<string, any>(); const delegateState: {resp: string | null} = {resp: null};
if(agents?.length) tools.push(...this.setupAgent(agents, agents, pendingDelegates, nestedAborts, options._agentDepth || 0)); if(agents?.length) tools.push(...this.setupAgent(agents, agents, history, nestedAborts, options._agentDepth || 0, delegateState));
// Memory // Memory
const mem = MemoryManager.normalize(options.memory); const mem = MemoryManager.normalize(options.memory);
@@ -332,48 +328,36 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam
if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'}); if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
// Time each tool call's real execution so its history entry gets its own duration/tps
const toolTimings = new Map<string, {duration: number, tps: number}>(); const toolTimings = new Map<string, {duration: number, tps: number}>();
tools = this.wrapToolTiming(tools, toolTimings); tools = this.wrapToolTiming(tools, toolTimings);
if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
prompts.unshift(options.system || this.ai.options.llm?.system || ''); prompts.unshift(options.system || this.ai.options.llm?.system || '');
request = this.models[m].ask(message, {...options, tools, system: prompts.filter(Boolean).join('\n\n')}); request = this.models[m].ask(message, {...options, tools, system: prompts.filter(Boolean).join('\n\n')});
let resp = await request; let resp = await request;
// Providers stamp duration/tps on assistant entries themselves (from real API usage). // Capture meta (duration / tps)
// Overwrite tool entries with actual tool-execution timing instead of the LLM call timing.
for(const h of history) { for(const h of history) {
if(h.role === 'tool' && toolTimings.has(h.id)) Object.assign(h, toolTimings.get(h.id)); if(h.role === 'tool' && toolTimings.has(h.id)) Object.assign(h, toolTimings.get(h.id));
} }
// Spice delegated agents response into history if(typeof resp === 'string' && !resp.trim() && delegateState.resp !== null) resp = delegateState.resp;
let lastDelegateResp: string | null = null;
if(pendingDelegates.size) {
for(let i = 0; i < history.length; i++) {
const h: any = history[i];
if(h.role !== 'tool' || !pendingDelegates.has(h.id)) continue;
const {resp: delegateResp, subHistory, duration, tps} = pendingDelegates.get(h.id)!;
pendingDelegates.delete(h.id);
const insert: LLMMessage[] = [...subHistory.filter(sh => sh.role === 'tool'), {role: 'assistant', content: delegateResp, timestamp: Date.now(), duration, tps}];
history.splice(i + 1, 0, ...insert);
lastDelegateResp = delegateResp;
i += insert.length;
}
}
// If the orchestrator added no commentary of its own, its answer IS the delegate's answer
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 !== 'memory_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) { if(options.compress && this.estimateTokens(history) >= options.compress.max) {
if(mem?.update) await this.memoryManager.memorize(history, mem.memory, {model: options.memoryModel || this.defaultModel, ...options}); if(mem?.update) await this.memoryManager.memorize(history, mem.memory, {model: options.memoryModel || this.defaultModel, ...options});
const compressed = await this.compressHistory(history, options.compress.max, options.compress.min, options); const compressed = await this.compressHistory(history, options.compress.max, options.compress.min, options);
if(options.history) options.history.splice(0, options.history.length, ...compressed); if(options.history) options.history.splice(0, options.history.length, ...compressed);
} }
const requestDuration = Date.now() - requestStart;
const totalTokens = history
.filter((h: any) => h.role === 'assistant' && h.duration && h.tps)
.reduce((sum: number, h: any) => sum + h.tps * (h.duration / 1000), 0);
const requestTps = requestDuration > 0 ? totalTokens / (requestDuration / 1000) : 0;
Object.assign(promise, {duration: requestDuration, tps: requestTps});
return resp; return resp;
})(); })();

View File

@@ -2,6 +2,16 @@ import {LLMRequest, LLMMessage} from './llm.ts';
import {AiTool} from './tools.ts'; import {AiTool} from './tools.ts';
import {KDPoint, KDTree} from './kd-tree.ts'; import {KDPoint, KDTree} from './kd-tree.ts';
const FACTS_HEADING = '## Facts';
const GENERIC_TEMPLATE = `# {{Title}}
## Summary
## Details
## Related`;
export class MemoryCache { export class MemoryCache {
private tree: KDTree<MemoryRef>; private tree: KDTree<MemoryRef>;
public memories: Memory[]; public memories: Memory[];
@@ -77,46 +87,35 @@ export type Memory = {
description: string; description: string;
content: string; content: string;
embedding: number[]; embedding: number[];
links: string[];
backlinks: string[];
} }
export type MemoryRef = { type MemoryRef = {
name: string; name: string;
description: string; description: string;
} }
export type FactBucket = { type FactBucket = {
subject: string; subject: string;
facts: string[]; facts: string[];
} }
export type MemoryNode = {
name: string;
missing: boolean;
links: string[];
backlinks: string[];
}
function extractLinks(content: string): string[] { function extractLinks(content: string): string[] {
if (!content) return []; if (!content) return [];
const matches = content.matchAll(/\[\[([^\]]+)\]\]/g); const matches = content.matchAll(/\[\[([^\]]+)\]\]/g);
return [...new Set([...matches].map(m => m[1].trim()))]; return [...new Set([...matches].map(m => m[1].trim()))];
} }
export function extractMetadata(content: string): {links: string[], backlinks: string[]} { export function rebuildGraph(memories: Memory[]): void {
const match = content.match(/^---\n([\s\S]*?)\n---/); for (const m of memories) m.links = extractLinks(m.content).filter(l => l !== m.name);
if (!match) return {links: [], backlinks: []}; for (const m of memories) m.backlinks = [];
for (const m of memories) {
const fm = match[1]; for (const link of m.links) {
const getList = (key: string): string[] => { const target = memories.find(t => t.name === link);
const m = fm.match(new RegExp(`^${key}:\\s*\\[(.*)\\]$`, 'm')); if (target) target.backlinks.push(m.name);
if (!m || !m[1].trim()) return []; }
return m[1].split(',').map(s => s.trim().replace(/^"|"$/g, '')).filter(Boolean); }
};
return {
links: getList('links'),
backlinks: getList('backlinks'),
};
} }
function dedupeFacts(facts: string[]): string[] { function dedupeFacts(facts: string[]): string[] {
@@ -147,23 +146,11 @@ function getWeekMonday(date: Date = new Date()): string {
return d.toISOString().slice(0, 10); return d.toISOString().slice(0, 10);
} }
function getWeekSunday(monday: string): string {
const d = new Date(`${monday}T00:00:00Z`);
d.setUTCDate(d.getUTCDate() + 6);
return d.toISOString().slice(0, 10);
}
export class MemoryManager { export class MemoryManager {
private recentlyTouched = new Map<string, number>(); private recentlyTouched = new Map<string, number>();
private pendingMemorizations = new Map<string, {
memories: Memory[] | MemoryCache,
tempMemoryName: string,
timestamp: number,
}>();
private queues = new Map<string, { private queues = new Map<string, {
pending: string[], dirty: boolean,
request: {abort?: () => void} | null, request: {abort?: () => void} | null,
task: Promise<void>, task: Promise<void>,
}>(); }>();
@@ -176,7 +163,7 @@ export class MemoryManager {
name: {type: 'string', description: 'Exact memory name', required: true}, name: {type: 'string', description: 'Exact memory name', required: true},
}, },
fn: (args: any) => { fn: (args: any) => {
const mems = memories instanceof MemoryCache ? memories.memories : memories; const mems = this.unwrap(memories);
const mem = mems.find(m => m.name === args.name); const mem = mems.find(m => m.name === args.name);
if (!mem) return 'Document not found'; if (!mem) return 'Document not found';
this.touch(mem.name); this.touch(mem.name);
@@ -205,110 +192,58 @@ export class MemoryManager {
return raw ? {memory: <Memory[] | MemoryCache>m, inject: true, tool: true, update: true} : {inject: true, tool: true, update: true, ...m}; return raw ? {memory: <Memory[] | MemoryCache>m, inject: true, tool: true, update: true} : {inject: true, tool: true, update: true, ...m};
} }
private async createTempMemory(conversation: string): Promise<Memory> { private unwrap(memories: Memory[] | MemoryCache): Memory[] {
const timestamp = Date.now(); return memories instanceof MemoryCache ? memories.memories : memories;
const content = `---
name: _temp_${timestamp}
description: Temporary memory - processing in background
tags: [_temporary]
links: []
backlinks: []
modified: ${new Date().toISOString()}
---
# Recent Conversation (Processing)
${conversation}`;
const [e] = await this.llm.embedding(content);
return {
name: `_temp_${timestamp}`,
description: 'Temporary memory - processing in background',
content,
embedding: e?.embedding || [],
};
} }
private applyHeader(content: string, header: string): string { private sync(memories: Memory[] | MemoryCache): void {
return `${header}\n\n${this.stripHeader(content)}`; if (memories instanceof MemoryCache) memories.rebuild();
} }
private async backgroundMemorization(conversation: string, memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string): Promise<void> { private parseFrontmatter(content: string): {fm: Map<string, string>, body: string} {
const mem = memories instanceof MemoryCache ? memories.memories : memories; const match = content.match(/^---\n([\s\S]*?)\n---\n?([\s\S]*)$/);
const monday = getWeekMonday(); if (!match) return {fm: new Map(), body: content};
const sunday = getWeekSunday(monday); const fm = new Map<string, string>();
const buckets = await this.factAgent(conversation, mem, options, monday); for (const line of match[1].split('\n')) {
if(!buckets.length) return; const i = line.indexOf(':');
const jobs = [...buckets].map(({subject, facts}) => { if (i === -1) continue;
let node = mem.find(m => m.name === subject); fm.set(line.slice(0, i).trim(), line.slice(i + 1).trim());
if(!node) {
node = {name: subject, description: '', content: '', embedding: [],};
mem.push(node);
} }
const week = subject.startsWith('Journal/') ? {monday, sunday} : undefined; return {fm, body: match[2]};
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 { private writeFrontmatter(fm: Map<string, string>, body: string): string {
const tags = node.name.split('/')[0]?.toLowerCase(); const lines = [...fm.entries()].map(([k, v]) => `${k}: ${v}`);
const lines = [ return `---\n${lines.join('\n')}\n---\n\n${body.trimStart()}`;
'---',
`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[] { private stripHeader(content: string): string {
const scored = memories return content.replace(/^---[\s\S]*?\n---\n?/, '').trimStart();
.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);
} }
/** private touchHeader(node: Memory, body: string): string {
* Coalescing queue: if a doc is already compiling, abort the in-flight run, merge its const {fm} = this.parseFrontmatter(node.content);
* facts with the new ones and restart. Never blocks a pending update, never drops facts. fm.set('name', node.name);
*/ fm.set('description', node.description || '');
private enqueue(node: Memory, facts: string[], memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string, week?: {monday: string, sunday: string}): Promise<void> { fm.set('modified', new Date().toISOString());
const key = node.name; return this.writeFrontmatter(fm, body);
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<void>} = {pending: [...facts], request: null, task: Promise.resolve()}; private ensureDoc(node: Memory): void {
this.queues.set(key, entry); if (node.content) return;
const m = memories instanceof MemoryCache ? memories.memories : memories; const title = node.name.split('/').pop() ?? node.name;
entry.task = (async () => { node.content = this.touchHeader(node, `# ${title}\n`);
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[] { private appendFacts(node: Memory, facts: string[]): void {
return memories.map(m => ({name: m.name, description: m.description})); this.ensureDoc(node);
const body = this.stripHeader(node.content);
const bullets = facts.map(f => `- ${f}`).join('\n');
const idx = body.indexOf(FACTS_HEADING);
const newBody = idx === -1
? `${body.trimEnd()}\n\n${FACTS_HEADING}\n${bullets}\n`
: `${body.slice(0, idx + FACTS_HEADING.length)}\n${bullets}${body.slice(idx + FACTS_HEADING.length)}`;
node.content = this.touchHeader(node, newBody);
} }
decay() { decay() {
@@ -318,71 +253,27 @@ ${conversation}`;
} }
} }
forget(name: string, memories: Memory[] | MemoryCache): boolean { touch(name: string, ttl = 2) {
const mem = memories instanceof MemoryCache ? memories.memories : memories; this.recentlyTouched.set(name, ttl);
const idx = mem.findIndex(m => m.name === name);
if (idx === -1) return false;
for (const node of mem) {
const {links, backlinks} = extractMetadata(node.content);
const newBacklinks = backlinks.filter(b => b !== name);
const newLinks = links.filter(l => l !== name);
if (newBacklinks.length !== backlinks.length || newLinks.length !== links.length) {
node.content = this.updateFrontmatter(node.content, {
links: newLinks,
backlinks: newBacklinks,
});
}
}
mem.splice(idx, 1);
if (memories instanceof MemoryCache) memories.rebuild();
return true;
} }
getTouched(): string[] { getTouched(): string[] {
return [...this.recentlyTouched.keys()]; return [...this.recentlyTouched.keys()];
} }
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise<Memory[]> { forget(name: string, memories: Memory[] | MemoryCache): boolean {
const conversation = history const mem = this.unwrap(memories);
.filter(h => h.role === 'user' || h.role === 'assistant') const idx = mem.findIndex(m => m.name === name);
.map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim(); if (idx === -1) return false;
if(!conversation) return [];
const trackingId = `${Date.now()}_${Math.random()}`; mem.splice(idx, 1);
const tempMemory = await this.createTempMemory(conversation); rebuildGraph(mem);
const mem = memories instanceof MemoryCache ? memories.memories : memories; this.sync(memories);
mem.push(tempMemory); return true;
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<Memory[]> { async recollect(query: string, memories: Memory[] | MemoryCache, limit = 5, graphDepth = 1): Promise<Memory[]> {
const mem: Memory[] = memories instanceof MemoryCache ? memories.memories : memories; const mem = this.unwrap(memories);
if (!mem.length) return []; if (!mem.length) return [];
const [e] = await this.llm.embedding(query); const [e] = await this.llm.embedding(query);
@@ -400,8 +291,7 @@ ${conversation}`;
for (const name of frontier) { for (const name of frontier) {
const node = mem.find(m => m.name === name); const node = mem.find(m => m.name === name);
if (!node) continue; if (!node) continue;
const {links} = extractMetadata(node.content); for (const link of node.links) {
for (const link of links) {
if (!found.has(link) && mem.find(m => m.name === link)) { if (!found.has(link) && mem.find(m => m.name === link)) {
found.add(link); found.add(link);
next.push(link); next.push(link);
@@ -419,114 +309,151 @@ ${conversation}`;
return ordered.map(n => mem.find(m => m.name === n)!).filter(Boolean); return ordered.map(n => mem.find(m => m.name === n)!).filter(Boolean);
} }
touch(name: string, ttl = 2) { private cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] {
this.recentlyTouched.set(name, ttl); 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);
} }
private updateFrontmatter(content: string, updates: {links?: string[], backlinks?: string[]}): string { private listNodes(memories: Memory[]): MemoryRef[] {
const match = content.match(/^---\n([\s\S]*?)\n---\n\n?([\s\S]*)$/); return memories.map(m => ({name: m.name, description: m.description}));
if (!match) return content;
const [, fm, body] = match;
let newFm = fm;
if (updates.links !== undefined) {
const linksList = updates.links.length ? `[${updates.links.map(l => `"${l}"`).join(', ')}]` : '[]';
newFm = newFm.replace(/^links:.*$/m, `links: ${linksList}`);
} }
if (updates.backlinks !== undefined) { async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise<Memory[]> {
const backlinksList = updates.backlinks.length ? `[${updates.backlinks.map(l => `"${l}"`).join(', ')}]` : '[]'; const conversation = history
newFm = newFm.replace(/^backlinks:.*$/m, `backlinks: ${backlinksList}`); .filter(h => h.role === 'user' || h.role === 'assistant')
.map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim();
if (!conversation) return [];
const uid = `${Date.now()}_${Math.random().toString(36).slice(2)}`;
// NOTE: adjust field names below (id/tool_call_id/name) to match your LLMMessage/tool-call schema.
const pending = {role: 'tool', name: 'memory_process', id: uid, content: 'Processing…'} as unknown as LLMMessage;
history.push(pending);
const mem = this.unwrap(memories);
const buckets = await this.factAgent(conversation, mem, options, getWeekMonday());
const touched: Memory[] = [];
for (const {subject, facts} of buckets) {
let node = mem.find(m => m.name === subject);
if (!node) {
node = {name: subject, description: '', content: '', embedding: [], links: [], backlinks: []};
mem.push(node);
}
this.appendFacts(node, facts);
const [e] = await this.llm.embedding(node.content);
if (e) node.embedding = e.embedding;
this.touch(node.name);
touched.push(node);
} }
newFm = newFm.replace(/^modified:.*$/m, `modified: ${new Date().toISOString()}`); if (touched.length) {
rebuildGraph(mem);
return `---\n${newFm}\n---\n\n${body}`; this.sync(memories);
(pending as any).content = `Saved to ${touched.map(n => `[[${n.name}]]`).join(', ')}`;
for (const node of touched) this.reconcile(node, memories, options).catch(() => {});
} else {
(pending as any).content = 'Nothing worth remembering.';
} }
private stripHeader(content: string): string { return touched;
return content.replace(/^---[\s\S]*?\n---\n?/, '').trimStart();
} }
private async docAgent(node: Memory, facts: string[], memories: Memory[], options: LLMRequest, tempName: string, week: {monday: string, sunday: string} | undefined, entry: {request: {abort?: () => void} | null}): Promise<boolean> { /** Manual/cron entry point. scope 'touched' only reconciles docs with a pending Facts inbox. */
const {links: oldLinks} = extractMetadata(node.content); async reconcileVault(memories: Memory[] | MemoryCache, options: LLMRequest, scope: 'touched' | 'all' = 'touched'): Promise<void> {
const mem = this.unwrap(memories);
const targets = scope === 'all' ? mem : mem.filter(m => m.content.includes(FACTS_HEADING));
await Promise.all(targets.map(node => this.reconcile(node, memories, options)));
this.sync(memories);
}
/**
* Coalescing queue: if a doc is already reconciling, mark it dirty and abort the in-flight
* request. The loop below always re-reads node.content fresh, so nothing is ever dropped.
*/
private reconcile(node: Memory, memories: Memory[] | MemoryCache, options: LLMRequest): Promise<void> {
const key = node.name;
const existing = this.queues.get(key);
if (existing) {
existing.dirty = true;
existing.request?.abort?.();
return existing.task;
}
const entry = {dirty: false, request: null, task: Promise.resolve()};
this.queues.set(key, entry);
const mem = this.unwrap(memories);
entry.task = (async () => {
do {
entry.dirty = false;
await this.reconcileDoc(node, mem, options, entry);
} while (entry.dirty);
})().finally(() => {
this.queues.delete(key);
rebuildGraph(mem);
this.sync(memories);
});
return entry.task;
}
private async reconcileDoc(node: Memory, memories: Memory[], options: LLMRequest, entry: {request: {abort?: () => void} | null}): Promise<void> {
const currentBody = this.stripHeader(node.content); const currentBody = this.stripHeader(node.content);
let update; let update;
try { try {
for(let i = 0; i < 3 && !update?.content; i++) { for (let i = 0; i < 2 && !update?.content; i++) {
const request = this.llm.ask(`New Facts:\n${facts.map(f => `- ${f}`).join('\n')}`, { const request = this.llm.ask(currentBody, {
model: options.model, model: options.model,
temperature: 0.3, temperature: 0.3,
schema: { schema: {
description: {type: 'string', description: 'One-line description of what this document covers, no formatting or emojis', required: true}, description: {type: 'string', description: 'One-line description of what this document covers, no formatting or emojis', required: true},
content: {type: 'string', description: 'Rewritten document in markdown, without the frontmatter block', required: true}, content: {type: 'string', description: 'Rewritten document body in markdown, without the frontmatter block', required: true},
}, },
system: `You are a knowledge base editor. Rewrite the current document below so it incorporates the new facts. system: `You are a knowledge base editor maintaining one document in an Obsidian-style vault.
If the document has a "${FACTS_HEADING}" section, integrate every bullet under it into the appropriate part of the document, then remove the "${FACTS_HEADING}" section entirely. If there is no such section, just tidy the document per the rules below.
Structure: follow this generic shape loosely, adapting section names/order to what the content actually needs (e.g. journal-style docs may want a timeline instead of "Details"):
\`\`\`markdown
${GENERIC_TEMPLATE}
\`\`\`
Formatting rules: Formatting rules:
- Use Obsidian-style markdown: # headings, **bold** to add emphasis, __italics__ for titles, terms, etc, bullet & numbered lists for grouped 1D data and tables for 2D data - Use Obsidian-style markdown: # headings, **bold** for emphasis, bullet & numbered lists for grouped 1D data, tables for 2D data
- Link related concepts with [[WikiLink]] notation using full paths like [[People/Sarah]] or [[Projects/Website]] - Link related concepts with [[WikiLink]] notation using full paths like [[People/Sarah]] or [[Projects/Website]]
- Create links for specific entities (person, place, project, program) and abstract concepts (quantum mechanics, entropy) but skip generics (car, red, dog) - Create links for specific entities (person, place, project, program) and abstract concepts, but skip generics (car, red, dog)
- Keep the document concise, factual, and human-readable - Keep the document concise, factual, and human-readable
- Resolve contradictions: the new facts always win — delete the outdated statement entirely, never keep both - Resolve contradictions: newer facts always win — delete the outdated statement entirely, never keep both
- Later facts in the list override earlier ones
- Do not add frontmatter blocks, filler, preamble, or AI commentary - Do not add frontmatter blocks, filler, preamble, or AI commentary
${week ? '- This is a weekly journal entry.\n' : ''}
All nodes: Other nodes in the vault (link to these instead of duplicating their content):
${this.listNodes(memories).map(n => n.name).join(', ') || 'none'} ${this.listNodes(memories).filter(n => n.name !== node.name).map(n => n.name).join(', ') || 'none'}
Current document: Current document:
\`\`\`markdown \`\`\`markdown
${currentBody} ${currentBody}
\`\`\``} \`\`\``,
); });
entry.request = request; entry.request = request;
update = await request; update = await request;
} }
} catch (err: any) { } catch (err: any) {
if (err?.name === 'AbortError') return false; if (err?.name === 'AbortError') return;
throw err; throw err;
} finally { } finally {
entry.request = null; entry.request = null;
} }
if(!update?.content) return false; if (!update?.content) return;
const newLinks = extractLinks(update.content).filter(l => l !== node.name && l !== tempName); node.description = node.name !== 'People/User' ? update.description : 'All information about the current user';
const newLinkSet = new Set(newLinks); node.content = this.touchHeader(node, update.content);
const oldLinkSet = new Set(oldLinks);
for (const added of newLinkSet) {
if (!oldLinkSet.has(added)) {
const target = memories.find(m => m.name === added);
if (target) {
const {backlinks} = extractMetadata(target.content);
if (!backlinks.includes(node.name)) {
target.content = this.updateFrontmatter(target.content, {
backlinks: [...backlinks, node.name],
});
}
}
}
}
for (const removed of oldLinkSet) {
if (!newLinkSet.has(removed)) {
const target = memories.find(m => m.name === removed);
if (target) {
const {backlinks} = extractMetadata(target.content);
target.content = this.updateFrontmatter(target.content, {
backlinks: backlinks.filter(b => b !== node.name),
});
}
}
}
const {backlinks} = extractMetadata(node.content);
node.description = node.name !== 'Person/User' ? update.description : 'All information about the current user';
node.content = this.applyHeader(update.content, this.buildHeader(node, week, newLinks, backlinks));
const [e] = await this.llm.embedding(node.content); const [e] = await this.llm.embedding(node.content);
if (e) node.embedding = e.embedding; if (e) node.embedding = e.embedding;
return true;
} }
private async factAgent(conversation: string, memories: Memory[], options: LLMRequest, weekKey: string): Promise<FactBucket[]> { private async factAgent(conversation: string, memories: Memory[], options: LLMRequest, weekKey: string): Promise<FactBucket[]> {
@@ -546,13 +473,13 @@ Rules:
When extracting facts, you MUST also decide the exact destination path: When extracting facts, you MUST also decide the exact destination path:
- Use an existing node name if the facts clearly belong there - Use an existing node name if the facts clearly belong there
- All information primary about the user should go under "People/User" - All information primarily about the user should go under "People/User"
- When required, create a new path following collection/subject format (e.g., People/Sarah, Projects/Oxide) - When required, create a new path following collection/subject format (e.g., People/Sarah, Projects/Oxide) — you are not limited to any fixed list of collections, use whatever fits
- For journal entries, use "Journal" - For journal entries, use "Journal"
Available nodes: Available nodes:
- Journal - Journal
${this.listNodes(memories).filter(n => !n.name.includes('_temp_') && !n.name.includes('Journal')).map(n => `- ${n.name}: ${n.description}`).join('\n') || 'None yet.'}`, ${this.listNodes(memories).filter(n => !n.name.includes('Journal')).map(n => `- ${n.name}: ${n.description}`).join('\n') || 'None yet.'}`,
tools: [{ tools: [{
name: 'facts_extract', name: 'facts_extract',
description: 'Submit facts with their destination', description: 'Submit facts with their destination',

View File

@@ -1,19 +1,28 @@
import {OpenAI as openAI} from 'openai'; import {OpenAI as openAI} from 'openai';
import {findByProp, objectMap, JSONSanitize, JSONAttemptParse, clean} from '@ztimson/utils'; import {findByProp, objectMap, JSONSanitize, JSONAttemptParse, clean, makeArray} from '@ztimson/utils';
import {AbortablePromise, Ai} from './ai.ts'; import {AbortablePromise, Ai} from './ai.ts';
import {LLMMessage, LLMRequest} from './llm.ts'; import {LLMMessage, LLMRequest} from './llm.ts';
import {LLMProvider} from './provider.ts'; import {LLMProvider} from './provider.ts';
import {TokenPool} from './token-pool.ts';
import {convertSchema} from './tools.ts'; import {convertSchema} from './tools.ts';
export class OpenAi extends LLMProvider { export class OpenAi extends LLMProvider {
client!: openAI; tokenPool!: TokenPool;
private clients = new Map<string, openAI>();
constructor(public readonly ai: Ai, public readonly host: string | null, public readonly token: string, public model: string) { constructor(public readonly ai: Ai, public readonly host: string | null, public readonly token: string | string[], public model: string) {
super(); super();
this.client = new openAI(clean({ const tokens = makeArray(token).filter(Boolean);
baseURL: host, this.tokenPool = new TokenPool(...(tokens.length ? tokens : [host ? 'ignored' : '']));
apiKey: token || (host ? 'ignored' : undefined) }
}));
private getClient(token: string): openAI {
let client = this.clients.get(token);
if(!client) {
client = new openAI(clean({baseURL: this.host, apiKey: token || undefined}));
this.clients.set(token, client);
}
return client;
} }
private toStandard(history: any[]): LLMMessage[] { private toStandard(history: any[]): LLMMessage[] {
@@ -117,7 +126,7 @@ export class OpenAi extends LLMProvider {
do { do {
requestParams.messages = history.map(({timestamp, ...m}) => m); requestParams.messages = history.map(({timestamp, ...m}) => m);
const callStart = Date.now(); const callStart = Date.now();
resp = await this.client.chat.completions.create(requestParams).catch(err => { resp = await this.tokenPool.run(token => this.getClient(token).chat.completions.create(requestParams)).catch(err => {
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`; err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
throw err; throw err;
}); });

65
src/token-pool.ts Normal file
View File

@@ -0,0 +1,65 @@
const DEFAULT_COOLDOWN = 15 * 60 * 1000;
type TokenState = {
token: string;
cooldownUntil: number; // 0 = available now
lastError?: {code: number, message: string};
};
export class TokenPoolExhaustedError extends Error {
constructor(public tokens: Record<string, {code: number, message: string}>) {
super(`All tokens exhausted:\n${Object.entries(tokens).map(([t, e]) => `${t}: [${e.code}] ${e.message}`).join('\n')}`);
this.name = 'TokenPoolExhaustedError';
}
}
export class TokenPool {
private states: TokenState[];
constructor(...tokens: string[]) {
this.states = tokens.map(token => ({token, cooldownUntil: 0}));
}
private preview(token: string): string {
return token.length <= 8 ? '****' : `${token.slice(0, 4)}...${token.slice(-4)}`;
}
/** Anthropic & OpenAI SDKs both attach `status` to thrown errors */
private statusCode(err: any): number {
return err?.status ?? err?.response?.status ?? err?.statusCode;
}
private retryAfter(err: any): number {
const headers = err?.headers || err?.response?.headers;
const raw = headers?.get?.('retry-after') ?? headers?.['retry-after'];
if(raw) {
const seconds = Number(raw);
if(!isNaN(seconds)) return Date.now() + seconds * 1000;
const date = new Date(raw).getTime();
if(!isNaN(date)) return date;
}
return Date.now() + DEFAULT_COOLDOWN;
}
async run<T>(fn: (token: string) => Promise<T>): Promise<T> {
const now = Date.now();
for(const state of this.states) {
if(state.cooldownUntil > now) continue;
try {
const result = await fn(state.token);
state.cooldownUntil = 0;
state.lastError = undefined;
return result;
} catch(err: any) {
const code = this.statusCode(err);
if(![401, 403, 429].includes(code)) throw err;
state.cooldownUntil = code === 429 ? this.retryAfter(err) : Date.now() + DEFAULT_COOLDOWN;
state.lastError = {code, message: err?.message || 'Unknown error'};
}
}
const failures: Record<string, {code: number, message: string}> = {};
this.states.forEach(s => { if(s.lastError) failures[this.preview(s.token)] = s.lastError; });
throw new TokenPoolExhaustedError(failures);
}
}

167
tests/llm.spec.ts Normal file
View File

@@ -0,0 +1,167 @@
import {describe, it, expect, vi, beforeEach} from 'vitest';
import LLM from '../src/llm';
const {FakeProvider, providerLog} = vi.hoisted(() => {
const providerLog: any[] = [];
class FakeProvider {
model: string;
constructor(...args: any[]) { this.model = args[args.length - 1]; }
ask(message: string, opts: any) {
let aborted = false;
const p = (async () => {
const script = (globalThis as any).__scripts?.[this.model];
const plan = script ? script(message, opts) : {text: ''};
providerLog.push({model: this.model, message, system: opts.system, tools: (opts.tools || []).map((t: any) => t.name)});
for (const c of plan.calls || []) {
if (aborted) break;
const tool = (opts.tools || []).find((t: any) => t.name === c.tool);
const id = c.id || `${c.tool}_${Math.random()}`;
const content = await tool.fn(c.args, opts.stream, null, id);
opts.history.push({role: 'tool', id, name: c.tool, args: c.args, content, timestamp: Date.now()});
}
const text = plan.text ?? '';
if (opts.stream && text) opts.stream({text, done: true});
opts.history.push({role: 'assistant', content: text, timestamp: Date.now(), duration: 10, tps: 5});
return text;
})();
return Object.assign(p, {abort: () => { aborted = true; }});
}
}
return {FakeProvider, providerLog};
});
vi.mock('../src/antrhopic.ts', () => ({Anthropic: FakeProvider}));
vi.mock('../src/open-ai.ts', () => ({OpenAi: FakeProvider}));
function makeAi(models: any) {
return {options: {llm: {models}}} as any;
}
beforeEach(() => {
providerLog.length = 0;
(globalThis as any).__scripts = {};
});
describe('LLM cross-provider interchangeability', () => {
it('runs identical tool calls the same way on an anthropic-backed model and an openai-backed model', async () => {
const ai = makeAi({
claude: {proto: 'anthropic', token: 'x'},
gpt: {proto: 'openai', token: 'y', host: 'http://local'},
});
const llm = new LLM(ai);
const calc = {
name: 'calc_add',
description: 'Add two numbers',
args: {a: {type: 'number', required: true}, b: {type: 'number', required: true}},
fn: (args: any) => String(args.a + args.b),
};
(globalThis as any).__scripts.claude = () => ({calls: [{tool: 'calc_add', args: {a: 2, b: 3}}], text: 'Result: 5'});
(globalThis as any).__scripts.gpt = () => ({calls: [{tool: 'calc_add', args: {a: 2, b: 3}}], text: 'Result: 5'});
const historyA: any[] = [], historyB: any[] = [];
const respA = await llm.ask('add 2 and 3', {model: 'claude', tools: [calc], history: historyA});
const respB = await llm.ask('add 2 and 3', {model: 'gpt', tools: [calc], history: historyB});
expect(respA).toBe('Result: 5');
expect(respB).toBe('Result: 5');
expect(providerLog.find(l => l.model === 'claude')!.tools).toContain('calc_add');
expect(providerLog.find(l => l.model === 'gpt')!.tools).toContain('calc_add');
// tool timing gets recomputed from real execution regardless of proto
for (const h of [historyA.find(h => h.name === 'calc_add'), historyB.find(h => h.name === 'calc_add')]) {
expect(h.content).toBe('5');
expect(typeof h.duration).toBe('number');
expect(typeof h.tps).toBe('number');
}
});
it('lets the same shared history flow across model + proto swaps with different system prompts', async () => {
const ai = makeAi({
claude: {proto: 'anthropic', token: 'x'},
gpt: {proto: 'openai', token: 'y', host: 'http://local'},
});
const llm = new LLM(ai);
const history: any[] = [];
(globalThis as any).__scripts.claude = () => ({text: 'Hi from claude'});
(globalThis as any).__scripts.gpt = () => ({text: 'Hi from gpt'});
const r1 = await llm.ask('hello', {model: 'claude', system: 'You are terse.', history});
const r2 = await llm.ask('follow up', {model: 'gpt', system: 'You are verbose.', history});
expect(r1).toBe('Hi from claude');
expect(r2).toBe('Hi from gpt');
expect(history.filter(h => h.role === 'assistant').map(h => h.content)).toEqual(['Hi from claude', 'Hi from gpt']);
expect(providerLog[0].system).toContain('You are terse.');
expect(providerLog[1].system).toContain('You are verbose.');
});
it('exposes MCP tools the same way no matter which proto backs the model', async () => {
const ai = makeAi({claude: {proto: 'anthropic', token: 'x'}, gpt: {proto: 'openai', token: 'y', host: 'http://local'}});
const llm = new LLM(ai);
const mcp = [{name: 'weather', host: 'http://mcp.local'}];
global.fetch = vi.fn(async (url: string, opts?: any) => {
if (url.endsWith('/tools')) {
return {json: async () => ({tools: [{name: 'lookup', description: 'Look up weather', inputSchema: {properties: {city: {type: 'string'}}, required: ['city']}}]})} as any;
}
const body = JSON.parse(opts.body);
return {json: async () => ({content: [{text: `Sunny in ${body.arguments.city}`}]})} as any;
}) as any;
for (const model of ['claude', 'gpt']) {
(globalThis as any).__scripts[model] = () => ({calls: [{tool: 'weather_lookup', args: {city: 'Rome'}}], text: 'done'});
const history: any[] = [];
await llm.ask('weather?', {model, mcp, history});
expect(history.find(h => h.name === 'weather_lookup')?.content).toBe('Sunny in Rome');
}
});
it('exposes and resolves skill documents identically across protos', async () => {
const ai = makeAi({claude: {proto: 'anthropic', token: 'x'}, gpt: {proto: 'openai', token: 'y', host: 'http://local'}});
const llm = new LLM(ai);
const skills = [{name: 'Onboarding', description: 'How to onboard a user', content: 'Step 1...'}];
for (const model of ['claude', 'gpt']) {
(globalThis as any).__scripts[model] = () => ({calls: [{tool: 'skill_read', args: {name: 'Onboarding'}}], text: 'done'});
const history: any[] = [];
await llm.ask('onboard me', {model, skills, history});
expect(history.find(h => h.name === 'skill_read')?.content).toContain('Step 1...');
}
});
it('delegate agent mutates the shared history directly and backfills the orchestrator response, across protos', async () => {
const ai = makeAi({claude: {proto: 'anthropic', token: 'x'}, gpt: {proto: 'openai', token: 'y', host: 'http://local'}});
const llm = new LLM(ai);
const history: any[] = [{role: 'user', content: 'research quantum computing'}];
const researcher = {name: 'researcher', system: 'You research topics.', delegate: true, model: 'gpt'};
(globalThis as any).__scripts.claude = () => ({calls: [{tool: 'agent_researcher', args: {}}], text: ''});
(globalThis as any).__scripts.gpt = () => ({text: 'Quantum computers use qubits.'});
const resp = await llm.ask('go', {model: 'claude', agents: [researcher], history});
expect(resp).toBe('Quantum computers use qubits.');
expect(history.some(h => h.role === 'assistant' && h.content === 'Quantum computers use qubits.')).toBe(true);
expect(history.find(h => h.name === 'agent_researcher')?.content).toBe('');
});
it('regular (non-delegate) subagent keeps its own isolated history separate from the parent, across protos', async () => {
const ai = makeAi({claude: {proto: 'anthropic', token: 'x'}, gpt: {proto: 'openai', token: 'y', host: 'http://local'}});
const llm = new LLM(ai);
const history: any[] = [];
const summarizer = {name: 'summarizer', system: 'You summarize text.', model: 'gpt'};
(globalThis as any).__scripts.claude = () => ({calls: [{tool: 'subagent_summarizer', args: {context: 'a long article', instructions: 'summarize it'}}], text: 'Summary: short version'});
(globalThis as any).__scripts.gpt = () => ({text: 'short version'});
const resp = await llm.ask('summarize this', {model: 'claude', agents: [summarizer], history});
expect(resp).toBe('Summary: short version');
expect(history.find(h => h.name === 'subagent_summarizer')?.content).toBe('short version');
// isolated history - subagent's own assistant turn never leaks into the parent
expect(history.some(h => h.role === 'assistant' && h.content === 'short version')).toBe(false);
});
});

256
tests/memory.spec.ts Normal file
View File

@@ -0,0 +1,256 @@
import {describe, it, expect, vi, beforeEach} from 'vitest';
import {MemoryManager, MemoryCache, rebuildGraph, Memory} from '../src/memory';
function makeMemory(overrides: Partial<Memory> = {}): Memory {
return {
name: 'Test/Doc',
description: '',
content: '',
embedding: [],
links: [],
backlinks: [],
...overrides,
};
}
function makeLLM() {
return {
embedding: vi.fn(async (_text: string) => [{embedding: [1, 0, 0]}]),
ask: vi.fn(async () => undefined),
};
}
describe('rebuildGraph', () => {
it('extracts [[WikiLinks]] from content, excluding self-links', () => {
const a = makeMemory({name: 'A', content: '[[B]] and [[A]] and [[C]]'});
const b = makeMemory({name: 'B', content: 'no links here'});
const mem = [a, b];
rebuildGraph(mem);
expect(a.links).toEqual(['B', 'C']);
expect(b.links).toEqual([]);
});
it('computes backlinks only for links that resolve to a real node', () => {
const a = makeMemory({name: 'A', content: '[[B]] [[Missing]]'});
const b = makeMemory({name: 'B', content: ''});
const mem = [a, b];
rebuildGraph(mem);
expect(b.backlinks).toEqual(['A']);
expect(mem.find(m => m.name === 'Missing')).toBeUndefined();
});
it('resets stale backlinks on every rebuild (no leftover from a removed link)', () => {
const a = makeMemory({name: 'A', content: '[[B]]'});
const b = makeMemory({name: 'B', content: ''});
const mem = [a, b];
rebuildGraph(mem);
expect(b.backlinks).toEqual(['A']);
a.content = 'no more links';
rebuildGraph(mem);
expect(b.backlinks).toEqual([]);
});
});
describe('MemoryCache', () => {
it('finds nearest neighbor by embedding via KD-tree search', () => {
const close = makeMemory({name: 'Close', embedding: [1, 0, 0]});
const far = makeMemory({name: 'Far', embedding: [0, 0, 1]});
const cache = new MemoryCache([close, far]);
const results = cache.search([1, 0, 0], 1);
expect(results[0].name).toBe('Close');
});
it('rebuilds the tree on add/update/remove', () => {
const cache = new MemoryCache([makeMemory({name: 'A', embedding: [1, 0, 0]})]);
cache.add(makeMemory({name: 'B', embedding: [0, 1, 0]}));
expect(cache.search([0, 1, 0], 1)[0].name).toBe('B');
cache.remove('B');
expect(cache.search([0, 1, 0], 1)[0]?.name).not.toBe('B');
});
});
describe('MemoryManager.forget', () => {
it('removes the node and recomputes backlinks for the rest of the graph', () => {
const llm = makeLLM();
const mgr = new MemoryManager(llm);
const a = makeMemory({name: 'A', content: '[[B]]'});
const b = makeMemory({name: 'B', content: '[[C]]'});
const c = makeMemory({name: 'C', content: ''});
const mem = [a, b, c];
rebuildGraph(mem);
expect(c.backlinks).toEqual(['B']);
const ok = mgr.forget('B', mem);
expect(ok).toBe(true);
expect(mem.find(m => m.name === 'B')).toBeUndefined();
expect(a.links).toEqual(['B']);
expect(c.backlinks).toEqual([]);
});
it('returns false for an unknown name', () => {
const mgr = new MemoryManager(makeLLM());
expect(mgr.forget('Nope', [makeMemory({name: 'A'})])).toBe(false);
});
});
describe('MemoryManager.recollect', () => {
it('orders vector matches first, then expands one hop via links', async () => {
const llm = makeLLM();
llm.embedding.mockResolvedValue([{embedding: [1, 0, 0]}]);
const mgr = new MemoryManager(llm);
const near = makeMemory({name: 'Near', embedding: [1, 0, 0], content: '[[Linked]]'});
const linked = makeMemory({name: 'Linked', embedding: [0, 0, 1], content: ''});
const far = makeMemory({name: 'Far', embedding: [0, 1, 0], content: ''});
const mem = [near, linked, far];
rebuildGraph(mem);
const result = await mgr.recollect('query', mem, 1, 1);
expect(result.map(r => r.name)).toEqual(['Near', 'Linked']);
});
it('returns [] when there are no memories', async () => {
const mgr = new MemoryManager(makeLLM());
expect(await mgr.recollect('q', [])).toEqual([]);
});
});
describe('MemoryManager.memorize (fast path)', () => {
let llm: ReturnType<typeof makeLLM>;
let mgr: MemoryManager;
beforeEach(() => {
llm = makeLLM();
mgr = new MemoryManager(llm);
});
it('pushes a pending tool message, then resolves it to links once facts land', async () => {
llm.ask.mockImplementation(async (_prompt: string, opts: any) => {
if (opts.tools) {
opts.tools[0].fn({destination: 'Projects/Oxide', facts: 'Uses a hybrid memory system'});
return undefined;
}
return {description: 'd', content: '# doc'};
});
const history: any[] = [{role: 'user', content: 'we use a hybrid memory system'}];
const touched = await mgr.memorize(history, [], {model: 'test'} as any);
const pending = history.find(h => h.name === 'memory_process');
expect(pending).toBeDefined();
expect(pending.content).toContain('[[Projects/Oxide]]');
expect(touched.map(t => t.name)).toEqual(['Projects/Oxide']);
});
it('creates a new node and appends facts under "## Facts" without calling the doc LLM', async () => {
llm.ask.mockImplementation(async (_prompt: string, opts: any) => {
if (opts.tools) opts.tools[0].fn({destination: 'People/Sarah', facts: 'Works at Acme, Likes hiking'});
return undefined;
});
const mem: Memory[] = [];
await mgr.memorize([{role: 'user', content: 'Sarah works at Acme and likes hiking'}] as any, mem, {model: 'test'} as any);
const node = mem.find(m => m.name === 'People/Sarah')!;
expect(node).toBeDefined();
expect(node.content).toContain('## Facts');
expect(node.content).toContain('- Works at Acme');
expect(node.content).toContain('- Likes hiking');
// doc reconciler LLM (schema call) should NOT have been awaited synchronously in this fast path assertion
});
it('routes "journal" destination to Journal/{weekMonday}', async () => {
llm.ask.mockImplementation(async (_prompt: string, opts: any) => {
if (opts.tools) opts.tools[0].fn({destination: 'journal', facts: 'Shipped v1'});
return undefined;
});
const mem: Memory[] = [];
const touched = await mgr.memorize([{role: 'user', content: 'shipped v1 today'}] as any, mem, {model: 'test'} as any);
expect(touched[0].name).toMatch(/^Journal\/\d{4}-\d{2}-\d{2}$/);
});
it('reports nothing to remember when no facts are extracted', async () => {
llm.ask.mockResolvedValue(undefined); // tools present but fn never called
const history: any[] = [{role: 'user', content: 'hey'}];
const touched = await mgr.memorize(history, [], {model: 'test'} as any);
expect(touched).toEqual([]);
expect(history.find(h => h.name === 'memory_process').content).toBe('Nothing worth remembering.');
});
it('returns [] and does nothing for an empty conversation', async () => {
const touched = await mgr.memorize([], [], {model: 'test'} as any);
expect(touched).toEqual([]);
expect(llm.ask).not.toHaveBeenCalled();
});
});
describe('MemoryManager reconcileVault', () => {
it('integrates the "## Facts" section via the doc LLM and removes it', async () => {
const llm = makeLLM();
llm.ask.mockResolvedValue({description: 'Tidy summary', content: '# Doc\n\nIntegrated fact.'});
const mgr = new MemoryManager(llm);
const node = makeMemory({
name: 'Projects/Oxide',
content: '---\nname: Projects/Oxide\n---\n\n# Doc\n\n## Facts\n- some raw fact\n',
});
const mem = [node];
await mgr.reconcileVault(mem, {model: 'test'} as any, 'all');
expect(node.content).not.toContain('## Facts');
expect(node.content).toContain('Integrated fact.');
expect(node.description).toBe('Tidy summary');
});
it('only targets docs with a pending Facts inbox when scope is "touched"', async () => {
const llm = makeLLM();
llm.ask.mockResolvedValue({description: 'd', content: '# clean'});
const mgr = new MemoryManager(llm);
const dirty = makeMemory({name: 'A', content: '## Facts\n- x'});
const clean = makeMemory({name: 'B', content: '# already tidy'});
await mgr.reconcileVault([dirty, clean], {model: 'test'} as any, 'touched');
expect(dirty.content).toContain('# clean'); // rewritten (frontmatter now wraps it)
expect(clean.content).toBe('# already tidy'); // untouched, never queued
});
});
describe('MemoryManager reconcile coalescing', () => {
it('coalesces a second call while one is in-flight: marks dirty, aborts, reuses the same task promise', () => {
const llm = makeLLM();
const abort = vi.fn();
let calls = 0;
llm.ask.mockImplementation(() => {
calls++;
const pending: any = new Promise(() => {}); // never resolves in this test
pending.abort = abort;
return pending;
});
const mgr: any = new MemoryManager(llm);
const node = makeMemory({name: 'Q', content: '# Q\n\n## Facts\n- f'});
const mem = [node];
const p1 = mgr.reconcile(node, mem, {model: 'test'});
const p2 = mgr.reconcile(node, mem, {model: 'test'});
expect(p2).toBe(p1); // same in-flight task, not a new queue entry
expect(abort).toHaveBeenCalledTimes(1); // second call aborted the in-flight request
expect(calls).toBe(1); // no second ask() fired synchronously — it'll rerun via the dirty loop
});
});