Compare commits

..

4 Commits

Author SHA1 Message Date
1aa6cdf329 Agent/subagent support
All checks were successful
Publish Library / Build NPM Project (push) Successful in 45s
Publish Library / Tag Version (push) Successful in 15s
2026-08-01 18:28:16 -04:00
d022a5ef4d Improved levenshtein fuzzy match
All checks were successful
Publish Library / Build NPM Project (push) Successful in 52s
Publish Library / Tag Version (push) Successful in 17s
2026-08-01 12:00:26 -04:00
a1d438a20a Tools can now emit "done" event and end chat early gracefully
All checks were successful
Publish Library / Build NPM Project (push) Successful in 1m0s
Publish Library / Tag Version (push) Successful in 9s
2026-07-31 17:49:06 -04:00
52a9e3aaa4 Fixed history poisoning on empty tool response
All checks were successful
Publish Library / Build NPM Project (push) Successful in 51s
Publish Library / Tag Version (push) Successful in 13s
2026-07-30 22:12:49 -04:00
8 changed files with 354 additions and 198 deletions

View File

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

View File

@@ -21,10 +21,10 @@ export class Anthropic extends LLMProvider {
messages.push(<any>{timestamp, ...h}); messages.push(<any>{timestamp, ...h});
} else { } else {
const textContent = h.content?.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n'); const textContent = h.content?.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
if(textContent) messages.push({timestamp, role: h.role, content: textContent}); if(textContent) messages.push({role: h.role, content: textContent, timestamp: timestamp});
h.content.forEach((c: any) => { h.content.forEach((c: any) => {
if(c.type == 'tool_use') { if(c.type == 'tool_use') {
messages.push({timestamp, role: 'tool', id: c.id, name: c.name, args: c.input, content: undefined}); messages.push({role: 'tool', id: c.id, name: c.name, args: c.input, timestamp: c.timestamp, content: undefined});
} else if(c.type == 'tool_result') { } else if(c.type == 'tool_result') {
const m: any = messages.findLast(m => (<any>m).id == c.tool_use_id); const m: any = messages.findLast(m => (<any>m).id == c.tool_use_id);
if(m) m[c.is_error ? 'error' : 'content'] = c.content; if(m) m[c.is_error ? 'error' : 'content'] = c.content;
@@ -46,7 +46,7 @@ export class Anthropic extends LLMProvider {
i++; i++;
} }
} }
return history.map(({timestamp, ...h}) => h); return history;
} }
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> { ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
@@ -83,8 +83,9 @@ export class Anthropic extends LLMProvider {
}; };
} }
let resp: any, isFirstMessage = true; let resp: any, isFirstMessage = true, terminal = false;
do { do {
requestParams.messages = history.map(({timestamp, ...m}) => m);
resp = await this.client.messages.create(requestParams).catch(err => { resp = await this.client.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;
@@ -113,7 +114,7 @@ export class Anthropic extends LLMProvider {
} }
} else if(chunk.type === 'content_block_stop') { } else if(chunk.type === 'content_block_stop') {
const last = resp.content.at(-1); const last = resp.content.at(-1);
if(last.input != null) last.input = last.input ? JSONAttemptParse(last.input, {}) : {}; if(last?.input != null) last.input = last.input ? JSONAttemptParse(last.input, {}) : {};
} else if(chunk.type === 'message_stop') { } else if(chunk.type === 'message_stop') {
break; break;
} }
@@ -123,31 +124,35 @@ export class Anthropic extends LLMProvider {
// Run tools // Run tools
const toolCalls = resp.content.filter((c: any) => c.type === 'tool_use'); const toolCalls = resp.content.filter((c: any) => c.type === 'tool_use');
if(toolCalls.length && !controller.signal.aborted) { if(toolCalls.length && !controller.signal.aborted) {
history.push({role: 'assistant', content: resp.content}); history.push({role: 'assistant', content: resp.content, timestamp: Date.now()});
const results = await Promise.all(toolCalls.map(async (toolCall: any) => { const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
const tool = tools.find(findByProp('name', toolCall.name)); const tool = tools.find(findByProp('name', toolCall.name));
if(options.stream) options.stream({tool: toolCall.name}); if(options.stream) options.stream({tool: toolCall.name});
if(!tool) return {tool_use_id: toolCall.id, is_error: true, content: 'Tool not found'}; if(!tool) return {tool_use_id: toolCall.id, is_error: true, content: 'Tool not found'};
try { try {
const result = await tool.fn(toolCall.input, options?.stream, this.ai); // Wrap stream so a tool's `done` ends turn gracefully
const toolStream = options.stream && ((chunk: any) => {
if(chunk.done) { terminal = true; return; }
options.stream!(chunk);
});
const result = await tool.fn(toolCall.input, toolStream, this.ai);
return {type: 'tool_result', tool_use_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result}; return {type: 'tool_result', tool_use_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result};
} catch (err: any) { } catch (err: any) {
return {type: 'tool_result', tool_use_id: toolCall.id, is_error: true, content: err?.message || err?.toString() || 'Unknown'}; return {type: 'tool_result', tool_use_id: toolCall.id, is_error: true, content: err?.message || err?.toString() || 'Unknown'};
} }
})); }));
history.push({role: 'user', content: results}); history.push({role: 'user', content: results, timestamp: Date.now()});
requestParams.messages = history; requestParams.messages = history;
} }
} while (!controller.signal.aborted && resp.content.some((c: any) => c.type === 'tool_use')); } while (!terminal && !controller.signal.aborted && resp.content.some((c: any) => c.type === 'tool_use'));
if(!terminal) {
const textContent = resp.content.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n'); const textContent = resp.content.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
history.push({role: 'assistant', content: textContent}); history.push({role: 'assistant', content: textContent, timestamp: Date.now()});
}
history = this.toStandard(history); history = this.toStandard(history);
if(options.stream) options.stream({done: true}); if(options.stream) options.stream({done: true});
if(options.history) options.history.splice(0, options.history.length, ...history); if(options.history) options.history.splice(0, options.history.length, ...history);
// Return parsed JSON if schema provided
const finalContent = history.at(-1)?.content; const finalContent = history.at(-1)?.content;
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent); res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
}), {abort: () => controller.abort()}); }), {abort: () => controller.abort()});

View File

@@ -3,8 +3,6 @@ export * from './antrhopic';
export * from './audio'; export * from './audio';
export * from './llm'; export * from './llm';
export * from './memory'; export * from './memory';
export * from './memory-cache';
export * from './memory-graph';
export * from './open-ai'; export * from './open-ai';
export * from './provider'; export * from './provider';
export * from './tools'; export * from './tools';

View File

@@ -1,17 +1,31 @@
import {snakeCase} from '@ztimson/utils';
import {AbortablePromise, Ai} from './ai.ts'; import {AbortablePromise, Ai} from './ai.ts';
import {Anthropic} from './antrhopic.ts'; import {Anthropic} from './antrhopic.ts';
import {MemoryCache} from './memory-cache.ts';
import {OpenAi} from './open-ai.ts'; import {OpenAi} from './open-ai.ts';
import {LLMProvider} from './provider.ts'; import {LLMProvider} from './provider.ts';
import {AiTool, AiToolArg} from './tools.ts'; import {AiTool, AiToolArg} from './tools.ts';
import {fileURLToPath} from 'url'; import {fileURLToPath} from 'url';
import {dirname, join} from 'path'; import {dirname, join} from 'path';
import {spawn} from 'node:child_process'; import {spawn} from 'node:child_process';
import {Memory, MemoryManager} from './memory.ts'; import {Memory, MemoryCache, MemoryManager, MemoryOptions} from './memory.ts';
export type AnthropicConfig = {proto: 'anthropic', token: string}; export type AnthropicConfig = {proto: 'anthropic', token: string};
export type OpenAiConfig = {proto: 'openai', host?: string, token: string}; export type OpenAiConfig = {proto: 'openai', host?: string, token: string};
export type Agent = {
name: string;
description?: string;
model?: string | null;
temperature?: number;
system: string;
delegate?: boolean;
skills?: Skill[] | null;
tools?: AiTool[] | null;
mcp?: McpServer[] | null;
/** Explicit whitelist of agents this agent may delegate to. Default: none - must opt-in, self is always excluded */
agents?: string[] | null;
}
export type LLMMessage = { export type LLMMessage = {
/** Message originator */ /** Message originator */
role: 'assistant' | 'system' | 'user'; role: 'assistant' | 'system' | 'user';
@@ -56,13 +70,17 @@ export type LLMRequest = {
/** Compress old messages in the chat to free up context */ /** Compress old messages in the chat to free up context */
compress?: {max: number; min: number}; compress?: {max: number; min: number};
/** User's memory documents - RAG injected automatically each turn */ /** User's memory documents - RAG injected automatically each turn */
memory?: Memory[] | MemoryCache; memory?: Memory[] | MemoryCache | MemoryOptions;
/** Model to use for memory operations */ /** Model to use for memory operations */
memoryModel?: string; memoryModel?: string;
/** Skill documents the AI can browse and read on demand */ /** Skill documents the AI can browse and read on demand */
skills?: Skill[]; skills?: Skill[];
/** MCP servers to connect and expose as tools */ /** MCP servers to connect and expose as tools */
mcp?: McpServer[]; mcp?: McpServer[];
/** Subagents exposed as delegatable/wrapped tools */
agents?: Agent[];
/** @internal recursion guard for nested agent delegation */
_agentDepth?: number;
} }
export type McpServer = { export type McpServer = {
@@ -83,6 +101,7 @@ export type Skill = {
content: string; content: string;
} }
const MAX_AGENT_DEPTH = 5;
class LLM { class LLM {
private memoryManager!: MemoryManager; private memoryManager!: MemoryManager;
@@ -100,6 +119,60 @@ class LLM {
this.memoryManager = new MemoryManager(this); this.memoryManager = new MemoryManager(this);
} }
/**
* Wrap agents as tools. Nested delegation is opt-in only (empty by default, like
* tools/skills/mcp) and an agent can never call itself even if explicitly whitelisted.
* Delegate results are queued in `pending` and spliced into history by `ask()` after
* the provider's own end-of-turn history sync has already run.
*/
private setupAgent(agents: Agent[] = [], allAgents: Agent[], pending: Map<string, {resp: string, subHistory: LLMMessage[]}[]>, aborts: (() => void)[], depth = 0): AiTool[] {
return agents.map(a => {
const toolName = `${a.delegate ? '' : 'sub'}agent_${snakeCase(a.name)}`;
return {
name: toolName,
description: `${a.delegate ? 'Delegate to ' : ''}Subagent: ${a.description || a.name}`,
args: {
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},
},
fn: async (args: any, stream: any) => {
if(depth >= MAX_AGENT_DEPTH) return 'Max agent delegation depth exceeded';
const subHistory: LLMMessage[] = [];
// Opt-in only, self always excluded regardless of whitelist
const nested = (a.agents || [])
.map(name => allAgents.find(x => x.name === name))
.filter((x): x is Agent => !!x && x.name !== a.name);
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 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.
${a.system}`,
model: a.model || undefined,
temperature: a.temperature,
stream: a.delegate ? stream : undefined,
history: subHistory,
mcp: a.mcp || undefined,
skills: a.skills || undefined,
tools: a.tools || undefined,
agents: nested,
_agentDepth: depth + 1,
} as any);
aborts.push(request.abort);
const resp = await request;
if(a.delegate) {
if(!pending.has(toolName)) pending.set(toolName, []);
pending.get(toolName)!.push({resp, subHistory});
return '';
}
return resp;
}
};
});
}
private async setupMcp(servers: McpServer[] = []): Promise<{prompt: string, tools: AiTool[]}> { private async setupMcp(servers: McpServer[] = []): Promise<{prompt: string, tools: AiTool[]}> {
if(!servers?.length) return {prompt: '', tools: []}; if(!servers?.length) return {prompt: '', tools: []};
const allTools: AiTool[] = []; const allTools: AiTool[] = [];
@@ -170,9 +243,11 @@ class LLM {
if(!this.models[m]) throw new Error(`Model does not exist: ${m}`); if(!this.models[m]) throw new Error(`Model does not exist: ${m}`);
let request: AbortablePromise<string> | null = null; let request: AbortablePromise<string> | null = null;
let aborted = false; let aborted = false;
const nestedAborts: (() => void)[] = [];
const abort = () => { const abort = () => {
aborted = true; aborted = true;
request?.abort?.(); request?.abort?.();
nestedAborts.forEach(a => a());
}; };
const promise = (async () => { const promise = (async () => {
@@ -196,22 +271,46 @@ class LLM {
tools.push(...s.tools); tools.push(...s.tools);
} }
// Agents
const agents = options.agents || this.ai.options?.llm?.agents;
const pendingDelegates = new Map<string, {resp: string, subHistory: LLMMessage[]}[]>();
if(agents?.length) tools.push(...this.setupAgent(agents, agents, pendingDelegates, nestedAborts, options._agentDepth || 0));
// Memory // Memory
if (options.memory) { const mem = MemoryManager.normalize(options.memory);
const mems = options.memory instanceof MemoryCache ? options.memory.memories : options.memory; if(mem) {
const mems = mem.memory instanceof MemoryCache ? mem.memory.memories : mem.memory;
if(mems.length) { if(mems.length) {
const relevant = await this.memoryManager.recollect(message, options.memory, 5); if(mem.inject) {
const pool = 15; // candidates considered, cheap since only refs are listed
const budget = mem.maxTokens ?? 2000; // actual content injected
const relevant = await this.memoryManager.recollect(message, mem.memory, pool);
let used = 0;
const preloaded: typeof relevant = [];
const listed: typeof relevant = [];
for(const r of relevant) {
const t = this.estimateTokens(r.content);
if(used + t <= budget || preloaded.length === 0) {
preloaded.push(r);
used += t;
} else listed.push(r);
}
prompts.unshift(`You have access to the following memory files: prompts.unshift(`You have access to the following memory files:
${mems.map(m => `- ${m.name}: ${m.description}`).join('\n')} ${mems.map(m => `- ${m.name}: ${m.description}`).join('\n')}
${relevant.length ? ` ${preloaded.length ? `
Relevant memories have been preloaded: Relevant memories have been preloaded:
${relevant.map(r => ` ${preloaded.map(r => `
**${r.name}** **${r.name}**
${r.description} ${r.description}
${r.content} ${r.content}
`).join('\n---\n')} `).join('\n---\n')}
` : ''}`.trim()); ` : ''}${listed.length ? `
tools.push(this.memoryManager.tools.read(options.memory)); Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.name).join(', ')}
` : ''}`.trim());
}
if(mem.tool) tools.push(this.memoryManager.tools.read(mem.memory));
} }
} }
@@ -219,16 +318,33 @@ class LLM {
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')});
const resp = await request; let resp = await request;
// Spice delegated agents response into history
let lastDelegateResp: string | null = null;
if(pendingDelegates.size) {
for(let i = 0; i < history.length; i++) {
const h = history[i];
if(h.role !== 'tool' || h.content !== '') continue;
const queue = pendingDelegates.get(h.name);
if(!queue?.length) continue;
const {resp: delegateResp, subHistory} = queue.shift()!;
const insert: LLMMessage[] = [...subHistory.filter(sh => sh.role === 'tool'), {role: 'assistant', content: delegateResp, timestamp: Date.now()}];
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 // Trim memory injections from history
if(options.memory) { if(mem?.tool) history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'recall'));
history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'recall'));
}
// Auto-memorize before compressing // Auto-memorize before compressing
if(options.compress && this.estimateTokens(history) >= options.compress.max) { if(options.compress && this.estimateTokens(history) >= options.compress.max) {
if(options.memory) await this.memoryManager.memorize(history, options.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);
} }
@@ -397,15 +513,33 @@ class LLM {
* @param {string} searchTerms Multiple search terms to check against target * @param {string} searchTerms Multiple search terms to check against target
* @returns {{avg: number, max: number, similarities: number[]}} Similarity values 0-1: 0 = unique, 1 = identical * @returns {{avg: number, max: number, similarities: number[]}} Similarity values 0-1: 0 = unique, 1 = identical
*/ */
fuzzyMatch(target: string, ...searchTerms: string[]) { fuzzyMatch(target, ...searchTerms) {
if(searchTerms.length < 2) throw new Error('Requires at least 2 strings to compare'); if (searchTerms.length < 2) throw new Error('Requires at least 2 strings to compare');
const vector = (text: string, dimensions: number = 10): number[] => { const levenshtein = (a, b) => {
return text.toLowerCase().split('').map((char, index) => const m = a.length, n = b.length;
(char.charCodeAt(0) * (index + 1)) % dimensions / dimensions).slice(0, dimensions); if (!m) return n;
if (!n) return m;
const dp = Array.from({length: m + 1}, (_, i) => [i, ...Array(n).fill(0)]);
for (let j = 0; j <= n; j++) dp[0][j] = j;
for (let i = 1; i <= m; i++) {
for (let j = 1; j <= n; j++) {
dp[i][j] = a[i - 1] === b[j - 1]
? dp[i - 1][j - 1]
: 1 + Math.min(dp[i - 1][j - 1], dp[i - 1][j], dp[i][j - 1]);
} }
const v = vector(target); }
const similarities = searchTerms.map(t => vector(t)).map(refVector => this.cosineSimilarity(v, refVector)); return dp[m][n];
return {avg: similarities.reduce((acc, s) => acc + s, 0) / similarities.length, max: Math.max(...similarities), similarities}; };
const similarity = (a, b) => {
a = a.toLowerCase(); b = b.toLowerCase();
return 1 - levenshtein(a, b) / Math.max(a.length, b.length, 1);
};
const similarities = searchTerms.map(t => similarity(target, t));
return {
avg: similarities.reduce((acc, s) => acc + s, 0) / similarities.length,
max: Math.max(...similarities),
similarities
};
} }
/** /**

View File

@@ -1,59 +0,0 @@
import {KDPoint, KDTree} from './kd-tree.ts';
import {Memory, MemoryRef} from './memory.ts';
export class MemoryCache {
private tree: KDTree<MemoryRef>;
public memories: Memory[];
get length() { return this.memories.length; }
constructor(memories: Memory[]) {
this.memories = memories;
this.tree = this.buildTree();
}
private buildTree(): KDTree<MemoryRef> {
const embedded = this.memories.filter(m => m.embedding?.length);
if (!embedded.length) return new KDTree<MemoryRef>(0);
const dims = embedded[0].embedding.length;
const points: KDPoint<MemoryRef>[] = embedded.map(m => ({
vector: m.embedding,
payload: {name: m.name, description: m.description},
}));
return new KDTree<MemoryRef>(dims, 'cosine', points);
}
search(query: number[], limit: number): MemoryRef[] {
const results = this.tree.knn(query, limit);
return results.map(r => r.point.payload);
}
add(memory: Memory): void {
this.memories.push(memory);
this.rebuild();
}
update(memory: Memory): void {
const idx = this.memories.findIndex(m => m.name === memory.name);
if (idx !== -1) {
this.memories[idx] = memory;
} else {
this.memories.push(memory);
}
this.rebuild();
}
remove(name: string): void {
const idx = this.memories.findIndex(m => m.name === name);
if (idx !== -1) {
this.memories.splice(idx, 1);
this.rebuild();
}
}
rebuild(): void {
this.tree = this.buildTree();
}
}

View File

@@ -1,69 +0,0 @@
import {MemoryCache} from './memory-cache.ts';
import {extractMetadata, Memory, MemoryNode} from './memory.ts';
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 => {
const {links, backlinks} = extractMetadata(m.content);
return {
name: m.name,
missing: false,
links,
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,6 +1,143 @@
import {LLMRequest, LLMMessage} from './llm.ts'; import {LLMRequest, LLMMessage} from './llm.ts';
import {MemoryCache} from './memory-cache.ts';
import {AiTool} from './tools.ts'; import {AiTool} from './tools.ts';
import {KDPoint, KDTree} from './kd-tree.ts';
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 => {
const {links, backlinks} = extractMetadata(m.content);
return {
name: m.name,
missing: false,
links,
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();
}
export class MemoryCache {
private tree: KDTree<MemoryRef>;
public memories: Memory[];
get length() { return this.memories.length; }
constructor(memories: Memory[]) {
this.memories = memories;
this.tree = this.buildTree();
}
private buildTree(): KDTree<MemoryRef> {
const embedded = this.memories.filter(m => m.embedding?.length);
if (!embedded.length) return new KDTree<MemoryRef>(0);
const dims = embedded[0].embedding.length;
const points: KDPoint<MemoryRef>[] = embedded.map(m => ({
vector: m.embedding,
payload: {name: m.name, description: m.description},
}));
return new KDTree<MemoryRef>(dims, 'cosine', points);
}
search(query: number[], limit: number): MemoryRef[] {
const results = this.tree.knn(query, limit);
return results.map(r => r.point.payload);
}
add(memory: Memory): void {
this.memories.push(memory);
this.rebuild();
}
update(memory: Memory): void {
const idx = this.memories.findIndex(m => m.name === memory.name);
if (idx !== -1) {
this.memories[idx] = memory;
} else {
this.memories.push(memory);
}
this.rebuild();
}
remove(name: string): void {
const idx = this.memories.findIndex(m => m.name === name);
if (idx !== -1) {
this.memories.splice(idx, 1);
this.rebuild();
}
}
rebuild(): void {
this.tree = this.buildTree();
}
}
export type MemoryOptions = {
/** Memory object */
memory: Memory[] | MemoryCache;
/** Inject N memories into the system prompt */
inject?: boolean;
/** expose recall tool to LLM */
tool?: boolean;
/** Update memory on compression */
update?: boolean;
/** Max context size of memories to inject to each call (removed immediately after use) */
maxTokens?: number;
}
export type Memory = { export type Memory = {
name: string; name: string;
@@ -83,8 +220,6 @@ function getWeekSunday(monday: string): string {
return d.toISOString().slice(0, 10); return d.toISOString().slice(0, 10);
} }
export class MemoryManager { export class MemoryManager {
private pendingMemorizations = new Map<string, { private pendingMemorizations = new Map<string, {
memories: Memory[] | MemoryCache, memories: Memory[] | MemoryCache,
@@ -128,6 +263,12 @@ export class MemoryManager {
constructor(private llm: any) {} constructor(private llm: any) {}
static normalize(m?: Memory[] | MemoryCache | MemoryOptions) {
if(!m) return null;
const raw = m instanceof MemoryCache || Array.isArray(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 async createTempMemory(conversation: string): Promise<Memory> {
const timestamp = Date.now(); const timestamp = Date.now();
const content = `--- const content = `---

View File

@@ -29,11 +29,11 @@ export class OpenAi extends LLMProvider {
})); }));
history.splice(i, 1, ...tools); history.splice(i, 1, ...tools);
i += tools.length - 1; i += tools.length - 1;
} else if(h.role === 'tool' && h.content) { } else if(h.role === 'tool') {
const record = history.find(h2 => h.tool_call_id == h2.id); const record = history.find(h2 => h.tool_call_id == h2.id);
if(record) { if(record) {
if(h.content.includes('"error":')) record.error = h.content; if(h.content?.includes('"error":')) record.error = h.content;
else record.content = h.content; else record.content = h.content || '';
} }
history.splice(i, 1); history.splice(i, 1);
i--; i--;
@@ -51,15 +51,16 @@ export class OpenAi extends LLMProvider {
content: null, content: null,
tool_calls: [{ id: h.id, type: 'function', function: { name: h.name, arguments: JSON.stringify(h.args) } }], tool_calls: [{ id: h.id, type: 'function', function: { name: h.name, arguments: JSON.stringify(h.args) } }],
refusal: null, refusal: null,
annotations: [] annotations: [],
timestamp: h.timestamp,
}, { }, {
role: 'tool', role: 'tool',
tool_call_id: h.id, tool_call_id: h.id,
content: h.error || h.content content: h.error || h.content,
timestamp: h.timestamp,
}); });
} else { } else {
const {timestamp, ...rest} = h; result.push(h);
result.push(rest);
} }
return result; return result;
}, [] as any[]); }, [] as any[]);
@@ -106,8 +107,9 @@ export class OpenAi extends LLMProvider {
}; };
} }
let resp: any, isFirstMessage = true; let resp: any, isFirstMessage = true, terminal = false;
do { do {
requestParams.messages = history.map(({timestamp, ...m}) => m);
resp = await this.client.chat.completions.create(requestParams).catch(err => { resp = await this.client.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;
@@ -116,7 +118,7 @@ export class OpenAi extends LLMProvider {
if(options.stream) { if(options.stream) {
if(!isFirstMessage) options.stream({text: '\n\n'}); if(!isFirstMessage) options.stream({text: '\n\n'});
else isFirstMessage = false; else isFirstMessage = false;
resp.choices = [{message: {role: 'assistant', content: '', tool_calls: []}}]; resp.choices = [{message: {role: 'assistant', content: '', tool_calls: [], timestamp: Date.now()}}];
for await (const chunk of resp) { for await (const chunk of resp) {
if(controller.signal.aborted) break; if(controller.signal.aborted) break;
if(chunk.choices[0].delta.content) { if(chunk.choices[0].delta.content) {
@@ -158,28 +160,32 @@ export class OpenAi extends LLMProvider {
const results = await Promise.all(toolCalls.map(async (toolCall: any) => { const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
const tool = tools?.find(findByProp('name', toolCall.function.name)); const tool = tools?.find(findByProp('name', toolCall.function.name));
if(options.stream) options.stream({tool: toolCall.function.name}); if(options.stream) options.stream({tool: toolCall.function.name});
if(!tool) return {role: 'tool', tool_call_id: toolCall.id, content: '{"error": "Tool not found"}'}; if(!tool) return {role: 'tool', tool_call_id: toolCall.id, content: '{"error": "Tool not found"}', timestamp: Date.now()};
try { try {
const args = JSONAttemptParse(toolCall.function.arguments, {}); const args = JSONAttemptParse(toolCall.function.arguments, {});
const result = await tool.fn(args, options.stream, this.ai); // Wrap stream so a tool's `done` ends turn gracefully
return {role: 'tool', tool_call_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result}; const toolStream = options.stream && ((chunk: any) => {
if(chunk.done) { terminal = true; return; }
options.stream!(chunk);
});
const result = await tool.fn(args, toolStream, this.ai);
return {role: 'tool', tool_call_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result, timestamp: Date.now()};
} catch (err: any) { } catch (err: any) {
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'})}; return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'}), timestamp: Date.now()};
} }
})); }));
history.push(...results); history.push(...results);
requestParams.messages = history; requestParams.messages = history;
} }
} while (!controller.signal.aborted && resp.choices?.[0]?.message?.tool_calls?.length); } while (!terminal && !controller.signal.aborted && resp.choices?.[0]?.message?.tool_calls?.length);
if(!terminal) {
const textContent = resp.choices[0].message.content?.trim() || ''; const textContent = resp.choices[0].message.content?.trim() || '';
history.push({role: 'assistant', content: textContent}); history.push({role: 'assistant', content: textContent, timestamp: Date.now()});
}
history = this.toStandard(history); history = this.toStandard(history);
if(options.stream) options.stream({done: true}); if(options.stream) options.stream({done: true});
if(options.history) options.history.splice(0, options.history.length, ...history); if(options.history) options.history.splice(0, options.history.length, ...history);
// Return parsed JSON if schema provided
const finalContent = history.at(-1)?.content; const finalContent = history.at(-1)?.content;
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent); res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
}), {abort: () => controller.abort()}); }), {abort: () => controller.abort()});