Compare commits

..

11 Commits
1.3.4 ... 1.4.5

Author SHA1 Message Date
d42f58d710 Memory refinement
All checks were successful
Publish Library / Build NPM Project (push) Successful in 54s
Publish Library / Tag Version (push) Successful in 11s
2026-08-05 12:22:13 -04:00
878a8794ee Rebuild graph edges on changes
All checks were successful
Publish Library / Build NPM Project (push) Successful in 46s
Publish Library / Tag Version (push) Successful in 19s
2026-08-04 17:05:58 -04:00
3f1289d993 Small agent tweaks
All checks were successful
Publish Library / Build NPM Project (push) Successful in 49s
Publish Library / Tag Version (push) Successful in 9s
2026-08-04 14:33:28 -04:00
077f75cdd9 Fixed delegate agent history... again
All checks were successful
Publish Library / Build NPM Project (push) Successful in 48s
Publish Library / Tag Version (push) Successful in 13s
2026-08-04 13:58:47 -04:00
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 850 additions and 791 deletions

View File

@@ -186,7 +186,7 @@ console.log(chunks);
// Manually compile history into memories at end of conversation // Manually compile history into memories at end of conversation
// Happens automatically when coverstaions are compressed // Happens automatically when coverstaions are compressed
await ai.language.updateMemory(history, memory); await ai.language.memorize(history, memory);
// Summarize text // Summarize text
const summary = await ai.language.summarize(longText, 200); const summary = await ai.language.summarize(longText, 200);

View File

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

View File

@@ -1,61 +1,52 @@
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 toStandard(history: any[]): LLMMessage[] { private getClient(token: string): anthropic {
const timestamp = Date.now(); let client = this.clients.get(token);
const messages: LLMMessage[] = []; if(!client) {
for(let h of history) { client = new anthropic({apiKey: token});
if(typeof h.content == 'string') { this.clients.set(token, client);
messages.push(<any>{timestamp, ...h});
} else {
const textContent = h.content?.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
if(textContent) messages.push({role: h.role, content: textContent, timestamp: timestamp, duration: h.duration, tps: h.tps});
h.content.forEach((c: any) => {
if(c.type == 'tool_use') {
messages.push({role: 'tool', id: c.id, name: c.name, args: c.input, timestamp: h.timestamp, content: undefined, duration: h.duration, tps: h.tps});
} else if(c.type == 'tool_result') {
const m: any = messages.findLast(m => (<any>m).id == c.tool_use_id);
if(m) m[c.is_error ? 'error' : 'content'] = c.content;
} }
}); return client;
}
}
return messages;
} }
private fromStandard(history: LLMMessage[]): any[] { /** Convert standard history -> Anthropic wire format */
for(let i = 0; i < history.length; i++) { private toWire(history: LLMMessage[]): any[] {
if(history[i].role == 'tool') { const wire: any[] = [];
const h: any = history[i]; for(const h of history) {
history.splice(i, 1, if(h.role === 'tool') {
wire.push(
{role: 'assistant', content: [{type: 'tool_use', id: h.id, name: h.name, input: h.args}]}, {role: 'assistant', content: [{type: 'tool_use', id: h.id, name: h.name, input: h.args}]},
{role: 'user', content: [{type: 'tool_result', tool_use_id: h.id, is_error: !!h.error, content: h.error || h.content}]} {role: 'user', content: [{type: 'tool_result', tool_use_id: h.id, is_error: !!h.error, content: h.error || h.content || ''}]}
) );
i++; } else {
wire.push({role: h.role, content: h.content});
} }
} }
return history; return wire;
} }
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> { ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
const controller = new AbortController(); const controller = new AbortController();
return Object.assign(new Promise<any>(async (res) => { return Object.assign(new Promise<any>(async (res, rej) => {
let history = this.fromStandard([ if(!options.history) options.history = [];
...(options.history || []).filter(h => h.role !== 'system'), const history = options.history;
{role: 'user', content: message, timestamp: Date.now()} if(message) history.push({role: 'user', content: message, timestamp: Date.now()});
]);
const tools = options.tools || this.ai.options.llm?.tools || []; const tools = options.tools || this.ai.options.llm?.tools || [];
const requestParams: any = { const requestParams: any = {
model: options.model || this.model, model: options.model || this.model,
@@ -69,54 +60,43 @@ export class Anthropic extends LLMProvider {
type: 'object', type: 'object',
properties: t.args ? objectMap(t.args, (key, value) => ({...value, required: undefined})) : {}, properties: t.args ? objectMap(t.args, (key, value) => ({...value, required: undefined})) : {},
required: t.args ? Object.entries(t.args).filter(t => t[1].required).map(t => t[0]) : [] required: t.args ? Object.entries(t.args).filter(t => t[1].required).map(t => t[0]) : []
}, }
fn: undefined
})), })),
messages: history,
stream: !!options.stream, stream: !!options.stream,
}; };
// Add structured output support
if(options.schema) { if(options.schema) {
requestParams.output_config = { requestParams.output_config = {format: {type: 'json_schema', schema: convertSchema(options.schema)}};
format: {
type: 'json_schema',
schema: convertSchema(options.schema)
}
};
} }
let resp: any, terminal = false, duration = 0, tps = 0; try {
let terminal = false;
do { do {
requestParams.messages = history.map(({timestamp, ...m}) => m); requestParams.messages = this.toWire(history.filter(h => h.role !== 'system'));
const callStart = Date.now(); const callStart = Date.now();
resp = await this.client.messages.create(requestParams).catch(err => { const resp: any = 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(requestParams.messages, null, 2)}`;
throw err; throw err;
}); });
let usage: any; let usage: any, content: any[] = [];
if(options.stream) { if(options.stream) {
resp.content = [];
for await (const chunk of resp) { for await (const chunk of resp) {
if(controller.signal.aborted) break; if(controller.signal.aborted) break;
if(chunk.type === 'content_block_start') { if(chunk.type === 'content_block_start') {
if(chunk.content_block.type === 'text') { if(chunk.content_block.type === 'text') content.push({type: 'text', text: ''});
resp.content.push({type: 'text', text: ''}); else if(chunk.content_block.type === 'tool_use') content.push({type: 'tool_use', id: chunk.content_block.id, name: chunk.content_block.name, input: ''});
} else if(chunk.content_block.type === 'tool_use') {
resp.content.push({type: 'tool_use', id: chunk.content_block.id, name: chunk.content_block.name, input: <any>''});
}
} else if(chunk.type === 'content_block_delta') { } else if(chunk.type === 'content_block_delta') {
if(chunk.delta.type === 'text_delta') { if(chunk.delta.type === 'text_delta') {
const text = chunk.delta.text; content.at(-1).text += chunk.delta.text;
resp.content.at(-1).text += text; options.stream({text: chunk.delta.text});
options.stream({text});
} else if(chunk.delta.type === 'input_json_delta') { } else if(chunk.delta.type === 'input_json_delta') {
resp.content.at(-1).input += chunk.delta.partial_json; content.at(-1).input += chunk.delta.partial_json;
} }
} else if(chunk.type === 'content_block_stop') { } else if(chunk.type === 'content_block_stop') {
const last = resp.content.at(-1); const last = content.at(-1);
if(last?.input != null) last.input = last.input ? JSONAttemptParse(last.input, {}) : {}; if(last?.type === 'tool_use') last.input = last.input ? JSONAttemptParse(last.input, {}) : {};
} else if(chunk.type === 'message_delta') { } else if(chunk.type === 'message_delta') {
if(chunk.usage) usage = chunk.usage; if(chunk.usage) usage = chunk.usage;
} else if(chunk.type === 'message_stop') { } else if(chunk.type === 'message_stop') {
@@ -125,49 +105,52 @@ export class Anthropic extends LLMProvider {
} }
} else { } else {
usage = resp.usage; usage = resp.usage;
content = resp.content;
} }
duration = Date.now() - callStart; const duration = Date.now() - callStart;
tps = usage?.output_tokens && duration > 0 ? usage.output_tokens / (duration / 1000) : 0; const tps = usage?.output_tokens && duration > 0 ? usage.output_tokens / (duration / 1000) : 0;
const toolCalls = resp.content.filter((c: any) => c.type === 'tool_use'); const toolCalls = 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, timestamp: Date.now(), duration, tps}); const text = content.filter((c: any) => c.type === 'text').map((c: any) => c.text).join('\n\n').trim();
const results = await Promise.all(toolCalls.map(async (toolCall: any) => { if(text) history.push({role: 'assistant', content: text, timestamp: Date.now(), duration, tps});
const tool = tools.find(findByProp('name', toolCall.name));
if(options.stream) options.stream({tool: toolCall.name}); const entries = toolCalls.map((tc: any) => {
if(!tool) return {tool_use_id: toolCall.id, is_error: true, content: 'Tool not found'}; const entry: any = {role: 'tool', id: tc.id, name: tc.name, args: tc.input, content: undefined, timestamp: Date.now()};
history.push(entry);
return {tc, entry};
});
await Promise.all(entries.map(async ({tc, entry}: any) => {
const tool = tools.find(findByProp('name', tc.name));
if(options.stream) options.stream({tool: tc.name});
if(!tool) { entry.error = 'Tool not found'; return; }
try { try {
const toolStream = options.stream && ((chunk: any) => { const toolStream = options.stream && ((chunk: any) => {
if(chunk.done) { terminal = true; return; } if(chunk.done) { terminal = true; return; }
options.stream!(chunk); options.stream!(chunk);
}); });
const result = await tool.fn(toolCall.input, toolStream, this.ai, toolCall.id); const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
return {type: 'tool_result', tool_use_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result}; entry.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'}; entry.error = err?.message || err?.toString() || 'Unknown';
} }
})); }));
history.push({role: 'user', content: results, timestamp: Date.now()}); } else {
requestParams.messages = history; terminal = true;
const text = content.filter((c: any) => c.type === 'text').map((c: any) => c.text).join('\n\n').trim();
if(text) history.push({role: 'assistant', content: text, timestamp: Date.now(), duration, tps});
} }
} while (!terminal && !controller.signal.aborted && resp.content.some((c: any) => c.type === 'tool_use')); } while(!terminal && !controller.signal.aborted);
if(!terminal) {
const textContent = resp.content.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
history.push({role: 'assistant', content: textContent.trim(), timestamp: Date.now(), duration, tps});
}
history = this.toStandard(history);
if(options.history) options.history.splice(0, options.history.length, ...history);
if(options.stream) options.stream({done: true}); if(options.stream) options.stream({done: true});
const turnStart = history.map(h => h.role).lastIndexOf('user'); const turnStart = history.map(h => h.role).lastIndexOf('user');
const finalContent = history.slice(turnStart + 1).reduce((str, h) => { const finalContent = history.slice(turnStart + 1).reduce((str, h) => h.role === 'assistant' ? str + (h.content || '') : str, '').trim();
if(h.role === 'assistant') return str + (h.content || '');
return str;
}, '').trim();
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent); res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
} catch(err) {
rej(err);
}
}), {abort: () => controller.abort()}); }), {abort: () => controller.abort()});
} }
} }

85
src/helpers.ts Normal file
View File

@@ -0,0 +1,85 @@
import {Memory, MemoryCache} from './memory.ts';
export type MemoryNode = {
name: string;
missing: boolean;
links: string[];
backlinks: string[];
}
export function extractLinks(content: string): string[] {
if (!content) return [];
const matches = content.matchAll(/\[\[([^\]|]+)(?:\|[^\]]*)?\]\]/g);
return [...new Set([...matches].map(m => m[1].trim()))];
}
export function rebuildGraph(memories: Memory[] | MemoryCache): MemoryNode[] {
const mems = memories instanceof MemoryCache ? memories.memories : memories;
const nameSet = new Set(mems.map(m => m.name));
for (const m of mems) m.links = extractLinks(m.content).filter(l => l !== m.name);
for (const m of mems) m.backlinks = [];
for (const m of mems) {
for (const link of m.links) {
const target = mems.find(t => t.name === link);
if (target) target.backlinks.push(m.name);
}
}
const nodes: MemoryNode[] = mems.map(m => ({
name: m.name,
missing: false,
links: m.links,
backlinks: m.backlinks,
}));
const ghosts = new Set<string>();
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: MemoryNode[]): string {
if (!nodes.length) return 'No memories yet.';
const groups = new Map<string, (MemoryNode & {label: string})[]>();
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

@@ -103,9 +103,10 @@ class BoundedMaxHeap<T> {
export class KDTree<T = unknown> { export class KDTree<T = unknown> {
private root: KDNode<T> | null = null; private root: KDNode<T> | null = null;
private _size = 0; private _size = 0;
private readonly dims: number;
private readonly distanceFn: (a: number[], b: number[]) => number; private readonly distanceFn: (a: number[], b: number[]) => number;
readonly dims: number;
/** /**
* @param dims Dimensionality of all vectors (must be consistent). * @param dims Dimensionality of all vectors (must be consistent).
* @param metric Distance metric to use. Default: "euclidean". * @param metric Distance metric to use. Default: "euclidean".

View File

@@ -1,4 +1,4 @@
import {snakeCase} from '@ztimson/utils'; import {clean, 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 {OpenAi} from './open-ai.ts'; import {OpenAi} from './open-ai.ts';
@@ -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;
@@ -32,6 +34,10 @@ export type LLMMessage = {
content: string | any; content: string | any;
/** Timestamp */ /** Timestamp */
timestamp?: number; timestamp?: number;
/** Response duration in ms */
duration?: number;
/** Tokens per second */
tps?: number;
} | { } | {
/** Tool call */ /** Tool call */
role: 'tool'; role: 'tool';
@@ -104,8 +110,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 +126,35 @@ 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: clean<any>({
context: {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true}, context: !a.delegate ? {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true} : undefined,
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
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 q = 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 wrapped in a tool call that will be analysis by an LLM'} const request = this.ask(q, {
As a subagent, focus on executing your task completely using available tools and returning only the final result - no commentary, questions, or dialogue. system: `You are a specialized subagent being called from an orchestrator
${a.delegate ? 'Your output streams directly to the user for the remainder of this turn. You are mid conversation' : 'You are wrapped in a tool call that will be analysis by an LLM'}
Dispense with greetings and focus on your instructions using available tools and returning only the final result unless specifically instructed to converse
${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 +163,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;
@@ -217,7 +216,7 @@ ${a.system}`,
if(!skills?.length) return {prompt: '', tools: []}; if(!skills?.length) return {prompt: '', tools: []};
const list = skills.map(s => `- ${s.name}: ${s.description}`).join('\n'); const list = skills.map(s => `- ${s.name}: ${s.description}`).join('\n');
return { return {
prompt: `You have access to the following skill documents, use \`read_skill\` to access them:\n${list}`, prompt: `You have access to the following skill documents, whenever there is overlap between a question and a skill file, use \`skill_read\` to get instructions and background knowledge:\n${list}`,
tools: [{ tools: [{
name: 'skill_read', name: 'skill_read',
description: 'Read the full content of a skill/knowledge document', description: 'Read the full content of a skill/knowledge document',
@@ -266,10 +265,14 @@ ${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 || [];
if(message) history.push({role: 'user', content: message, timestamp: Date.now()});
// MCP // MCP
const mcp = options.mcp || this.ai.options?.llm?.mcp; const mcp = options.mcp || this.ai.options?.llm?.mcp;
@@ -289,8 +292,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);
@@ -313,18 +316,30 @@ ${a.system}`,
} else listed.push(r); } else listed.push(r);
} }
prompts.unshift(`You have access to the following memory files: prompts.unshift(`You have a background memory process which has prefetched relevant information${mem.update ? ' and will create new memories from this conversation' : ''} for you
${mems.map(m => `- ${m.name}: ${m.description}`).join('\n')} Assume it is perfect and never mention this process to anyone ever
Always use your memories to craft a personalized response, they contain links / [[wiki links]] which you use navigate between them
${mem.tool ? `You can access memory files via the \`memory_search\` and \`memory_recall\` tools
When you need information about the user, \`memory_recall\` \`People/User\` before asking (fetch if not included bellow)
When you need information not provided, attempt 1-3 \`memory_search\` calls with unique queries before asking` : ''}
${preloaded.length ? ` ${preloaded.length ? `
Relevant memories have been preloaded: Prefetched Memories (Most relevant first):
${preloaded.map(r => `
**${r.name}** ${preloaded.map(r => `Memory: ${r.name}
${r.description} Description: ${r.description}
Linked: ${[r.links, ...r.backlinks].join(', ')}
\`\`\`
${r.content} ${r.content}
`).join('\n---\n')} \`\`\``).join('\n\n')}` : ''}
` : ''}${listed.length ? `
Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.name).join(', ')} ${mem.tool && listed.length ? listed.map(r => `Memory: ${r.name}
` : ''}`.trim()); Description: ${r.description}
Linked: ${[r.links, ...r.backlinks].join(', ')}
<!-- Truncated -->`).join('\n\n') : ''}
${mem.tool ? `Full memory list:
${mems.map(m => `- ${m.name}: ${m.description}`).join('\n')}` : ''}`.trim())
} }
if(mem.tool) tools.push(this.memoryManager.tools.read(mem.memory)); if(mem.tool) tools.push(this.memoryManager.tools.read(mem.memory));
} }
@@ -332,62 +347,42 @@ 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('', {...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;
})(); })();
return Object.assign(promise, {abort}); return Object.assign(promise, {abort});
} }
/**
* Digest full conversation history into memory documents.
* Call on session end to persist the conversation.
*/
async updateMemory(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest = {}): Promise<Memory[]> {
return this.memoryManager.memorize(history, memories, {model: this.defaultModel, ...options});
}
/** /**
* Compress chat history to reduce context size * Compress chat history to reduce context size
* @param {LLMMessage[]} history Chatlog that will be compressed * @param {LLMMessage[]} history Chatlog that will be compressed
@@ -567,6 +562,14 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam
}; };
} }
/**
* Digest full conversation history into memory documents.
* Call on session end to persist the conversation.
*/
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest = {}): Promise<Memory[]> {
return this.memoryManager.memorize(history, memories, {model: this.defaultModel, ...options});
}
/** /**
* Create a summary of some text * Create a summary of some text
* @param {string} text Text to summarize * @param {string} text Text to summarize

View File

@@ -1,122 +1,35 @@
import {MemoryNode, rebuildGraph} from './helpers.ts';
import {LLMRequest, LLMMessage} from './llm.ts'; 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';
export class MemoryCache { const FACTS_HEADING = '## Facts';
private tree: KDTree<MemoryRef>;
public memories: Memory[];
get length() { return this.memories.length; } const GENERIC_TEMPLATE = `# {{Title}}
constructor(memories: Memory[]) { ## Summary
this.memories = memories;
this.tree = this.buildTree();
}
private buildTree(): KDTree<MemoryRef> { ## Details
const embedded = this.memories.filter(m => m.embedding?.length);
if (!embedded.length) return new KDTree<MemoryRef>(0);
const dims = embedded[0].embedding.length; ## Related`;
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;
description: string; description: string;
content: string; content: string;
embedding: number[]; embedding: number[];
}
export type MemoryRef = {
name: string;
description: string;
}
export type FactBucket = {
subject: string;
facts: string[];
}
export type MemoryNode = {
name: string;
missing: boolean;
links: string[]; links: string[];
backlinks: string[]; backlinks: string[];
} }
function extractLinks(content: string): string[] { type MemoryRef = {
if(!content) return []; name: string;
const matches = content.matchAll(/\[\[([^\]]+)\]\]/g); description: string;
return [...new Set([...matches].map(m => m[1].trim()))];
} }
export function extractMetadata(content: string): {links: string[], backlinks: string[]} { type FactBucket = {
const match = content.match(/^---\n([\s\S]*?)\n---/); subject: string;
if (!match) return {links: [], backlinks: []}; facts: string[];
const fm = match[1];
const getList = (key: string): string[] => {
const m = fm.match(new RegExp(`^${key}:\\s*\\[(.*)\\]$`, 'm'));
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[] {
@@ -139,51 +52,144 @@ function cosineDistance(a: number[], b: number[]): number {
return denom === 0 ? 1 : 1 - dot / denom; return denom === 0 ? 1 : 1 - dot / denom;
} }
function getWeekMonday(date: Date = new Date()): string { function cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] {
const d = new Date(Date.UTC(date.getFullYear(), date.getMonth(), date.getDate())); return memories
const day = d.getUTCDay(); .filter(m => m.embedding?.length)
const diff = day === 0 ? -6 : 1 - day; .map(m => ({ref: {name: m.name, description: m.description}, distance: cosineDistance(query, m.embedding)}))
d.setUTCDate(d.getUTCDate() + diff); .sort((a, b) => a.distance - b.distance)
return d.toISOString().slice(0, 10); .slice(0, limit)
.map(s => s.ref);
} }
function getWeekSunday(monday: string): string { export class MemoryCache {
const d = new Date(`${monday}T00:00:00Z`); private tree!: KDTree<MemoryRef>;
d.setUTCDate(d.getUTCDate() + 6); public memories: Memory[];
return d.toISOString().slice(0, 10); public nodes: MemoryNode[] = [];
get length() { return this.memories.length; }
constructor(memories: Memory[]) {
this.memories = memories;
this.rebuild();
}
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[] {
if (!this.tree || this.tree.dims === 0) return [];
return this.tree.knn(query, limit).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.nodes = rebuildGraph(this.memories);
this.tree = this.buildTree();
}
}
class MemoryAccessor {
readonly list: Memory[];
private readonly cache: MemoryCache | null;
constructor(memories: Memory[] | MemoryCache) {
this.cache = memories instanceof MemoryCache ? memories : null;
this.list = this.cache ? this.cache.memories : <Memory[]>memories;
}
find(name: string): Memory | undefined {
return this.list.find(m => m.name === name);
}
commit(): MemoryNode[] {
if (this.cache) {
this.cache.rebuild();
return this.cache.nodes;
}
return rebuildGraph(this.list);
}
ghosts(): string[] {
const nodes = this.cache ? this.cache.nodes : rebuildGraph(this.list);
return nodes.filter(n => n.missing).map(n => n.name);
}
search(vector: number[], limit: number): MemoryRef[] {
return this.cache ? this.cache.search(vector, limit) : cosineSearch(vector, this.list, limit);
}
forget(name: string): boolean {
const idx = this.list.findIndex(m => m.name === name);
if (idx === -1) return false;
this.list.splice(idx, 1);
this.commit();
return true;
}
async backfillEmbeddings(llm: any): Promise<number> {
const missing = this.list.filter(m => !m.embedding?.length);
if (!missing.length) return 0;
await Promise.all(missing.map(async node => {
const [e] = await llm.embedding(node.content);
if (e) node.embedding = e.embedding;
}));
this.commit();
return missing.length;
}
}
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 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>,
}>(); }>();
tools = { tools = {
read: (memories: Memory[] | MemoryCache): AiTool => ({
name: 'memory_recall',
description: 'Read the full content of a memory document',
args: {
name: {type: 'string', description: 'Exact memory name', required: true},
},
fn: (args: any) => {
const mems = memories instanceof MemoryCache ? memories.memories : memories;
const mem = mems.find(m => m.name === args.name);
if (!mem) return 'Document not found';
this.touch(mem.name);
return mem.content;
},
}),
forget: (memories: Memory[] | MemoryCache): AiTool => ({ forget: (memories: Memory[] | MemoryCache): AiTool => ({
name: 'memory_forget', name: 'memory_forget',
description: 'Permanently delete a memory document and clean up all references to it', description: 'Permanently delete a memory document and clean up all references to it',
@@ -195,381 +201,335 @@ export class MemoryManager {
return result ? `Forgotten: ${args.name}` : `Not found: ${args.name}`; return result ? `Forgotten: ${args.name}` : `Not found: ${args.name}`;
}, },
}), }),
read: (memories: Memory[] | MemoryCache): AiTool => ({
name: 'memory_recall',
description: 'Read the full content of a memory document',
args: {
name: {type: 'string', description: 'Exact memory name', required: true},
},
fn: (args: any) => {
const mem = this.access(memories).find(args.name);
if (!mem) return 'Document not found';
this.touch(mem.name);
return mem.content;
},
}),
search: (memories: Memory[] | MemoryCache): AiTool => ({
name: 'memory_search',
description: 'Use embeddings to find the MOST relevant memories, even if NOT relevant',
args: {
query: {type: 'string', description: 'What to look for in the memories', required: true},
limit: {type: 'number', description: 'Number of memories to return', default: 1},
},
fn: async ({query, limit}) => {
const mem = await this.recollect(query, memories, limit)
return mem.map(m => `Memory: ${m.name}
Description: ${m.description}
Links: ${[...m.links, ...m.backlinks].join(', ')}
\`\`\`
${m.content}
\`\`\``).join('\n\n');
},
}),
}; };
constructor(private llm: any) {} constructor(private llm: any) {}
static normalize(m?: Memory[] | MemoryCache | MemoryOptions) { static normalize(m?: Memory[] | MemoryCache | MemoryOptions) {
if(!m) return null; if (!m) return null;
const raw = m instanceof MemoryCache || Array.isArray(m); 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}; 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 access(memories: Memory[] | MemoryCache): MemoryAccessor {
const timestamp = Date.now(); return new MemoryAccessor(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 appendFacts(node: Memory, facts: string[]): void {
return `${header}\n\n${this.stripHeader(content)}`; 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);
} }
private async backgroundMemorization(conversation: string, memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string): Promise<void> { private ensureDoc(node: Memory): void {
const mem = memories instanceof MemoryCache ? memories.memories : memories; if (node.content) return;
const monday = getWeekMonday(); const title = node.name.split('/').pop() ?? node.name;
const sunday = getWeekSunday(monday); node.content = this.touchHeader(node, `# ${title}\n`);
const buckets = await this.factAgent(conversation, mem, options, monday);
if(!buckets.length) return;
const jobs = [...buckets].map(({subject, facts}) => {
let node = mem.find(m => m.name === subject);
if(!node) {
node = {name: subject, description: '', content: '', embedding: [],};
mem.push(node);
} }
const week = subject.startsWith('Journal/') ? {monday, sunday} : undefined;
return this.enqueue(node, facts, mem, options, tempName, week); private async factAgent(conversation: string, store: MemoryAccessor, options: LLMRequest, weekKey: string): Promise<FactBucket[]> {
const ghosts = store.ghosts();
const response = await this.llm.ask(conversation, {
model: options.model,
temperature: 0.2,
system: `You are a fact extractor to build obsidian knowledge vaults.
Analyze this conversation and extract facts worth remembering long-term.
Rules:
- Always extract facts that the user explicitly told you to remember
- ONLY extract current facts the USER explicitly stated about themselves, their work, projects or decisions that were MADE during this conversation
- DO NOT extract greetings, pleasantries, or generic exchanges
- DO NOT extract deltas or changes in facts; ONLY the end fact
- DO NOT extract anything the AI/assistant itself said
- If nothing worth remembering was said, return an empty buckets array
When extracting facts, you MUST also decide the exact destination path:
- Reuse node names (including ghost) as much as possible IF the facts belongs there
- 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) — you are not limited to any fixed list of collections, use whatever fits
- For journal entries, use "Journal"
Available nodes:
- Journal
${this.listNodes(store.list).filter(n => !n.name.includes('Journal')).map(n => `- ${n.name}: ${n.description}`).join('\n') || 'None yet.'}
${ghosts.length ? `${ghosts.map(g => `- ${g}: (Ghost)`).join('\n')}` : ''}`,
schema: {
buckets: {type: 'array', description: 'Groups of facts to remember, each assigned to a different node. Return an empty array if there is nothing worth storing in an obsidian vault', items: {
type: 'object', items: {
subject: {type: 'string', description: 'Exact existing node name OR new path (e.g. "People/Sarah", "Projects/Oxide"), or "Journal"', required: true},
facts: {
type: 'array',
description: 'Facts to store at this destination',
items: {type: 'string', description: 'A single fact'},
},
},
},
},
},
}); });
await Promise.all(jobs);
const buckets = new Map<string, string[]>();
for(const bucket of response.buckets ?? []) {
const subject = bucket.subject.trim().toLowerCase() === 'journal'
? `Journal/${weekKey}` : bucket.subject.trim();
const facts = buckets.get(subject) ?? [];
facts.push(...dedupeFacts(bucket.facts));
buckets.set(subject, facts);
} }
private buildHeader(node: Memory, week?: {monday: string, sunday: string}, links: string[] = [], backlinks: string[] = []): string { return buckets.entries().toArray().map(([subject, facts]) => ({subject, facts}));
const tags = node.name.split('/')[0]?.toLowerCase();
const lines = [
'---',
`name: ${node.name}`,
`description: ${node.description || ''}`,
tags ? `tags: [${tags}]` : '',
links.length ? `links: [${links.map(l => `"${l}"`).join(', ')}]` : 'links: []',
backlinks.length ? `backlinks: [${backlinks.map(l => `"${l}"`).join(', ')}]` : 'backlinks: []',
week ? `week: ${week.monday} ${week.sunday}` : '',
`modified: ${new Date().toISOString()}`,
'---',
].filter(Boolean);
return lines.join('\n');
} }
private cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] { private getWeekMonday(date: Date = new Date()): string {
const scored = memories const d = new Date(Date.UTC(date.getFullYear(), date.getMonth(), date.getDate()));
.filter(m => m.embedding?.length) const day = d.getUTCDay();
.map(m => ({ const diff = day === 0 ? -6 : 1 - day;
ref: {name: m.name, description: m.description}, d.setUTCDate(d.getUTCDate() + diff);
distance: cosineDistance(query, m.embedding), return d.toISOString().slice(0, 10);
}))
.sort((a, b) => a.distance - b.distance)
.slice(0, limit);
return scored.map(s => s.ref);
}
/**
* Coalescing queue: if a doc is already compiling, abort the in-flight run, merge its
* facts with the new ones and restart. Never blocks a pending update, never drops facts.
*/
private enqueue(node: Memory, facts: string[], memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string, week?: {monday: string, sunday: string}): Promise<void> {
const key = node.name;
const existing = this.queues.get(key);
if (existing) {
existing.pending.push(...facts);
existing.request?.abort?.();
return existing.task;
}
const entry: {pending: string[], request: {abort?: () => void} | null, task: Promise<void>} = {pending: [...facts], request: null, task: Promise.resolve()};
this.queues.set(key, entry);
const m = memories instanceof MemoryCache ? memories.memories : memories;
entry.task = (async () => {
while (entry.pending.length) {
const batch = dedupeFacts(entry.pending.splice(0, entry.pending.length));
const written = await this.docAgent(node, batch, m, options, tempName, week, entry);
if (!written) entry.pending.unshift(...batch);
}
})().finally(() => {
this.queues.delete(key);
if(!this.queues.size && memories instanceof MemoryCache) memories.rebuild();
});
return entry.task;
} }
private listNodes(memories: Memory[]): MemoryRef[] { private listNodes(memories: Memory[]): MemoryRef[] {
return memories.map(m => ({name: m.name, description: m.description})); return memories.map(m => ({name: m.name, description: m.description}));
} }
decay() { private reconcile(node: Memory, memories: Memory[] | MemoryCache, options: LLMRequest): Promise<void> {
for(const [name, ttl] of this.recentlyTouched) { const key = node.name;
if(ttl <= 1) this.recentlyTouched.delete(name); const existing = this.queues.get(key);
else this.recentlyTouched.set(name, ttl - 1); if (existing) {
} existing.dirty = true;
existing.request?.abort?.();
return existing.task;
} }
forget(name: string, memories: Memory[] | MemoryCache): boolean { const entry = {dirty: false, request: null, task: Promise.resolve()};
const mem = memories instanceof MemoryCache ? memories.memories : memories; this.queues.set(key, entry);
const idx = mem.findIndex(m => m.name === name); const store = this.access(memories);
if (idx === -1) return false; entry.task = (async () => {
do {
for (const node of mem) { entry.dirty = false;
const {links, backlinks} = extractMetadata(node.content); await this.docAgent(node, store.list, options, entry);
const newBacklinks = backlinks.filter(b => b !== name); } while (entry.dirty);
const newLinks = links.filter(l => l !== name); })().finally(() => {
this.queues.delete(key);
if (newBacklinks.length !== backlinks.length || newLinks.length !== links.length) { store.commit();
node.content = this.updateFrontmatter(node.content, {
links: newLinks,
backlinks: newBacklinks,
}); });
} return entry.task;
} }
mem.splice(idx, 1); private async docAgent(node: Memory, memories: Memory[], options: LLMRequest, entry: {request: {abort?: () => void} | null}): Promise<void> {
const currentBody = this.stripHeader(node.content);
if (memories instanceof MemoryCache) memories.rebuild(); let update;
return true;
}
getTouched(): string[] {
return [...this.recentlyTouched.keys()];
}
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise<Memory[]> {
const conversation = history
.filter(h => h.role === 'user' || h.role === 'assistant')
.map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim();
if(!conversation) return [];
const trackingId = `${Date.now()}_${Math.random()}`;
const tempMemory = await this.createTempMemory(conversation);
const mem = memories instanceof MemoryCache ? memories.memories : memories;
mem.push(tempMemory);
if (memories instanceof MemoryCache) memories.rebuild();
this.pendingMemorizations.set(trackingId, {
memories,
tempMemoryName: tempMemory.name,
timestamp: Date.now(),
});
try { try {
await this.backgroundMemorization(conversation, memories, options, tempMemory.name); for (let i = 0; i < 2 && !update?.content; i++) {
const finalMem = memories instanceof MemoryCache ? memories.memories : memories; const request = this.llm.ask(currentBody, {
return finalMem.filter(m => !m.name.startsWith('_temp_')); model: options.model,
temperature: 0.3,
schema: {
description: {type: 'string', description: 'One-line description of what this document covers, no formatting or emojis', required: true},
content: {type: 'string', description: 'Rewritten document body in markdown, without the frontmatter block', required: true},
},
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:
- 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]]
- 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
- Resolve contradictions: newer facts always win — delete the outdated statement entirely, never keep both
- Do not add frontmatter blocks, filler, preamble, or AI commentary
Other nodes in the vault (link to these instead of duplicating their content):
${this.listNodes(memories).filter(n => n.name !== node.name).map(n => n.name).join(', ') || 'none'}
Current document:
\`\`\`markdown
${currentBody}
\`\`\``,
});
entry.request = request;
update = await request;
}
} catch (err: any) {
if (err?.name === 'AbortError') return;
throw err;
} finally { } finally {
const pending = this.pendingMemorizations.get(trackingId); entry.request = null;
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[]> { if (!update?.content) return;
const mem: Memory[] = memories instanceof MemoryCache ? memories.memories : memories; node.description = node.name !== 'People/User' ? update.description : 'All information about the current user';
if (!mem.length) return []; node.content = this.touchHeader(node, update.content);
const [e] = await this.llm.embedding(node.content);
const [e] = await this.llm.embedding(query); if (e) node.embedding = e.embedding;
if (!e) return [];
let vectorResults: MemoryRef[];
if (memories instanceof MemoryCache) vectorResults = memories.search(e.embedding, limit);
else vectorResults = this.cosineSearch(e.embedding, mem, limit);
const found = new Set<string>(vectorResults.map(r => r.name));
if (graphDepth > 0) {
const frontier = [...found];
for (let depth = 0; depth < graphDepth; depth++) {
const next: string[] = [];
for (const name of frontier) {
const node = mem.find(m => m.name === name);
if (!node) continue;
const {links} = extractMetadata(node.content);
for (const link of links) {
if (!found.has(link) && mem.find(m => m.name === link)) {
found.add(link);
next.push(link);
}
}
}
frontier.splice(0, frontier.length, ...next);
if (!frontier.length) break;
}
} }
const vectorOrder = vectorResults.map(r => r.name); private parseFrontmatter(content: string): {fm: Map<string, string>, body: string} {
const graphExpansions = [...found].filter(n => !vectorOrder.includes(n)); const match = content.match(/^---\n([\s\S]*?)\n---\n?([\s\S]*)$/);
const ordered = [...vectorOrder, ...graphExpansions]; if (!match) return {fm: new Map(), body: content};
return ordered.map(n => mem.find(m => m.name === n)!).filter(Boolean); const fm = new Map<string, string>();
for (const line of match[1].split('\n')) {
const i = line.indexOf(':');
if (i === -1) continue;
fm.set(line.slice(0, i).trim(), line.slice(i + 1).trim());
} }
return {fm, body: match[2]};
touch(name: string, ttl = 2) {
this.recentlyTouched.set(name, ttl);
}
private updateFrontmatter(content: string, updates: {links?: string[], backlinks?: string[]}): string {
const match = content.match(/^---\n([\s\S]*?)\n---\n\n?([\s\S]*)$/);
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) {
const backlinksList = updates.backlinks.length ? `[${updates.backlinks.map(l => `"${l}"`).join(', ')}]` : '[]';
newFm = newFm.replace(/^backlinks:.*$/m, `backlinks: ${backlinksList}`);
}
newFm = newFm.replace(/^modified:.*$/m, `modified: ${new Date().toISOString()}`);
return `---\n${newFm}\n---\n\n${body}`;
} }
private stripHeader(content: string): string { private stripHeader(content: string): string {
return content.replace(/^---[\s\S]*?\n---\n?/, '').trimStart(); 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> { private touchHeader(node: Memory, body: string): string {
const {links: oldLinks} = extractMetadata(node.content); const {fm} = this.parseFrontmatter(node.content);
const currentBody = this.stripHeader(node.content); fm.set('name', node.name);
let update; fm.set('description', node.description || '');
try { fm.set('modified', new Date().toISOString());
for(let i = 0; i < 3 && !update?.content; i++) { return this.writeFrontmatter(fm, body);
const request = this.llm.ask(`New Facts:\n${facts.map(f => `- ${f}`).join('\n')}`, {
model: options.model,
temperature: 0.3,
schema: {
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},
},
system: `You are a knowledge base editor. Rewrite the current document below so it incorporates the new facts.
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
- 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)
- Keep the document concise, factual, and human-readable
- Resolve contradictions: the new 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
${week ? '- This is a weekly journal entry.\n' : ''}
All nodes:
${this.listNodes(memories).map(n => n.name).join(', ') || 'none'}
Current document:
\`\`\`markdown
${currentBody}
\`\`\``}
);
entry.request = request;
update = await request;
}
} catch (err: any) {
if (err?.name === 'AbortError') return false;
throw err;
} finally {
entry.request = null;
} }
if(!update?.content) return false; private writeFrontmatter(fm: Map<string, string>, body: string): string {
const newLinks = extractLinks(update.content).filter(l => l !== node.name && l !== tempName); const lines = [...fm.entries()].map(([k, v]) => `${k}: ${v}`);
const newLinkSet = new Set(newLinks); return `---\n${lines.join('\n')}\n---\n\n${body.trimStart()}`;
const oldLinkSet = new Set(oldLinks); }
for (const added of newLinkSet) { decay() {
if (!oldLinkSet.has(added)) { for (const [name, ttl] of this.recentlyTouched) {
const target = memories.find(m => m.name === added); if (ttl <= 1) this.recentlyTouched.delete(name);
if (target) { else this.recentlyTouched.set(name, ttl - 1);
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); touch(name: string, ttl = 2) {
node.description = node.name !== 'Person/User' ? update.description : 'All information about the current user'; this.recentlyTouched.set(name, ttl);
node.content = this.applyHeader(update.content, this.buildHeader(node, week, newLinks, backlinks)); }
forget(name: string, memories: Memory[] | MemoryCache): boolean {
return this.access(memories).forget(name);
}
async recollect(query: string, memories: Memory[] | MemoryCache, limit = 5, graphDepth = 1): Promise<Memory[]> {
const store = this.access(memories);
if (!store.list.length) return [];
await store.backfillEmbeddings(this.llm);
const [e] = await this.llm.embedding(query);
if (!e) return [];
const vectorResults = store.search(e.embedding, limit);
const found = new Set<string>(vectorResults.map(r => r.name));
if (graphDepth > 0) {
let frontier = [...found];
for (let depth = 0; depth < graphDepth && frontier.length; depth++) {
const next: string[] = [];
for (const name of frontier) {
const node = store.find(name);
if (!node) continue;
for (const link of node.links) {
if (!found.has(link) && store.find(link)) {
found.add(link);
next.push(link);
}
}
}
frontier = next;
}
}
const vectorOrder = vectorResults.map(r => r.name);
const graphExpansions = [...found].filter(n => !vectorOrder.includes(n));
return [...vectorOrder, ...graphExpansions].map(n => store.find(n)!).filter(Boolean);
}
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise<Memory[]> {
const conversation = history
.filter(h => h.role === 'user' || h.role === 'assistant')
.map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim();
if (!conversation) return [];
const uid = `${Date.now()}_${Math.random().toString(36).slice(2)}`;
const pending = {role: 'tool', name: 'memory_process', id: uid, content: conversation} as unknown as LLMMessage;
history.push(pending);
const store = this.access(memories);
const buckets = await this.factAgent(conversation, store, options, this.getWeekMonday());
const touched: Memory[] = [];
for (const {subject, facts} of buckets) {
let node = store.find(subject);
if (!node) {
node = {name: subject, description: '', content: '', embedding: [], links: [], backlinks: []};
store.list.push(node);
}
this.appendFacts(node, facts);
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; this.touch(node.name);
touched.push(node);
} }
private async factAgent(conversation: string, memories: Memory[], options: LLMRequest, weekKey: string): Promise<FactBucket[]> { if (touched.length) {
const buckets = new Map<string, string[]>(); store.commit();
await this.llm.ask(conversation, { (pending as any).content = `Saved to ${touched.map(n => `[[${n.name}]]`).join(', ')}`;
model: options.model, await Promise.all(touched.map(node => this.reconcile(node, memories, options).catch(() => {})));
temperature: 0.2, } else {
system: `You are a fact extractor. Analyze this conversation and extract facts worth remembering long-term. (pending as any).content = 'Nothing worth remembering.';
}
Rules: (touched as any).uid = uid;
- ONLY extract current facts the USER explicitly stated about themselves, their work, or their projects return touched;
- ONLY extract decisions that were MADE during this conversation }
- DO NOT extract anything the AI said, its capabilities, or meta-conversation about the AI
- DO NOT extract greetings, pleasantries, or generic exchanges
- DO NOT extract deltas or changes in facts; ONLY the end fact
- If nothing worth remembering was said, do not call any tools
When extracting facts, you MUST also decide the exact destination path: async reconcileVault(memories: Memory[] | MemoryCache, options: LLMRequest, scope: 'touched' | 'all' = 'touched'): Promise<void> {
- Use an existing node name if the facts clearly belong there const store = this.access(memories);
- All information primary about the user should go under "People/User" const targets = scope === 'all' ? store.list : store.list.filter(m => m.content.includes(FACTS_HEADING));
- When required, create a new path following collection/subject format (e.g., People/Sarah, Projects/Oxide) await Promise.all(targets.map(node => this.reconcile(node, memories, options)));
- For journal entries, use "Journal" store.commit();
Available nodes:
- Journal
${this.listNodes(memories).filter(n => !n.name.includes('_temp_') && !n.name.includes('Journal')).map(n => `- ${n.name}: ${n.description}`).join('\n') || 'None yet.'}`,
tools: [{
name: 'facts_extract',
description: 'Submit facts with their destination',
args: {
destination: {type: 'string', description: 'Exact existing node name OR new path (e.g. "People/Sarah", "Projects/Oxide")', required: true},
facts: {type: 'string', description: 'Comma-separated facts', required: true},
},
fn: (args: any) => {
const subject = args.destination.trim().toLowerCase() === 'journal'
? `Journal/${weekKey}` : args.destination.trim();
const facts = buckets.get(subject) ?? [];
facts.push(...dedupeFacts(String(args.facts).split(',')));
buckets.set(subject, facts);
return 'Recorded';
},
}],
});
return buckets.entries().toArray().map(([subject, facts]) => ({subject, facts}));
} }
} }

View File

@@ -1,88 +1,62 @@
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 toStandard(history: any[]): LLMMessage[] { private getClient(token: string): openAI {
for(let i = 0; i < history.length; i++) { let client = this.clients.get(token);
const h = history[i]; if(!client) {
if(h.role === 'assistant' && h.tool_calls) { client = new openAI(clean({baseURL: this.host, apiKey: token || undefined}));
const items: any[] = []; this.clients.set(token, client);
if(h.content) items.push({role: 'assistant', content: h.content, timestamp: h.timestamp, duration: h.duration, tps: h.tps});
items.push(...h.tool_calls.map((tc: any) => ({
role: 'tool',
id: tc.id,
name: tc.function.name,
args: JSONAttemptParse(tc.function.arguments, {}),
timestamp: h.timestamp,
duration: h.duration,
tps: h.tps
})));
history.splice(i, 1, ...items);
i += items.length - 1;
} else if(h.role === 'tool') {
const record = history.find(h2 => h.tool_call_id == h2.id);
if(record) {
if(h.content?.includes('"error":')) record.error = h.content;
else record.content = h.content || '';
} }
history.splice(i, 1); return client;
i--;
}
if(!history[i]?.timestamp) history[i].timestamp = Date.now();
}
return history;
} }
private fromStandard(history: LLMMessage[]): any[] { /** Convert standard history -> OpenAI wire format */
return history.reduce((result, h) => { private toWire(history: LLMMessage[], system?: string): any[] {
const wire: any[] = [];
if(system) wire.push({role: 'system', content: system});
for(const h of history) {
if(h.role === 'tool') { if(h.role === 'tool') {
result.push({ wire.push({
role: 'assistant', role: 'assistant',
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,
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 {
result.push(h); wire.push({role: h.role, content: h.content});
} }
return result; }
}, [] as any[]); return wire;
} }
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> { ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
const controller = new AbortController(); const controller = new AbortController();
return Object.assign(new Promise<any>(async (res, rej) => { return Object.assign(new Promise<any>(async (res, rej) => {
const base = (options.history || []).filter(h => h.role !== 'system'); if(!options.history) options.history = [];
let history = this.fromStandard([ const history = options.history;
...(options.system ? [{role: <any>'system', content: options.system, timestamp: Date.now()}] : []), if(message) history.push({role: 'user', content: message, timestamp: Date.now()});
...base,
{role: 'user', content: message, timestamp: Date.now()}
]);
const tools = options.tools || this.ai.options.llm?.tools || []; const tools = options.tools || this.ai.options.llm?.tools || [];
const requestParams: any = { const requestParams: any = {
model: options.model || this.model, model: options.model || this.model,
messages: history,
stream: !!options.stream, stream: !!options.stream,
max_completion_tokens: options.max_tokens || this.ai.options.llm?.max_tokens || undefined, max_completion_tokens: options.max_tokens || this.ai.options.llm?.max_tokens || undefined,
temperature: options.temperature || this.ai.options.llm?.temperature || undefined, temperature: options.temperature || this.ai.options.llm?.temperature || undefined,
@@ -102,56 +76,42 @@ export class OpenAi extends LLMProvider {
if(options.schema) { if(options.schema) {
const schema = convertSchema(options.schema); const schema = convertSchema(options.schema);
requestParams.response_format = { requestParams.response_format = {type: 'json_schema', json_schema: {name: 'response', strict: true, schema}};
type: 'json_schema',
json_schema: {
name: 'response',
strict: true,
schema
} }
};
}
if(options.stream) requestParams.stream_options = {include_usage: true}; if(options.stream) requestParams.stream_options = {include_usage: true};
let resp: any, terminal = false, duration = 0, tps = 0;
try {
let terminal = false;
do { do {
requestParams.messages = history.map(({timestamp, ...m}) => m); requestParams.messages = this.toWire(history.filter(h => h.role !== 'system'), options.system);
const callStart = Date.now(); const callStart = Date.now();
resp = await this.client.chat.completions.create(requestParams).catch(err => { const resp: any = 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(requestParams.messages, null, 2)}`;
throw err; throw err;
}); });
let usage: any; let usage: any, msg: any = {content: '', tool_calls: []};
if(options.stream) { if(options.stream) {
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.usage) usage = chunk.usage; if(chunk.usage) usage = chunk.usage;
if(chunk.choices[0]?.delta?.content) { if(chunk.choices[0]?.delta?.content) {
resp.choices[0].message.content += chunk.choices[0].delta.content; msg.content += chunk.choices[0].delta.content;
options.stream({text: chunk.choices[0].delta.content}); options.stream({text: chunk.choices[0].delta.content});
} }
if(chunk.choices[0]?.delta?.tool_calls) { if(chunk.choices[0]?.delta?.tool_calls) {
for(const deltaTC of chunk.choices[0].delta.tool_calls) { for(const deltaTC of chunk.choices[0].delta.tool_calls) {
const existing = resp.choices[0].message.tool_calls.find(tc => tc.index === deltaTC.index); const existing = msg.tool_calls.find((tc: any) => tc.index === deltaTC.index);
if(existing) { if(existing) {
if(deltaTC.id) existing.id = deltaTC.id; if(deltaTC.id) existing.id = deltaTC.id;
if(deltaTC.type) existing.type = deltaTC.type; if(deltaTC.function?.name) existing.function.name = deltaTC.function.name;
if(deltaTC.function) { if(deltaTC.function?.arguments) existing.function.arguments += deltaTC.function.arguments;
if(!existing.function) existing.function = {};
if(deltaTC.function.name) existing.function.name = deltaTC.function.name;
if(deltaTC.function.arguments) existing.function.arguments = (existing.function.arguments || '') + deltaTC.function.arguments;
}
} else { } else {
resp.choices[0].message.tool_calls.push({ msg.tool_calls.push({
index: deltaTC.index, index: deltaTC.index,
id: deltaTC.id || '', id: deltaTC.id || '',
type: deltaTC.type || 'function', function: {name: deltaTC.function?.name || '', arguments: deltaTC.function?.arguments || ''}
function: {
name: deltaTC.function?.name || '',
arguments: deltaTC.function?.arguments || ''
}
}); });
} }
} }
@@ -159,51 +119,51 @@ export class OpenAi extends LLMProvider {
} }
} else { } else {
usage = resp.usage; usage = resp.usage;
msg = resp.choices[0].message;
} }
duration = Date.now() - callStart; const duration = Date.now() - callStart;
tps = usage?.completion_tokens && duration > 0 ? usage.completion_tokens / (duration / 1000) : 0; const tps = usage?.completion_tokens && duration > 0 ? usage.completion_tokens / (duration / 1000) : 0;
if(resp.error) throw new Error(resp.error); const toolCalls = msg.tool_calls || [];
const toolCalls = resp.choices[0].message.tool_calls || [];
if(toolCalls.length && !controller.signal.aborted) { if(toolCalls.length && !controller.signal.aborted) {
history.push({...resp.choices[0].message, duration, tps}); if(msg.content?.trim()) history.push({role: 'assistant', content: msg.content.trim(), timestamp: Date.now(), duration, tps});
const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
const tool = tools?.find(findByProp('name', toolCall.function.name)); const entries = toolCalls.map((tc: any) => {
if(options.stream) options.stream({tool: toolCall.function.name}); const entry: any = {role: 'tool', id: tc.id, name: tc.function.name, args: JSONAttemptParse(tc.function.arguments, {}), content: undefined, timestamp: Date.now()};
if(!tool) return {role: 'tool', tool_call_id: toolCall.id, content: '{"error": "Tool not found"}', timestamp: Date.now()}; history.push(entry);
return {tc, entry};
});
await Promise.all(entries.map(async ({tc, entry}: any) => {
const tool = tools.find(findByProp('name', tc.function.name));
if(options.stream) options.stream({tool: tc.function.name});
if(!tool) { entry.error = 'Tool not found'; return; }
try { try {
const args = JSONAttemptParse(toolCall.function.arguments, {});
const toolStream = options.stream && ((chunk: any) => { const toolStream = options.stream && ((chunk: any) => {
if(chunk.done) { terminal = true; return; } if(chunk.done) { terminal = true; return; }
options.stream!(chunk); options.stream!(chunk);
}); });
const result = await tool.fn(args, toolStream, this.ai, toolCall.id); const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
return {role: 'tool', tool_call_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result, timestamp: Date.now()}; entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
} catch (err: any) { } catch(err: any) {
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'}), timestamp: Date.now()}; entry.error = err?.message || err?.toString() || 'Unknown';
} }
})); }));
history.push(...results); } else {
requestParams.messages = history; terminal = true;
const text = (msg.content || '').trim();
if(text) history.push({role: 'assistant', content: text, timestamp: Date.now(), duration, tps});
} }
} while (!terminal && !controller.signal.aborted && resp.choices?.[0]?.message?.tool_calls?.length); } while(!terminal && !controller.signal.aborted);
if(!terminal) {
const textContent = resp.choices[0].message.content || '';
history.push({role: 'assistant', content: textContent.trim(), timestamp: Date.now(), duration, tps});
}
history = this.toStandard(history);
if(options.history) options.history.splice(0, options.history.length, ...history.filter(h => h.role !== 'system'));
if(options.stream) options.stream({done: true}); if(options.stream) options.stream({done: true});
const turnStart = history.map(h => h.role).lastIndexOf('user'); const turnStart = history.map(h => h.role).lastIndexOf('user');
const finalContent = history.slice(turnStart + 1).reduce((str, h) => { const finalContent = history.slice(turnStart + 1).reduce((str, h) => h.role === 'assistant' ? str + (h.content || '') : str, '').trim();
if(h.role === 'assistant') return str + (h.content || '');
return str;
}, '').trim();
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent); res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
} catch(err) {
rej(err);
}
}), {abort: () => controller.abort()}); }), {abort: () => controller.abort()});
} }
} }

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);
}
}