Compare commits

...
62 Commits
Author SHA1 Message Date
ztimson 2921b208da More memory optimizations
Publish Library / Build NPM Project (push) Successful in 1m6s
Publish Library / Tag Version (push) Successful in 14s
2026-09-25 00:33:20 -04:00
ztimson b5aec246ac Memory refinement WIP 2026-09-24 14:25:06 -04:00
ztimson 2d6debad86 Memorization prompt tightening
Publish Library / Build NPM Project (push) Successful in 44s
Publish Library / Tag Version (push) Successful in 13s
2026-09-20 11:27:01 -04:00
ztimson 6bed8f20b5 Recursive agents update
Publish Library / Build NPM Project (push) Successful in 36s
Publish Library / Tag Version (push) Successful in 7s
2026-09-20 00:43:52 -04:00
ztimson dc45a99b04 Bump 1.6.13
Publish Library / Build NPM Project (push) Successful in 41s
Publish Library / Tag Version (push) Successful in 14s
2026-09-19 19:30:06 -04:00
ztimson 263a65c192 Fix open-ai early termination & memory improvements
Publish Library / Build NPM Project (push) Successful in 48s
Publish Library / Tag Version (push) Successful in 7s
2026-09-19 19:27:00 -04:00
ztimson 1e8c7c6662 Fix open-ai early termination
Publish Library / Build NPM Project (push) Successful in 1m47s
Publish Library / Tag Version (push) Successful in 8s
2026-09-19 13:22:48 -04:00
ztimson 1f1a4662d4 Entity based notes
Publish Library / Build NPM Project (push) Successful in 41s
Publish Library / Tag Version (push) Successful in 14s
2026-09-18 22:37:05 -04:00
ztimson ee4147e24e Fixed opanai early termination from tool calls
Publish Library / Build NPM Project (push) Successful in 46s
Publish Library / Tag Version (push) Successful in 15s
2026-09-18 16:07:58 -04:00
ztimson d29c0ca389 Fixed opanai early termination from tool calls
Publish Library / Build NPM Project (push) Successful in 55s
Publish Library / Tag Version (push) Successful in 9s
2026-09-18 02:15:02 -04:00
ztimson 4203cb34ef Better fact organization
Publish Library / Build NPM Project (push) Successful in 40s
Publish Library / Tag Version (push) Successful in 10s
2026-09-14 12:22:49 -04:00
ztimson d42c240362 Memorization optimziations
Publish Library / Build NPM Project (push) Successful in 1m2s
Publish Library / Tag Version (push) Successful in 10s
2026-08-31 12:38:40 -04:00
ztimson c1a16096ae Keep message progress on abort
Publish Library / Build NPM Project (push) Successful in 43s
Publish Library / Tag Version (push) Successful in 14s
2026-08-29 21:18:28 -04:00
ztimson ff0ee0b60e Patched memory merging
Publish Library / Build NPM Project (push) Successful in 45s
Publish Library / Tag Version (push) Successful in 10s
2026-08-28 16:48:46 -04:00
ztimson 0a6f1e4d62 Refined memory management prompts
Publish Library / Build NPM Project (push) Successful in 1m18s
Publish Library / Tag Version (push) Successful in 20s
2026-08-25 10:03:36 -04:00
ztimson 08a351e028 Better memory management
Publish Library / Build NPM Project (push) Successful in 59s
Publish Library / Tag Version (push) Successful in 11s
2026-08-24 14:42:10 -04:00
ztimson 85c01d3ef1 Added official file support
Publish Library / Build NPM Project (push) Successful in 30s
Publish Library / Tag Version (push) Successful in 10s
2026-08-17 15:50:48 -04:00
ztimson 5826573d5c Added official file support
Publish Library / Build NPM Project (push) Successful in 50s
Publish Library / Tag Version (push) Successful in 13s
2026-08-17 15:16:32 -04:00
ztimson 797a40a566 Added official file support
Publish Library / Build NPM Project (push) Successful in 58s
Publish Library / Tag Version (push) Successful in 13s
2026-08-16 15:40:50 -04:00
ztimson 7308927a3c max token rename
Publish Library / Build NPM Project (push) Successful in 35s
Publish Library / Tag Version (push) Successful in 14s
2026-08-05 16:16:30 -04:00
ztimson 04f038ba65 Memory prompt refinement
Publish Library / Build NPM Project (push) Successful in 38s
Publish Library / Tag Version (push) Successful in 19s
2026-08-05 13:14:21 -04:00
ztimson d42f58d710 Memory refinement
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
ztimson 878a8794ee Rebuild graph edges on changes
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
ztimson 3f1289d993 Small agent tweaks
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
ztimson 077f75cdd9 Fixed delegate agent history... again
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
ztimson 566d84fd7a Added memory graph traversal helpers
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
ztimson 4230b534fc bump 1.4.0
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
ztimson 119f8472f2 token pools
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
ztimson 9c04e58c63 Pass deligate subagents full history, improved memory managment 2026-08-04 12:24:23 -04:00
ztimson 7fbb42c26a improved subagent instructions 2026-08-04 12:03:31 -04:00
ztimson be08db8e2c Attach tps to response promise
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
ztimson 497f051c62 bump 1.3.5
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
ztimson 62fbe73b22 Added tps + duration to AI history
Publish Library / Build NPM Project (push) Successful in 52s
Publish Library / Tag Version (push) Successful in 11s
2026-08-04 09:26:57 -04:00
ztimson d53b1c6328 Removed <tool> blocks from responses
Publish Library / Build NPM Project (push) Successful in 39s
Publish Library / Tag Version (push) Successful in 11s
2026-08-03 20:23:22 -04:00
ztimson 89619e211e Fixed message history and response
Publish Library / Build NPM Project (push) Successful in 59s
Publish Library / Tag Version (push) Successful in 22s
2026-08-03 19:30:39 -04:00
ztimson afc6653364 fixed openai system calls in history breaking anthropic calls
Publish Library / Build NPM Project (push) Successful in 55s
Publish Library / Tag Version (push) Successful in 21s
2026-08-02 22:35:17 -04:00
ztimson 68e72445a2 Keep recent memories in context
Publish Library / Build NPM Project (push) Successful in 53s
Publish Library / Tag Version (push) Successful in 17s
2026-08-01 21:42:05 -04:00
ztimson 1aa6cdf329 Agent/subagent support
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
ztimson d022a5ef4d Improved levenshtein fuzzy match
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
ztimson a1d438a20a Tools can now emit "done" event and end chat early gracefully
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
ztimson 52a9e3aaa4 Fixed history poisoning on empty tool response
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
ztimson a7aec4ee29 Improved memory prompt slightly
Publish Library / Build NPM Project (push) Successful in 1m9s
Publish Library / Tag Version (push) Successful in 19s
2026-07-30 16:00:03 -04:00
ztimson dda2d4c2a3 Bump 1.2.8
Publish Library / Build NPM Project (push) Successful in 49s
Publish Library / Tag Version (push) Successful in 7s
2026-07-29 22:35:29 -04:00
ztimson 58e0e488e4 Added Geo, FS and flarescraperr tools
Publish Library / Tag Version (push) Has been cancelled
Publish Library / Build NPM Project (push) Has been cancelled
2026-07-29 22:34:51 -04:00
ztimson 8dfcd06752 More memory fixes
Publish Library / Build NPM Project (push) Successful in 43s
Publish Library / Tag Version (push) Successful in 14s
2026-07-29 22:11:09 -04:00
ztimson 14f6cdd313 Personal file memory organization instructions
Publish Library / Build NPM Project (push) Successful in 33s
Publish Library / Tag Version (push) Successful in 12s
2026-07-27 22:47:48 -04:00
ztimson 73d6ee0f2a Personal file memory organization instructions
Publish Library / Build NPM Project (push) Successful in 45s
Publish Library / Tag Version (push) Successful in 12s
2026-07-27 22:39:06 -04:00
ztimson bee4085666 updatememory awaits full result
Publish Library / Tag Version (push) Has been cancelled
Publish Library / Build NPM Project (push) Has been cancelled
2026-07-27 22:34:36 -04:00
ztimson 3b5c71de7c Improved memory management
Publish Library / Build NPM Project (push) Successful in 40s
Publish Library / Tag Version (push) Successful in 14s
2026-07-27 20:10:09 -04:00
ztimson 8229e02a52 Improved memory management
Publish Library / Build NPM Project (push) Successful in 44s
Publish Library / Tag Version (push) Successful in 11s
2026-07-27 14:25:24 -04:00
ztimson a6fb8ae828 New memory system
Publish Library / Build NPM Project (push) Successful in 1m5s
Publish Library / Tag Version (push) Successful in 11s
2026-07-27 03:59:39 -04:00
ztimson d1230bcaad Updated wiki tool
Publish Library / Build NPM Project (push) Successful in 55s
Publish Library / Tag Version (push) Successful in 14s
2026-07-26 12:18:57 -04:00
ztimson 2d49c9aa80 Removed redundant llama protocol (Use openai)
Publish Library / Build NPM Project (push) Successful in 44s
Publish Library / Tag Version (push) Successful in 13s
2026-07-11 19:33:02 -04:00
ztimson 9a39f00f94 Diarization fix
Publish Library / Build NPM Project (push) Successful in 44s
Publish Library / Tag Version (push) Successful in 13s
2026-07-11 18:36:24 -04:00
ztimson 436757daad Added new json output support
Publish Library / Build NPM Project (push) Failing after 1m2s
Publish Library / Tag Version (push) Has been skipped
2026-07-11 18:27:55 -04:00
ztimson 69b3297bb3 Proper error handling for OCR
Publish Library / Build NPM Project (push) Successful in 1m25s
Publish Library / Tag Version (push) Successful in 10s
2026-06-09 11:21:12 -04:00
ztimson 710c6ce52c Proper error handling for OCR
Publish Library / Tag Version (push) Has been cancelled
Publish Library / Build NPM Project (push) Has been cancelled
2026-06-09 11:20:50 -04:00
ztimson 4ac3036000 Proper error handling for OCR
Publish Library / Build NPM Project (push) Successful in 43s
Publish Library / Tag Version (push) Successful in 11s
2026-06-09 09:41:09 -04:00
ztimson 3121d542d4 OCR
Publish Library / Build NPM Project (push) Successful in 1m4s
Publish Library / Tag Version (push) Successful in 17s
2026-06-09 08:29:46 -04:00
ztimson 51ab8f2538 Memory / history fixes
Publish Library / Build NPM Project (push) Successful in 52s
Publish Library / Tag Version (push) Successful in 14s
2026-06-07 21:35:26 -04:00
ztimson 7dd3307a07 Update LLM models at runtime
Publish Library / Build NPM Project (push) Successful in 39s
Publish Library / Tag Version (push) Successful in 15s
2026-06-07 15:50:54 -04:00
ztimson 209d3b120b Export memory types
Publish Library / Build NPM Project (push) Successful in 1m7s
Publish Library / Tag Version (push) Successful in 13s
2026-06-07 13:06:45 -04:00
21 changed files with 3019 additions and 949 deletions
+2 -2
View File
@@ -119,7 +119,7 @@ const ai = new Ai({
system: 'You are a helpful assistant.', system: 'You are a helpful assistant.',
compress: {max: 90_000, min: 50_000}, // Compress chat history to min tokens when max is reached compress: {max: 90_000, min: 50_000}, // Compress chat history to min tokens when max is reached
temperature: 0.8, temperature: 0.8,
max_tokens: 100_000, maxTokens: 100_000,
memoryModel: 'gpt-4o', // Cheap model for managing memories in background, defaults to current model memoryModel: 'gpt-4o', // Cheap model for managing memories in background, defaults to current model
models: { models: {
'claude-3-5-sonnet': {proto: 'anthropic', token: process.env.ANTHROPIC_TOKEN}, 'claude-3-5-sonnet': {proto: 'anthropic', token: process.env.ANTHROPIC_TOKEN},
@@ -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);
-21
View File
@@ -1,21 +0,0 @@
import {Ai} from './dist/index.mjs';
const ai = new Ai({
path: './',
llm: {
system: 'You are a testbed for developing an AI library',
models: {
'qwen/qwen3.5-9b': {proto: 'openai', host: 'http://127.0.0.1:1234/v1'}
}
}
});
const skills = [{
name: 'Momentum',
description: 'Learn how to use the Momentum API',
content: 'You can initialize it with: new Momentum(url);'
}];
const history = [], memory = [];
console.log(await ai.language.ask('Can you tell me how to use momentum?', {history, skills}));
console.log(history, memory);
+444 -280
View File
File diff suppressed because it is too large Load Diff
+4 -3
View File
@@ -1,6 +1,6 @@
{ {
"name": "@ztimson/ai-utils", "name": "@ztimson/ai-utils",
"version": "1.0.0", "version": "1.7.2",
"description": "AI Utility library", "description": "AI Utility library",
"author": "Zak Timson", "author": "Zak Timson",
"license": "MIT", "license": "MIT",
@@ -26,12 +26,13 @@
}, },
"dependencies": { "dependencies": {
"@anthropic-ai/sdk": "^0.102.0", "@anthropic-ai/sdk": "^0.102.0",
"@tensorflow/tfjs": "^4.22.0",
"@huggingface/transformers": "^4.2.0", "@huggingface/transformers": "^4.2.0",
"@tensorflow/tfjs": "^4.22.0",
"@ztimson/node-utils": "^1.0.7", "@ztimson/node-utils": "^1.0.7",
"@ztimson/utils": "^0.29.4", "@ztimson/utils": "^0.30.8",
"cheerio": "^1.2.0", "cheerio": "^1.2.0",
"openai": "^6.42.0", "openai": "^6.42.0",
"pdf-parse": "^2.4.5",
"tesseract.js": "^7.0.0" "tesseract.js": "^7.0.0"
}, },
"devDependencies": { "devDependencies": {
+3 -3
View File
@@ -1,10 +1,10 @@
import * as os from 'node:os'; import * as os from 'node:os';
import LLM, {AnthropicConfig, OllamaConfig, OpenAiConfig, LLMRequest} from './llm'; import LLM, {AnthropicConfig, OpenAiConfig, LLMRequest} from './llm';
import { Audio } from './audio.ts'; import { Audio } from './audio.ts';
import {Vision} from './vision.ts'; import {Vision} from './vision.ts';
export type AbortablePromise<T> = Promise<T> & { export type AbortablePromise<T> = Promise<T> & {
abort: () => any abort: (keep?: boolean) => any
}; };
export type AiOptions = { export type AiOptions = {
@@ -18,7 +18,7 @@ export type AiOptions = {
embedder?: string; embedder?: string;
/** Large language models, first is default */ /** Large language models, first is default */
llm?: Omit<LLMRequest, 'model'> & { llm?: Omit<LLMRequest, 'model'> & {
models: {[model: string]: AnthropicConfig | OllamaConfig | OpenAiConfig}; models: {[model: string]: AnthropicConfig | OpenAiConfig};
} }
/** OCR model: eng, eng_best, eng_fast */ /** OCR model: eng, eng_best, eng_fast */
ocr?: string; ocr?: string;
+100 -76
View File
@@ -1,63 +1,65 @@
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';
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({timestamp, role: h.role, content: textContent});
h.content.forEach((c: any) => {
if(c.type == 'tool_use') {
messages.push({timestamp, role: 'tool', id: c.id, name: c.name, args: c.input, content: undefined});
} 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[] { private toWireContent(content: any): any {
for(let i = 0; i < history.length; i++) { if(!Array.isArray(content)) return content;
if(history[i].role == 'tool') { return content.map(c => c.type === 'image'
const h: any = history[i]; ? {type: 'image', source: {type: 'base64', media_type: c.mime, data: c.data}}
history.splice(i, 1, : {type: 'text', text: c.text});
}
/** Convert standard history -> Anthropic wire format */
private toWire(history: LLMMessage[]): any[] {
const wire: any[] = [];
for(const h of history) {
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: this.toWireContent(h.content)});
} }
} }
return history.map(({timestamp, ...h}) => h); return wire;
} }
ask(message: string, options: LLMRequest = {}): AbortablePromise<string> { 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([...options.history || [], {role: 'user', content: message, timestamp: Date.now()}]); if(!options.history) options.history = [];
const history = options.history;
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,
max_tokens: options.max_tokens || this.ai.options.llm?.max_tokens || 4096, max_tokens: options.maxTokens || this.ai.options.llm?.maxTokens || 4096,
system: options.system || this.ai.options.llm?.system || '', system: options.system || this.ai.options.llm?.system || '',
temperature: options.temperature || this.ai.options.llm?.temperature || 0.7, temperature: options.temperature || this.ai.options.llm?.temperature || undefined,
tools: tools.map(t => ({ tools: tools.map(t => ({
name: t.name, name: t.name,
description: t.description, description: t.description,
@@ -65,75 +67,97 @@ 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,
}; };
let resp: any, isFirstMessage = true; if(options.schema) {
requestParams.output_config = {format: {type: 'json_schema', schema: convertSchema(options.schema)}};
}
try {
let terminal = false;
do { do {
resp = await this.client.messages.create(requestParams).catch(err => { requestParams.messages = this.toWire(history.filter(h => h.role !== 'system'));
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
const callStart = Date.now();
const resp: any = await this.tokenPool.run(token => this.getClient(token).messages.create(requestParams)).catch(err => {
err.message += `\n\nMessages:\n${JSON.stringify(requestParams.messages, null, 2)}`;
throw err; throw err;
}); });
// Streaming mode let usage: any, content: any[] = [];
if(options.stream) { if(options.stream) {
if(!isFirstMessage) options.stream({text: '\n\n'});
else isFirstMessage = false;
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') {
if(chunk.usage) usage = chunk.usage;
} else if(chunk.type === 'message_stop') { } else if(chunk.type === 'message_stop') {
break; break;
} }
} }
} else {
usage = resp.usage;
content = resp.content;
} }
const duration = Date.now() - callStart;
const tps = usage?.output_tokens && duration > 0 ? usage.output_tokens / (duration / 1000) : 0;
// Run tools const toolCalls = 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}); 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 result = await tool.fn(toolCall.input, options?.stream, this.ai); const toolStream = options.stream && ((chunk: any) => {
return {type: 'tool_result', tool_use_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result}; if(chunk.done) { terminal = true; return; }
} catch (err: any) { options.stream!(chunk);
return {type: 'tool_result', tool_use_id: toolCall.id, is_error: true, content: err?.message || err?.toString() || 'Unknown'}; });
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
} catch(err: any) {
entry.error = err?.message || err?.toString() || 'Unknown';
} }
})); }));
history.push({role: 'user', content: results}); } 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 (!controller.signal.aborted && resp.content.some((c: any) => c.type === 'tool_use')); } while(!terminal && !controller.signal.aborted);
history.push({role: 'assistant', content: resp.content.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n')});
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);
res(history.at(-1)?.content); const turnStart = history.map(h => h.role).lastIndexOf('user');
const finalContent = history.slice(turnStart + 1).reduce((str, h) => h.role === 'assistant' ? str + (h.content || '') : str, '').trim();
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
} catch(err) {
rej(err);
}
}), {abort: () => controller.abort()}); }), {abort: () => controller.abort()});
} }
} }
+10 -3
View File
@@ -141,11 +141,18 @@ print(json.dumps(segments))
if(!llm) return transcript; if(!llm) return transcript;
let chunks = this.ai.language.chunk(transcript, 500, 0); let chunks = this.ai.language.chunk(transcript, 500, 0);
if(chunks.length > 4) chunks = [...chunks.slice(0, 3), <string>chunks.at(-1)]; if(chunks.length > 4) chunks = [...chunks.slice(0, 3), <string>chunks.at(-1)];
const names = await this.ai.language.json(chunks.join('\n'), '{1: "Detected Name", 2: "Second Name"}', { await this.ai.language.ask(chunks.join('\n'), {
system: 'Use the following transcript to identify speakers. Only identify speakers you are positive about, dont mention speakers you are unsure about in your response', system: 'Read the following transcript and attempt to identify every speaker. For every positively identified speaker, call the \`identify\` tool with the speaker\'s ID number & the identified name exactly once.',
temperature: 0.1, temperature: 0.1,
tools: [
{name: 'identify', description: 'Identify a speaker', args: {
speaker: {type: 'number', description: 'Speaker number', required: true},
name: {type: 'string', description: 'Inferred name', required: true},
}, fn: ({speaker, name}) => {
transcript = transcript.replaceAll(`[Speaker ${speaker}]`, `[${name}]`);
}}
]
}); });
Object.entries(names).forEach(([speaker, name]) => transcript = transcript.replaceAll(`[Speaker ${speaker}]`, `[${name}]`));
return transcript; return transcript;
} }
+6
View File
@@ -2,7 +2,13 @@ export * from './ai';
export * from './antrhopic'; export * from './antrhopic';
export * from './audio'; export * from './audio';
export * from './llm'; export * from './llm';
export * from './memory/graph';
export * from './memory/kd-tree';
export * from './memory/memory';
export * from './memory/memory-state';
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';
export * from './utils';
+437 -102
View File
@@ -1,24 +1,72 @@
import {clean, makeUnique, 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/memory-state.ts';
import {Memory, MemoryManager, MemoryOptions} from './memory/memory.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} from './tools.ts'; import {AiTool, AiToolArg} from './tools.ts';
import {fileURLToPath} from 'url'; import {fileURLToPath} from 'url';
import {dirname, join} from 'path';
import {spawn} from 'node:child_process'; import {spawn} from 'node:child_process';
import {Memory, MemoryManager} from './memory.ts'; import {mkdtempSync} from 'node:fs';
import fs from 'node:fs/promises';
import {tmpdir} from 'node:os';
import {dirname, join, basename, extname} from 'path';
import { PDFParse } from 'pdf-parse';
import {stripHeader} from './utils.ts';
export type AnthropicConfig = {proto: 'anthropic', token: string}; const MAX_AGENT_DEPTH = 5;
export type OllamaConfig = {proto: 'ollama', host: string}; const PDF_OCR_PAGE_THRESHOLD = 12; // above this many pages, OCR scanned pages instead of feeding images to the model
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 AgentRef = {
name: string;
description?: string;
delegate?: boolean;
fn: () => Agent | null | Promise<Agent | null>;
}
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;
agents?: AgentRef[] | null;
}
export type LLMFile = {
/** Path to file on disk */
path?: string;
/** File content: raw text, base64-encoded binary, or a Buffer */
content?: string | Buffer;
/** Original filename, used to infer type from extension */
name?: string;
/** Mime type override, inferred from extension if omitted */
mime?: string;
/** @internal set once extraction has run, skips re-processing next turn */
extracted?: boolean;
};
export type LLMMessage = { export type LLMMessage = {
/** Message originator */ /** Message originator */
role: 'assistant' | 'system' | 'user'; role: 'assistant' | 'system' | 'user';
/** Message content */ /** Message content */
content: string | any; content: string | any;
/** Files attached to request */
files?: LLMFile[];
/** 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';
@@ -34,15 +82,21 @@ export type LLMMessage = {
error?: undefined | string; error?: undefined | string;
/** Timestamp */ /** Timestamp */
timestamp?: number; timestamp?: number;
/** Response duration in ms */
duration?: number;
/** Tokens per second */
tps?: number;
} }
export type LLMRequest = { export type LLMRequest = {
/** Return a parsed JSON object that matches the schema */
schema?: AiToolArg;
/** System prompt */ /** System prompt */
system?: string; system?: string;
/** Message history */ /** Message history */
history?: LLMMessage[]; history?: LLMMessage[];
/** Max tokens for request */ /** Max tokens for request */
max_tokens?: number; maxTokens?: number;
/** 0 = Rigid Logic, 1 = Balanced, 2 = Hyper Creative **/ /** 0 = Rigid Logic, 1 = Balanced, 2 = Hyper Creative **/
temperature?: number; temperature?: number;
/** Available tools */ /** Available tools */
@@ -54,13 +108,19 @@ 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[]; 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, resolved lazily via their `fn` */
agents?: AgentRef[];
/** Attach files to request */
files?: LLMFile[];
/** @internal recursion guard for nested agent delegation */
_agentDepth?: number;
} }
export type McpServer = { export type McpServer = {
@@ -81,8 +141,12 @@ export type Skill = {
content: string; content: string;
} }
class LLM { class LLM {
private static AUDIO_EXT = ['wav','mp3','m4a','flac','ogg','aac','wma'];
private static IMAGE_EXT = ['png','jpg','jpeg','bmp','gif','tiff','webp'];
private static TEXT_EXT = ['txt','md','csv','json','xml','html','js','ts','py','yaml','yml','log'];
private static PDF_EXT = ['pdf'];
private memoryManager!: MemoryManager; private memoryManager!: MemoryManager;
defaultModel!: string; defaultModel!: string;
@@ -93,12 +157,170 @@ class LLM {
Object.entries(ai.options.llm.models).forEach(([model, config]) => { Object.entries(ai.options.llm.models).forEach(([model, config]) => {
if(!this.defaultModel) this.defaultModel = model; if(!this.defaultModel) this.defaultModel = model;
if(config.proto == 'anthropic') this.models[model] = new Anthropic(this.ai, config.token, model); if(config.proto == 'anthropic') this.models[model] = new Anthropic(this.ai, config.token, model);
else if(config.proto == 'ollama') this.models[model] = new OpenAi(this.ai, config.host, 'not-needed', model);
else if(config.proto == 'openai') this.models[model] = new OpenAi(this.ai, config.host || null, config.token, model); else if(config.proto == 'openai') this.models[model] = new OpenAi(this.ai, config.host || null, config.token, model);
}); });
this.memoryManager = new MemoryManager(this); this.memoryManager = new MemoryManager(this);
} }
private async loadBuffer(file: LLMFile, asText: boolean): Promise<Buffer> {
if(file.path) return fs.readFile(file.path);
if(Buffer.isBuffer(file.content)) return file.content;
if(typeof file.content === 'string') return Buffer.from(file.content, asText ? 'utf-8' : 'base64');
throw new Error('No path or content provided');
}
private async writeTemp(name: string, buffer: Buffer): Promise<string> {
const path = join(mkdtempSync(join(tmpdir(), 'ai-file-')), name);
await fs.writeFile(path, buffer);
return path;
}
/**
* Extract text from a PDF. Pages with no text layer (scanned/image-only) are handled as either:
* - Rendered to images and returned alongside the text so the (vision-capable) model can read them directly
* - OCR'd via Tesseract when the doc is too large to reasonably pass as images
*/
private async resolvePdf(buffer: Buffer): Promise<{text: string, images: {mime: string, data: string}[]}> {
const parser = new PDFParse({data: buffer});
try {
const {text, pages} = await parser.getText();
const scanned = (pages || []).filter(p => !p.text?.trim());
if(!scanned.length) return {text: text.trim() || '[Empty PDF]', images: []};
const total = pages.length;
const pageNums = scanned.map(p => p.num);
const {pages: shots} = await parser.getScreenshot({partial: pageNums});
if(total <= PDF_OCR_PAGE_THRESHOLD) {
return {
text: text.trim(),
images: shots.map(s => ({mime: 'image/png', data: Buffer.from(s.data).toString('base64')}))
};
}
const ocrText = await Promise.all(shots.map(async (s, i) => {
const path = await this.writeTemp(`page-${pageNums[i]}.png`, Buffer.from(s.data));
try {
return await this.ai.vision.ocr(path) || '';
} finally {
fs.rm(dirname(path), {recursive: true, force: true}).catch(() => {});
}
}));
return {text: [text.trim(), ...ocrText].filter(Boolean).join('\n\n'), images: []};
} finally {
await parser.destroy();
}
}
private async resolveFile(file: LLMFile): Promise<{text?: string, images?: {mime: string, data: string}[]}> {
const name = file.name || (file.path ? basename(file.path) : 'file');
// Already resolved on a previous turn, reuse cached text
if(file.extracted) return {text: `<file name="${name}">\n${file.content}\n</file>`};
const ext = extname(name).slice(1).toLowerCase();
const mime = file.mime || '';
const isAudio = mime.startsWith('audio/') || LLM.AUDIO_EXT.includes(ext);
const isImage = mime.startsWith('image/') || LLM.IMAGE_EXT.includes(ext);
const isPdf = mime === 'application/pdf' || LLM.PDF_EXT.includes(ext);
const isText = mime.startsWith('text/') || LLM.TEXT_EXT.includes(ext);
let tmpDir: string | null = null;
try {
if(isImage) {
const data = (await this.loadBuffer(file, false)).toString('base64');
return {images: [{mime: mime || `image/${ext === 'jpg' ? 'jpeg' : ext}`, data}]};
}
if(isPdf) {
const {text, images} = await this.resolvePdf(await this.loadBuffer(file, false));
// Only cache/skip re-processing when we didn't need to hand off images (OCR'd or fully text-based)
if(!images.length) {
file.content = text;
file.extracted = true;
delete file.path;
}
return {text: `<file name="${name}">\n${text || '[Scanned PDF - see attached page images]'}\n</file>`, images};
}
let text: string;
if(isAudio) {
let path = file.path;
if(!path) {
const buffer = await this.loadBuffer(file, false);
path = await this.writeTemp(name, buffer);
tmpDir = dirname(path);
}
text = await this.ai.audio.asr(path) || '';
} else if(isText) {
text = (await this.loadBuffer(file, true)).toString('utf-8');
} else {
text = typeof file.content === 'string' ? file.content : `[Binary file, unable to extract: ${name}]`;
}
file.content = text;
file.extracted = true;
delete file.path;
return {text: `<file name="${name}">\n${text}\n</file>`};
} catch(err: any) {
return {text: `<file name="${name}">Failed to process: ${err.message}</file>`};
} finally {
if(tmpDir) fs.rm(tmpDir, {recursive: true, force: true}).catch(() => {});
}
}
private async resolveFiles(files: LLMFile[]): Promise<{text: string, images: {mime: string, data: string}[]}> {
const resolved = await Promise.all(files.map(f => this.resolveFile(f)));
return {
text: resolved.filter(r => r.text).map(r => r.text).join('\n\n'),
images: resolved.flatMap(r => r.images || [])
};
}
private setupAgent(stubs: AgentRef[] = [], history: LLMMessage[], aborts: ((keep?: boolean) => void)[], depth = 0, delegateState: {resp: string | null}): AiTool[] {
return stubs.map(stub => {
const toolName = `${stub.delegate ? '' : 'sub'}agent_${snakeCase(stub.name)}`;
return {
name: toolName,
description: `${stub.delegate ? 'Delegate to ' : ''}Subagent: ${stub.description || stub.name}`,
args: clean<any>({
context: !stub.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},
}),
fn: async (args: any, stream: any, ai: any, id?: string) => {
if(depth >= MAX_AGENT_DEPTH) return 'Max agent delegation depth exceeded';
const a = await stub.fn();
if(!a) return `Agent "${stub.name}" could not be resolved`;
const q = a.delegate ? '' : `${args.instructions}${args.context ? `\n\n<context>${args.context}</context>` : ''}`;
const request = this.ask(q, {
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}`,
model: a.model || undefined,
temperature: a.temperature,
stream: a.delegate ? stream : undefined,
history: a.delegate ? history : [],
mcp: a.mcp || undefined,
skills: a.skills || undefined,
tools: a.tools || undefined,
agents: a.agents || [],
_agentDepth: depth + 1,
} as any);
aborts.push(request.abort);
const resp = await request;
if(a.delegate) {
delegateState.resp = resp;
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[] = [];
@@ -132,7 +354,7 @@ class LLM {
const list = allTools.map(t => `- ${t.name}: ${t.description}`).join('\n'); const list = allTools.map(t => `- ${t.name}: ${t.description}`).join('\n');
return { return {
prompt: `You have access to the following MCP tools:\n${list}`, prompt: `## MCP\nYou have access to the following MCP tools:\n${list}`,
tools: allTools tools: allTools
}; };
} }
@@ -141,9 +363,9 @@ class LLM {
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: `## Skills\nYou 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: 'read_skill', name: 'skill_read',
description: 'Read the full content of a skill/knowledge document', description: 'Read the full content of a skill/knowledge document',
args: { args: {
name: {type: 'string', description: 'Exact skill name', required: true} name: {type: 'string', description: 'Exact skill name', required: true}
@@ -157,10 +379,23 @@ class LLM {
} }
} }
private wrapToolTiming(tools: AiTool[], timings: Map<string, {duration: number, tps: number}>): AiTool[] {
return tools.map(t => ({
...t,
fn: async (args: any, stream: any, ai: any, id?: string) => {
const start = Date.now();
const result = await t.fn(args, stream, ai, id);
const duration = Date.now() - start;
const tps = duration > 0 ? this.estimateTokens(result) / (duration / 1000) : 0;
if(id) timings.set(id, {duration, tps});
return result;
}
}));
}
ask(message: string, options: LLMRequest = {}): AbortablePromise<string> { ask(message: string, options: LLMRequest = {}): AbortablePromise<string> {
options = <any>{ options = <any>{
system: '', system: '',
temperature: 0.8,
...this.ai.options.llm, ...this.ai.options.llm,
models: undefined, models: undefined,
history: [], history: [],
@@ -168,11 +403,42 @@ class LLM {
} }
const m = options.model || this.defaultModel; const m = options.model || this.defaultModel;
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 abort = () => {}; let request: AbortablePromise<string> | null = null;
return Object.assign(new Promise<string>(async res => { let aborted = false;
let keepOnAbort = true;
const nestedAborts: ((keep?: boolean) => void)[] = [];
const abort = (keep = true) => {
aborted = true;
keepOnAbort = keep;
request?.abort?.(keep);
nestedAborts.forEach(a => a(keep));
};
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[] = [options.system || this.ai.options.llm?.system || '']; const prompts: string[] = [];
if(!options.history) options.history = []; let history = options.history || [];
const historyStart = history.length;
const files = options.files || [];
if(message || files.length) history.push({role: 'user', content: message || '', timestamp: Date.now()});
// Accumulate streamed text so it can be committed to history if aborted mid-generation
let partialText = '';
const onStream = options.stream;
const stream = (chunk: {text?: string, tool?: string, done?: true}) => {
if(chunk.text) partialText += chunk.text;
return onStream?.(chunk);
};
/** Commit (keep) or discard this turn's progress on abort, then throw */
const abortNow = (): never => {
if(keepOnAbort) { if(partialText) history.push({role: 'assistant', content: partialText, timestamp: Date.now()}); }
else history.splice(historyStart, history.length - historyStart);
throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
};
// MCP // MCP
const mcp = options.mcp || this.ai.options?.llm?.mcp; const mcp = options.mcp || this.ai.options?.llm?.mcp;
@@ -190,44 +456,116 @@ class LLM {
tools.push(...s.tools); tools.push(...s.tools);
} }
// Agents
const agents = options.agents || this.ai.options?.llm?.agents;
const delegateState: {resp: string | null} = {resp: null};
if(agents?.length) tools.push(...this.setupAgent(agents, history, nestedAborts, options._agentDepth || 0, delegateState));
// Memory // Memory
if(options.memory) { const mem = MemoryManager.normalize(options.memory);
const relevant = await this.memoryManager.recollect(message, options.memory); if(mem) {
if(relevant.length) { const mems = mem.memory instanceof MemoryCache ? mem.memory.memories : mem.memory;
const context = relevant.map(m => `### ${m.name}\n${m.content}`).join('\n\n'); if(mems.length) {
options.history.push({ if(mem.inject) {
id: 'auto_recall_' + Math.random().toString(), role: 'tool', name: 'recall', args: {}, const pool = 15;
content: `Knowledge Documents:\n\n${context}` const budget = mem.maxTokens ?? 2000;
}); const relevant = await this.memoryManager.recollect(message, mem.memory, pool);
}
prompts.unshift('You have access to a knowledge base. Relevant documents are injected automatically before each message. Use this knowledge to inform your responses.'); 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);
} }
const resp = await this.models[m].ask(message, {...options, tools, system: prompts.filter(Boolean).join('\n\n')}); prompts.unshift(`## Memory
You have a background memory process which has prefetched relevant information${mem.update ? ' and will create new memories from this conversation' : ''} for you
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 distinct queries before asking` : ''}
// Trim memory injections from history ${preloaded.length ? `### Prefetched Memories (Most relevant first):
if(options.memory) {
options.history.splice(0, options.history.length, ...options.history.filter(h => ${preloaded.map(r => `Memory: ${r.name}
h.role !== 'tool' || h.name !== 'recall')); Description: ${r.description}
Linked: ${makeUnique([...r.links, ...r.backlinks]).join(', ')}
\`\`\`
${stripHeader(r.content)}
\`\`\``).join('\n\n')}` : ''}
${mem.tool && listed.length ? '\n' + listed.map(r => `Memory: ${r.name}
Description: ${r.description}
Linked: ${makeUnique([...r.links, ...r.backlinks]).join(', ')}
<!-- Truncated -->`).join('\n\n') : ''}`.trim())
}
if(mem.tool) tools.push(...this.memoryManager.tools.read(mem.memory));
}
} }
// Auto-memorize before compressing if(aborted) abortNow();
if(options.compress) {
if(options.memory) await this.memoryManager.memorize(options.history, options.memory, options); const lastMsg = history[history.length - 1];
const compressed = await this.compressHistory(options.history, options.compress.max, options.compress.min, options); if(files.length && lastMsg?.role === 'user') lastMsg.files = files;
options.history.splice(0, options.history.length, ...compressed); const restores: {msg: LLMMessage, content: any}[] = [];
for(const msg of history) {
if(msg.role !== 'user' || !msg.files?.length) continue;
const {text, images} = await this.resolveFiles(msg.files);
if(!text && !images.length) continue;
restores.push({msg, content: msg.content});
const merged = text ? [msg.content, text].filter(Boolean).join('\n\n') : msg.content;
msg.content = images.length
? [...images.map(i => ({type: 'image', mime: i.mime, data: i.data})), {type: 'text', text: merged}]
: merged;
} }
return res(resp); const toolTimings = new Map<string, {duration: number, tps: number}>();
}), {abort}); tools = this.wrapToolTiming(tools, toolTimings);
if(aborted) abortNow();
prompts.unshift(options.system || this.ai.options.llm?.system || '');
request = this.models[m].ask('', {...options, tools, stream, system: prompts.filter(Boolean).join('\n\n')});
let resp: string;
try {
resp = await request;
} catch(err: any) {
if(aborted) return abortNow();
throw err;
} }
/** // Strip the file injection shim
* Digest full conversation history into memory documents. restores.forEach(({msg, content}) => msg.content = content);
* Call on session end to persist the conversation.
*/ // Capture meta (duration / tps)
async updateMemory(history: LLMMessage[], memories: Memory[], options: LLMRequest = {}): Promise<void> { for(const h of history) {
await this.memoryManager.memorize(history, memories, {model: this.defaultModel, ...options}); if(h.role === 'tool' && toolTimings.has(h.id)) Object.assign(h, toolTimings.get(h.id));
}
if(typeof resp === 'string' && !resp.trim() && delegateState.resp !== null) resp = delegateState.resp;
if(mem?.tool) history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'memory_recall'));
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});
const compressed = await this.compressHistory(history, options.compress.max, options.compress.min, options);
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 Object.assign(promise, {abort});
} }
/** /**
@@ -258,24 +596,6 @@ class LLM {
return h; return h;
} }
/**
* Compare the difference between embeddings (calculates the angle between two vectors)
* @param {number[]} v1 First embedding / vector comparison
* @param {number[]} v2 Second embedding / vector for comparison
* @returns {number} Similarity values 0-1: 0 = unique, 1 = identical
*/
cosineSimilarity(v1: number[], v2: number[]): number {
if (v1.length !== v2.length) throw new Error('Vectors must be same length');
let dotProduct = 0, normA = 0, normB = 0;
for (let i = 0; i < v1.length; i++) {
dotProduct += v1[i] * v2[i];
normA += v1[i] * v1[i];
normB += v2[i] * v2[i];
}
const denominator = Math.sqrt(normA) * Math.sqrt(normB);
return denominator === 0 ? 0 : dotProduct / denominator;
}
/** /**
* Chunk text into parts for AI digestion * Chunk text into parts for AI digestion
* @param {object | string} target Item that will be chunked (objects get converted) * @param {object | string} target Item that will be chunked (objects get converted)
@@ -380,49 +700,41 @@ 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
};
} }
/** /**
* Ask a question with JSON response * Digest full conversation history into memory documents.
* @param {string} text Text to process * Call on session end to persist the conversation.
* @param {string} schema JSON schema the AI should match
* @param {LLMRequest} options Configuration options and chat history
* @returns {Promise<{} | {} | RegExpExecArray | null>}
*/ */
async json(text: string, schema: string, options?: LLMRequest): Promise<any> { async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest = {}): Promise<Memory[]> {
let system = `Your job is to convert input to JSON using tool calls. Call the \`submit\` tool at least once with JSON matching this schema:\n\`\`\`json\n${schema}\n\`\`\`\n\nResponses are ignored`; return this.memoryManager.memorize(history, memories, {model: this.defaultModel, ...options});
if(options?.system) system += '\n\n' + options.system;
return new Promise(async (resolve, reject) => {
let done = false;
const resp = await this.ask(text, {
temperature: 0.3,
...options,
system,
tools: [{
name: 'submit',
description: 'Submit JSON',
args: {json: {type: 'string', description: 'Javascript parsable JSON string', required: true}},
fn: (args) => {
try {
const json = JSON.parse(args.json);
resolve(json);
done = true;
} catch { return 'Invalid JSON'; }
return 'Saved';
}
}, ...(options?.tools || [])],
});
if(!done) reject(`AI failed to create JSON:\n${resp}`);
});
} }
/** /**
@@ -458,6 +770,29 @@ class LLM {
if(!done) reject(`AI failed to create summary:\n${resp}`); if(!done) reject(`AI failed to create summary:\n${resp}`);
}); });
} }
addModel(name: string, config: AnthropicConfig | OpenAiConfig, setDefault = false) {
if(config.proto == 'anthropic') this.models[name] = new Anthropic(this.ai, config.token, name);
else if(config.proto == 'openai') this.models[name] = new OpenAi(this.ai, config.host || null, config.token, name);
if(setDefault || !this.defaultModel) this.defaultModel = name;
}
removeModel(name: string) {
delete this.models[name];
if(this.defaultModel === name) {
this.defaultModel = Object.keys(this.models)[0] ?? '';
}
}
setModels(models: {[model: string]: AnthropicConfig | OpenAiConfig}, replace = true) {
if(replace) this.models = {};
Object.entries(models).forEach(([model, config]) => {
if(!this.defaultModel) this.defaultModel = model;
if(config.proto == 'anthropic') this.models[model] = new Anthropic(this.ai, config.token, model);
else if(config.proto == 'openai') this.models[model] = new OpenAi(this.ai, config.host || null, config.token, model);
});
this.defaultModel = Object.keys(this.models)[0] ?? '';
}
} }
export default LLM; export default LLM;
-177
View File
@@ -1,177 +0,0 @@
// memory.ts
import {LLMRequest, LLMMessage} from './llm.ts';
/** Background information the AI will be fed as a knowledge document */
export type Memory = {
/** Memory subject */
name: string;
/** Short description of what this document contains - used for RAG retrieval */
description: string;
/** Full markdown content of the document */
content: string;
/** Embedding vector of the description - used for similarity search */
embedding: number[];
}
export type MemoryCollection = {
/** Memory subject */
name: string;
/** Short description - required if isNew */
description?: string;
/** Extracted facts to merge */
facts: string[];
}
export class MemoryManager {
tools = {
edit: (memory: Memory) => ({
name: 'edit',
description: 'Edit a memory. Omit start/end to append. Pass start only to replace from that line on. Pass start+end to replace a specific range. start=0 replaces the whole document.',
args: {
content: {type: 'string', description: 'New content', required: true},
start: {type: 'number', description: 'First line to replace (0-indexed, inclusive). Omit to append.'},
end: {type: 'number', description: 'Last line to replace (0-indexed, inclusive). Omit to replace from start to end of doc.'},
},
fn: (args: any) => {
const lines = memory.content ? memory.content.split('\n') : [];
const newLines = args.content.split('\n');
if(args.start === undefined) lines.push(...newLines);
else if(args.end === undefined) lines.splice(args.start, lines.length - args.start, ...newLines);
else lines.splice(args.start, args.end - args.start + 1, ...newLines);
memory.content = lines.join('\n');
return `Updated memory:\n${memory.content}`;
}
}),
extract: (pools: MemoryCollection[]) => ({
name: 'extract',
description: 'Extract a list of facts to group into a single memory',
args: {
name: {type: 'string', description: 'Exact name of an existing memory, or a new name if none fits ([pro]nouns only)', required: true},
description: {type: 'string', description: 'One sentence description of the memory subject, only required if new'},
facts: {type: 'string', description: 'Comma separated list of extracted facts', required: true},
},
fn: (args: any) => {
pools.push({
name: args.name,
description: args.description,
facts: args.facts.split(',').map((f: string) => f.trim()).filter(Boolean),
});
return 'Success';
}}),
read: (memories: Memory[]) => ({
name: 'read',
description: 'Read entire memory',
args: {
name: {type: 'string', description: 'Exact memory name', required: true},
},
fn: (args: any) => {
const mem = memories.find(m => m.name === args.name);
if(!mem) return 'Document not found';
return `Name: ${mem.name}\nDescription: ${mem.description}\n\n${mem.content}`;
}
}),
}
constructor(private llm: any, private model?: string) {}
/**
* Extracts facts from conversation and groups them into individual memories
* @param {string} conversation Full conversation formatted as [role]: content
* @param {Memory[]} memories The user's memory documents
* @param {LLMRequest} options LLM options
* @returns {Promise<MemoryCollection[]>} Fact pools grouped by target document
*/
private async extract(conversation: string, memories: Memory[], options: LLMRequest): Promise<MemoryCollection[]> {
const existingDocs = memories.map(m => `Name: ${m.name}\nDescription: ${m.description}`).join('\n\n');
const pools: MemoryCollection[] = [];
await this.llm.ask(conversation, {
model: this.model || options.model,
temperature: 0.2,
system: `You are a fact extractor. Analyze this conversation and extract facts worth remembering long term.
Rules:
- ONLY extract facts the USER explicitly stated about themselves or their business
- ONLY extract decisions that were MADE during this conversation
- DO NOT extract anything the AI said, its name, capabilities, or how it introduced itself
- DO NOT extract greetings, pleasantries or generic exchanges
- If nothing worth remembering was said, call NO tools
For each fact decide whether it belongs in an existing document or needs a new one, then call the \`extract\` tool.
Existing documents:\n${existingDocs || 'None yet.'}`,
tools: [this.tools.extract(pools)]
});
return pools;
}
/**
* Bot 2 - Editor: merges a pool of facts into a specific document using surgical line-based edits.
* Receives full document content and uses read + amend tools to make precise edits.
* @param {MemoryCollection} newMem The fact pool to merge
* @param {Memory[]} memories The user's memory documents
* @param {LLMRequest} options LLM options
*/
private async edit(newMem: MemoryCollection, memories: Memory[], options: LLMRequest): Promise<void> {
const existing = memories.find(m => m.name === newMem.name);
const mem: Memory = existing || {name: newMem.name, description: newMem.description || '', content: '', embedding: []};
const isNew = !existing;
await this.llm.ask(newMem.facts.map(f => `- ${f}`).join('\n'),
{
model: this.model || options.model,
temperature: 0.2,
system: `You are a document editor. Merge the users list of facts into the following document using the \`edit\` tool; call it as many times as necessary.
Name: ${mem.name}
Description: ${mem.description}
${mem.content}`,
tools: [this.tools.edit(mem)]
}
);
if(isNew || mem.description !== existing?.description) {
const [e] = await this.llm.embedding(mem.description);
mem.embedding = e.embedding;
}
if(isNew) memories.push(mem);
else {
const idx = memories.findIndex(m => m.name === newMem.name);
if(idx >= 0) memories[idx] = mem;
}
}
/**
* Find relevant memory documents for a query using description embeddings
* @param {string} query The query to search against
* @param {Memory[]} memories The user's memory documents
* @param {number} limit Max number of results to return
* @returns {Promise<Memory[]>} The most relevant memory documents
*/
async recollect(query: string, memories: Memory[], limit = 5): Promise<Memory[]> {
const [e] = await this.llm.embedding(query);
return memories
.filter(m => m.embedding?.length)
.map(m => ({...m, score: this.llm.cosineSimilarity(m.embedding, e.embedding)}))
.toSorted((a: any, b: any) => b.score - a.score)
.slice(0, limit);
}
/**
* Two-stage memory pipeline: classify facts from conversation history then surgically merge them into documents.
* Bot 1 (classify) extracts and groups facts cheaply. Bot 2 (edit) runs per-document in parallel with full content access.
* @param {LLMMessage[]} history Full conversation history to digest
* @param {Memory[]} memories The user's memory documents — mutated in place
* @param {LLMRequest} options LLM options
*/
async memorize(history: LLMMessage[], memories: Memory[], options: LLMRequest): Promise<void> {
const conversation = history
.filter(h => h.role === 'user' || h.role === 'assistant')
.map(h => `[${h.role}]: ${h.content}`)
.join('\n\n');
if(!conversation.trim()) return;
const pools = await this.extract(conversation, memories, options);
if(!pools.length) return;
await Promise.all(pools.map(pool => this.edit(pool, memories, options)));
}
}
+128
View File
@@ -0,0 +1,128 @@
import {MemoryCache} from './memory-state.ts';
import type {Memory} 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 patchGraph(mems: Memory[], nodes: MemoryNode[], changed: Memory[]): MemoryNode[] {
const nameSet = new Set(mems.map(m => m.name));
const byName = new Map(nodes.map(n => [n.name, n]));
const ensureNode = (name: string): MemoryNode => {
let n = byName.get(name);
if (!n) {
n = {name, missing: !nameSet.has(name), links: [], backlinks: []};
byName.set(name, n);
}
return n;
};
for (const m of changed) {
const node = ensureNode(m.name);
node.missing = false; // real memory, promotes any pre-existing ghost entry
const oldLinks = m.links ?? [];
const newLinks = extractLinks(m.content).filter(l => l !== m.name);
for (const target of oldLinks.filter(l => !newLinks.includes(l))) {
const t = byName.get(target);
if (!t) continue;
t.backlinks = t.backlinks.filter(n => n !== m.name);
if (t.missing && !t.backlinks.length) byName.delete(target); // fully dereferenced ghost
}
for (const target of newLinks.filter(l => !oldLinks.includes(l))) {
const t = ensureNode(target);
if (!t.backlinks.includes(m.name)) t.backlinks.push(m.name);
}
m.links = newLinks;
node.links = newLinks;
}
for (const m of mems) {
const n = byName.get(m.name);
if (n) m.backlinks = n.backlinks;
}
return [...byName.values()];
}
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();
}
+353
View File
@@ -0,0 +1,353 @@
import {cosineDistance, euclideanDistance} from '../utils.ts';
export type DistanceMetric = "euclidean" | "cosine";
export interface KDPoint<T = unknown> {
vector: number[];
payload: T;
}
export interface KNNResult<T = unknown> {
point: KDPoint<T>;
distance: number;
}
interface KDNode<T> {
point: KDPoint<T>;
axis: number;
left: KDNode<T> | null;
right: KDNode<T> | null;
deleted?: boolean;
}
/**
* Keeps the k closest candidates in memory, evicts the furthest when full
*/
class BoundedMaxHeap<T> {
private heap: KNNResult<T>[] = [];
constructor(private readonly k: number) {}
get size(): number { return this.heap.length; }
get worstDistance(): number {
return this.heap.length < this.k ? Infinity : this.heap[0].distance;
}
push(item: KNNResult<T>): void {
if (this.heap.length < this.k) {
this.heap.push(item);
this.bubbleUp(this.heap.length - 1);
} else if (item.distance < this.heap[0].distance) {
this.heap[0] = item;
this.sinkDown(0);
}
}
toSortedArray(): KNNResult<T>[] {
return [...this.heap].sort((a, b) => a.distance - b.distance);
}
private bubbleUp(i: number): void {
while (i > 0) {
const parent = (i - 1) >> 1;
if (this.heap[parent].distance >= this.heap[i].distance) break;
[this.heap[parent], this.heap[i]] = [this.heap[i], this.heap[parent]];
i = parent;
}
}
private sinkDown(i: number): void {
const n = this.heap.length;
while (true) {
let largest = i;
const l = 2 * i + 1, r = 2 * i + 2;
if (l < n && this.heap[l].distance > this.heap[largest].distance) largest = l;
if (r < n && this.heap[r].distance > this.heap[largest].distance) largest = r;
if (largest === i) break;
[this.heap[largest], this.heap[i]] = [this.heap[i], this.heap[largest]];
i = largest;
}
}
}
/**
* K-D Tree for efficient nearest-neighbor search over high-dimensional vectors / embeddings.
*
* Supports:
* - Insertion of labeled points
* - Lazy (tombstone) removal, physically purged on rebalance()
* - k-nearest-neighbor (KNN) search
* - Radius search (all points within a given distance)
* - Euclidean and cosine distance metrics
* - Bulk construction (balanced tree) for best query performance
*/
export class KDTree<T = unknown> {
private root: KDNode<T> | null = null;
private _size = 0;
private _tombstones = 0;
private readonly distanceFn: (a: number[], b: number[]) => number;
readonly dims: number;
/**
* @param dims Dimensionality of all vectors (must be consistent).
* @param metric Distance metric to use. Default: "euclidean".
* @param points Optional initial set of points. Builds a balanced tree
* in O(n log² n) — prefer this over inserting one-by-one
* when you have a large corpus.
*/
constructor(
dims: number,
metric: DistanceMetric = "euclidean",
points?: KDPoint<T>[]
) {
this.dims = dims;
this.distanceFn = metric === "cosine" ? cosineDistance : euclideanDistance;
if (points && points.length > 0) {
this.validateAll(points);
this.root = this.buildBalanced([...points], 0);
this._size = points.length;
}
}
/** Total number of live points stored in the tree (excludes tombstoned). */
get size(): number { return this._size; }
/** Fraction of physical nodes that are tombstoned (pending removal on next rebalance). */
get tombstoneRatio(): number {
const total = this._size + this._tombstones;
return total ? this._tombstones / total : 0;
}
// ── Insertion ──────────────────────────────────────────────────────────────
/**
* Insert a single point. O(log n) average, O(n) worst case on skewed data.
* For bulk loading prefer passing points to the constructor.
*/
insert(point: KDPoint<T>): void {
this.validate(point);
this.root = this.insertNode(this.root, point, 0);
this._size++;
}
// ── Removal ────────────────────────────────────────────────────────────────
/**
* Lazily remove all live points whose payload matches `predicate`.
* O(n) traversal, but avoids a full tree rebuild. Call `rebalance()`
* periodically (e.g. once tombstoneRatio crosses ~0.25) to reclaim space
* and restore optimal query depth.
* @returns number of points removed
*/
remove(predicate: (payload: T) => boolean): number {
let removed = 0;
const visit = (node: KDNode<T> | null): void => {
if (!node) return;
if (!node.deleted && predicate(node.point.payload)) {
node.deleted = true;
removed++;
}
visit(node.left);
visit(node.right);
};
visit(this.root);
this._size -= removed;
this._tombstones += removed;
return removed;
}
// ── KNN search ─────────────────────────────────────────────────────────────
/**
* Find the k nearest live neighbors to `query`.
* Returns results sorted by distance ascending.
*/
knn(query: number[], k: number): KNNResult<T>[] {
if (k <= 0) throw new RangeError("k must be a positive integer");
this.validateVector(query);
const heap = new BoundedMaxHeap<T>(k);
this.searchKNN(this.root, query, k, heap, 0);
return heap.toSortedArray();
}
/**
* Nearest single neighbor. Convenience wrapper around knn(query, 1).
* Returns null if the tree is empty.
*/
nearest(query: number[]): KNNResult<T> | null {
const results = this.knn(query, 1);
return results[0] ?? null;
}
// ── Radius search ──────────────────────────────────────────────────────────
/**
* Return all live points whose distance to `query` is ≤ `radius`,
* sorted by distance ascending.
*/
radiusSearch(query: number[], radius: number): KNNResult<T>[] {
if (radius < 0) throw new RangeError("radius must be non-negative");
this.validateVector(query);
const results: KNNResult<T>[] = [];
this.searchRadius(this.root, query, radius, results, 0);
results.sort((a, b) => a.distance - b.distance);
return results;
}
// ── Conversion ─────────────────────────────────────────────────────────────
/** Collect all live points in the tree (order not guaranteed). */
toArray(): KDPoint<T>[] {
const out: KDPoint<T>[] = [];
this.collect(this.root, out);
return out;
}
/**
* Rebuild the tree from its current live points as a balanced tree.
* Physically purges tombstones and restores O(log n) query time.
*/
rebalance(): void {
const points = this.toArray();
this.root = points.length ? this.buildBalanced(points, 0) : null;
this._size = points.length;
this._tombstones = 0;
}
// ── Private: build ─────────────────────────────────────────────────────────
private buildBalanced(points: KDPoint<T>[], depth: number): KDNode<T> {
const axis = depth % this.dims;
points.sort((a, b) => a.vector[axis] - b.vector[axis]);
const mid = Math.floor(points.length / 2);
return {
point: points[mid],
axis,
left: points.slice(0, mid).length
? this.buildBalanced(points.slice(0, mid), depth + 1)
: null,
right: points.slice(mid + 1).length
? this.buildBalanced(points.slice(mid + 1), depth + 1)
: null,
};
}
// ── Private: insert ────────────────────────────────────────────────────────
private insertNode(
node: KDNode<T> | null,
point: KDPoint<T>,
depth: number
): KDNode<T> {
if (node === null) {
return { point, axis: depth % this.dims, left: null, right: null };
}
const axis = depth % this.dims;
if (point.vector[axis] < node.point.vector[axis]) {
node.left = this.insertNode(node.left, point, depth + 1);
} else {
node.right = this.insertNode(node.right, point, depth + 1);
}
return node;
}
// ── Private: KNN traversal ─────────────────────────────────────────────────
private searchKNN(
node: KDNode<T> | null,
query: number[],
k: number,
heap: BoundedMaxHeap<T>,
depth: number
): void {
if (node === null) return;
if (!node.deleted) {
const dist = this.distanceFn(query, node.point.vector);
heap.push({ point: node.point, distance: dist });
}
const axis = node.axis;
const diff = query[axis] - node.point.vector[axis];
const [near, far] = diff <= 0
? [node.left, node.right]
: [node.right, node.left];
this.searchKNN(near, query, k, heap, depth + 1);
const shouldExplore =
this.distanceFn === cosineDistance
? true
: Math.abs(diff) < heap.worstDistance;
if (shouldExplore) {
this.searchKNN(far, query, k, heap, depth + 1);
}
}
// ── Private: radius traversal ──────────────────────────────────────────────
private searchRadius(
node: KDNode<T> | null,
query: number[],
radius: number,
results: KNNResult<T>[],
depth: number
): void {
if (node === null) return;
if (!node.deleted) {
const dist = this.distanceFn(query, node.point.vector);
if (dist <= radius) {
results.push({ point: node.point, distance: dist });
}
}
const axis = node.axis;
const diff = query[axis] - node.point.vector[axis];
const [near, far] = diff <= 0
? [node.left, node.right]
: [node.right, node.left];
this.searchRadius(near, query, radius, results, depth + 1);
const shouldExplore =
this.distanceFn === cosineDistance ? true : Math.abs(diff) <= radius;
if (shouldExplore) {
this.searchRadius(far, query, radius, results, depth + 1);
}
}
// ── Private: collect ───────────────────────────────────────────────────────
private collect(node: KDNode<T> | null, out: KDPoint<T>[]): void {
if (node === null) return;
if (!node.deleted) out.push(node.point);
this.collect(node.left, out);
this.collect(node.right, out);
}
// ── Private: validation ────────────────────────────────────────────────────
private validateVector(v: number[]): void {
if (v.length !== this.dims) {
throw new TypeError(
`Vector length ${v.length} does not match tree dimensionality ${this.dims}`
);
}
}
private validate(point: KDPoint<T>): void {
this.validateVector(point.vector);
}
private validateAll(points: KDPoint<T>[]): void {
for (const p of points) this.validate(p);
}
}
+181
View File
@@ -0,0 +1,181 @@
import {MemoryNode, patchGraph, rebuildGraph} from './graph.ts';
import {KDTree} from './kd-tree.ts';
import type {Memory, MemoryRef, MemoryStore} from './memory.ts';
import {cosineDistance, embedMemoryFields} from '../utils.ts';
const TREE_TOMBSTONE_LIMIT = 0.25;
export function memoryStore(memories: MemoryStore): {
list: Memory[];
cache: MemoryCache | null;
find: (name: string) => Memory | undefined;
ghosts: () => string[];
search: (vector: number[], limit: number) => MemoryRef[];
forget: (name: string) => boolean;
rebuild: (changed?: Memory[]) => MemoryNode[];
backfillEmbeddings: (llm: any) => Promise<number>;
} {
if(memories instanceof MemoryCache) {
return {
list: memories.memories,
cache: memories,
find: name => memories.find(name),
ghosts: () => memories.ghosts(),
search: (vector, limit) => memories.search(vector, limit),
forget: name => memories.remove(name),
rebuild: changed => memories.rebuild(changed),
backfillEmbeddings: llm => memories.backfillEmbeddings(llm),
};
}
return {
list: memories,
cache: null,
find: name => memories.find(m => m.name === name),
ghosts: () => rebuildGraph(memories).filter(n => n.missing).map(n => n.name),
search: (vector, limit) => memories
.filter(m => m.embedding?.length)
.map(m => ({
name: m.name,
description: m.description,
distance: cosineDistance(vector, m.embedding),
}))
.sort((a, b) => a.distance - b.distance)
.slice(0, limit),
forget: name => {
const idx = memories.findIndex(m => m.name === name);
if(idx === -1) return false;
memories.splice(idx, 1);
return true;
},
rebuild: changed => rebuildGraph(memories),
backfillEmbeddings: async llm => {
const missing = memories.filter(m => !m.embedding?.length);
if(!missing.length) return 0;
await Promise.all(missing.map(async node => {
await embedMemoryFields(node, llm);
}));
return missing.length;
},
};
}
export class MemoryCache {
private tree!: KDTree<MemoryRef>;
private indexed = new Map<string, number[]>();
public memories: Memory[];
public nodes: MemoryNode[] = [];
get length() {
return this.memories.length;
}
constructor(memories: Memory[]) {
this.memories = memories;
this.tree = new KDTree<MemoryRef>(0);
this.rebuild();
}
find(name: string): Memory | undefined {
return this.memories.find(m => m.name === name);
}
private syncTree(): void {
const current = new Set(this.memories.map(m => m.name));
for(const [name, emb] of [...this.indexed]) {
const mem = this.memories.find(m => m.name === name);
if(!mem || !current.has(name) || mem.embedding !== emb) {
this.tree.remove(p => p.name === name);
this.indexed.delete(name);
}
}
for(const mem of this.memories) {
if(!mem.embedding?.length || this.indexed.has(mem.name)) continue;
if(this.tree.dims === 0) {
this.tree = new KDTree<MemoryRef>(mem.embedding.length, 'cosine');
}
if(mem.embedding.length !== this.tree.dims) continue;
this.tree.insert({
vector: mem.embedding,
payload: {
name: mem.name,
description: mem.description,
},
});
this.indexed.set(mem.name, mem.embedding);
}
if(this.tree.tombstoneRatio > TREE_TOMBSTONE_LIMIT) {
this.tree.rebalance();
}
}
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,
distance: r.distance,
}));
}
add(memory: Memory): void {
this.memories.push(memory);
this.rebuild([memory]);
}
update(memory: Memory): void {
const existing = this.find(memory.name);
if(existing) Object.assign(existing, memory);
else this.memories.push(memory);
this.rebuild([existing ?? memory]);
}
remove(name: string): boolean {
const idx = this.memories.findIndex(m => m.name === name);
if(idx === -1) return false;
this.memories.splice(idx, 1);
this.rebuild();
return true;
}
ghosts(): string[] {
return this.nodes.filter(n => n.missing).map(n => n.name);
}
rebuild(changed?: Memory[]): MemoryNode[] {
this.nodes = changed?.length && this.nodes.length
? patchGraph(this.memories, this.nodes, changed)
: rebuildGraph(this.memories);
this.syncTree();
return this.nodes;
}
commit(changed?: Memory[]): MemoryNode[] {
return this.rebuild(changed);
}
async backfillEmbeddings(llm: any): Promise<number> {
const missing = this.memories.filter(m => !m.embedding?.length);
if(!missing.length) return 0;
await Promise.all(missing.map(node => embedMemoryFields(node, llm)));
this.commit(missing);
return missing.length;
}
}
+348
View File
@@ -0,0 +1,348 @@
import {AiTool} from '../tools.ts';
import type {LLMMessage, LLMRequest} from '../llm.ts';
import {MemoryCache, memoryStore} from './memory-state.ts';
import {cosineDistance, embedMemoryFields, stripHeader, updateMemory} from '../utils.ts';
const FACT_SIMILARITY_THRESHOLD = 0.62;
const DUPLICATE_THRESHOLD = 0.68;
const PROTECTED_MEMORIES = ['People/User'];
const COLLECTION_WORDS = ['project', 'projects', 'people', 'person', 'managed', 'guides', 'guide', 'research', 'class', 'classes'];
export type Memory = {
name: string;
description: string;
content: string;
embedding: number[];
titleEmbedding?: number[];
bodyEmbeddings?: number[][];
links: string[];
backlinks: string[];
}
export type MemoryRef = {
name: string;
description: string;
distance?: number;
}
export type MemoryOptions = {
memory: Memory[] | MemoryCache;
inject?: boolean;
tool?: boolean;
update?: boolean;
maxTokens?: number;
}
export type MemoryStore = Memory[] | MemoryCache;
/** Create an empty memory shell. */
function emptyNode(name: string, description = ''): Memory {
return {name, description, content: `# ${name.split('/').pop()}\n`, embedding: [], links: [], backlinks: []};
}
function renderNode(node: Memory): string {
return `### ${node.name}
Description: ${node.description}
Links: ${[...node.links, ...node.backlinks].join(', ') || 'none'}
\`\`\`markdown
${node.content}
\`\`\``;
}
function factSimilarity(a: Memory, b: Memory): number {
return !a.bodyEmbeddings?.length || !b.bodyEmbeddings?.length ? 0 : Math.max(...a.bodyEmbeddings.flatMap(av => b.bodyEmbeddings!.map(bv => 1 - cosineDistance(av, bv))));
}
function words(text: string): string[] {
return [...new Set(text.toLowerCase().replace(/[[\]()/_-]/g, ' ').replace(/[^a-z0-9\s]/g, '').split(/\s+/).filter(w => w && !COLLECTION_WORDS.includes(w)))];
}
function jaccard(a: string[], b: string[]): number {
const bs = new Set(b), hit = a.filter(x => bs.has(x)).length, total = new Set([...a, ...b]).size;
return total ? hit / total : 0;
}
function duplicateScore(a: Memory, b: Memory): number {
const name = Math.max(
jaccard(words(a.name), words(b.name)),
jaccard(words(a.name.split('/').pop() || a.name), words(b.name.split('/').pop() || b.name)),
);
const desc = jaccard(words(a.description), words(b.description));
const body = factSimilarity(a, b);
const emb = a.embedding?.length && b.embedding?.length && a.embedding.length === b.embedding.length ? 1 - cosineDistance(a.embedding, b.embedding) : 0;
return Math.max(body, name * 0.9 + desc * 0.06 + emb * 0.04, emb * 0.55 + name * 0.35 + desc * 0.1);
}
function homeScore(node: Memory): number {
return (PROTECTED_MEMORIES.includes(node.name) ? 1e9 : 0)
+ (node.name.includes('/') ? 4 : 0)
+ (node.description && node.description !== 'Persistent memory document' ? 1 : 0)
+ Math.min(stripHeader(node.content).length / 1000, 5);
}
function pickMerge(a: Memory, b: Memory, touched: Set<string>): [drop: Memory, home: Memory] {
const as = homeScore(a), bs = homeScore(b);
if(touched.has(a.name) && !touched.has(b.name)) return as > bs + 2 ? [b, a] : [a, b];
if(touched.has(b.name) && !touched.has(a.name)) return bs > as + 2 ? [a, b] : [b, a];
return as <= bs ? [a, b] : [b, a];
}
/** Build memory tools and memory index text. */
export function memoryTools(llm: any, memories: MemoryStore): {tools: AiTool[]; list: string} {
const store = memoryStore(memories);
const names = new Map<string, string>();
for(const node of store.list)
if(!names.has(node.name)) names.set(node.name, `${node.name} - ${node.description}`);
for(const name of store.ghosts())
if(!names.has(name)) names.set(name, `${name} - ghost node`);
return {
list: [...names.values()].join('\n'),
tools: [
{
name: 'memory_search',
description: 'Semantically search memories for most relevant',
args: {
query: {type: 'string', description: 'Search query', required: true},
limit: {type: 'number', description: 'Maximum results, default 5', default: 5},
},
fn: async ({query, limit = 5}) => {
if(!query?.trim()) return 'Search query is required.';
const [chunk] = await llm.embedding(query, {maxTokens: 8000, overlapTokens: 0});
if(!chunk?.embedding) return 'Failed to create embedding from query';
const results = store.search(chunk.embedding, limit).map(ref => store.find(ref.name)).filter((node): node is Memory => !!node);
return results.length ? results.map(renderNode).join('\n\n---\n\n') : 'No relevant memories found.';
},
},
{
name: 'memory_read',
description: 'Read an entire memory document by name',
args: {name: {type: 'string', description: 'Exact document name', required: true}},
fn: async ({name}) => {
const node = store.find(name);
return node ? renderNode(node) : store.ghosts().includes(name) ? `"${name}" is a ghost node with no document of its own.` : `Not found: "${name}".`;
},
},
{
name: 'memory_delete',
description: 'Delete a duplicate or merged memory',
args: {name: {type: 'string', description: 'Exact document name', required: true}},
fn: async ({name}) => {
store.forget(name);
return `Removed: ${name}`;
},
},
{
name: 'memory_write',
description: 'Create or replace a memory document.',
args: {
name: {type: 'string', description: 'Document name following the entity naming convention.', required: true},
description: {type: 'string', description: 'One factual sentence describing the entire document subject', required: true},
content: {type: 'string', description: 'Complete Markdown document body, including the # title', required: true},
},
fn: async (args: any) => {
const name = String(args.name || '').trim();
if(!name) return 'A document name is required.';
const description = String(args.description || '').trim();
if(!description) return 'A document description is required.';
const content = String(args.content || '').trim();
if(!content) return 'Document content is required.';
let node = store.find(name);
if(!node) {
node = emptyNode(name, description);
if(store.cache) store.cache.add(node);
else store.list.push(node);
}
node.description = name === 'People/User' ? 'All information about the current user' : description.replace(/\s+/g, ' ').trim();
node.content = updateMemory(node, content);
await embedMemoryFields(node, llm);
store.cache?.commit([node]);
return `Updated ${name}`;
},
},
],
};
}
export class MemoryManager {
private memorized = new WeakMap<LLMMessage[], LLMMessage>();
constructor(private llm: any) {}
static normalize(memory?: Memory[] | MemoryCache | MemoryOptions): MemoryOptions | null {
if(!memory) return null;
if(Array.isArray(memory) || memory instanceof MemoryCache) return {memory, inject: true, tool: false, update: false};
if(typeof memory === 'object' && 'memory' in memory) return {inject: true, tool: false, update: false, ...memory};
return null;
}
private memorySystem(list: string): string {
return `You maintain notes written in markdown used for memories from recent conversations using your tools.
Only preserve durable information worth remembering established by the USER.
Do not store assistant guesses, speculation, suggestions, commentary, temporary state, or details that are not worth remembering.
## Rules
- ALWAYS READ a target memory before changing it, \`memory_write\` does a full replace, it DOES NOT append!
- Memories should contain the final state, not deltas
- New conversational context is authoritative when it contracts existing information; reconcile it
- Only remove information when stale, contradicted or duplicated; always preserve existing information, formatting and keep related information together
- Only merge memories when two or more nodes are clearly about the same thing; only split a memory when it is clearly about two distinct subjects
- Use [[WikiLinks]] liberally to record aliases and relationships between entities, even ones without pages yet (ghost nodes)
- Use headings, subheadings, lists, tables and other markdown formatting to make documents clean
- Maintain a \`## Todo List\` of checkboxes AS THE FIRST SUBHEADING when an entity has tasks
- Only create todo items for USER tasks, not AI work
- Only store each in one place, no duplicates
- Use \`People/User\` for personal tasks or as a fallback
## Naming
- Every fact should be grouped with the owning entity
- Always follow the naming convention \`Collection/(Pro)Noun\`
- Facts about the user belong under People/User
- Reuse existing memories when they are clearly the same entity including aliases and ghost references.
- Only create deeper paths when there is a real parent/child entity relationship: \`School/Class/Chapter\`
Valid Examples:
- People/User
- People/John Smith
- Projects/Momentum
- Projects/Momentum/Marketing
- Research/Object Recognition
- Guides/HAM Radio SOP
## Workflow
1. Create groups of durable information and todos based on the owning entity & naming rules above
2. For each group:
1. Read the existing memory(s)
2. Merge the information & todos based on the rules above
3. Write the entire patched document
Available memories:
${list || 'No memory documents exist yet.'}`;
}
private touchedNames(history: LLMMessage[]): string[] {
return [...new Set(history
.filter((h: any) => h.role === 'tool' && h.name === 'memory_write' && !h.error)
.map((h: any) => String(h.args?.name || h.content?.match(/^Updated (.+)$/)?.[1] || '').trim())
.filter(Boolean))];
}
private async backfillEmbeddings(store: ReturnType<typeof memoryStore>): Promise<void> {
const missing = store.list.filter(m => !m.embedding?.length || !m.titleEmbedding?.length || !m.bodyEmbeddings?.length);
await Promise.all(missing.map(m => embedMemoryFields(m, this.llm)));
store.cache?.commit(missing);
}
private closestDuplicate(node: Memory, store: ReturnType<typeof memoryStore>): Memory | null {
return store.list
.filter(m => m.name !== node.name && !m.name.startsWith('Journal/') && !node.name.startsWith('Journal/'))
.map(m => ({node: m, score: duplicateScore(node, m)}))
.filter(x => x.score >= DUPLICATE_THRESHOLD || factSimilarity(node, x.node) >= FACT_SIMILARITY_THRESHOLD)
.sort((a, b) => b.score - a.score)[0]?.node || null;
}
private async rehomeDeleted(drop: Memory, home: Memory, memories: MemoryStore, options: LLMRequest): Promise<void> {
const store = memoryStore(memories);
const backup = structuredClone(drop);
store.forget(drop.name);
try {
const memory = memoryTools(this.llm, memories);
await this.llm.ask(`A duplicate memory document was removed automatically.
Deleted document:
${renderNode(backup)}
Closest surviving home:
${renderNode(home)}
Reinsert every durable unique fact, useful relationship, alias, and user todo from the deleted document into the best remaining memory document.
Usually this should be "${home.name}", but use another existing memory if it is a better home.
Read before writing. Write full replacement documents only.
Do NOT recreate "${backup.name}" unless the deletion was wrong and it is clearly a distinct persistent entity.`, {
model: options.memoryModel || options.model,
temperature: 0.2,
maxTokens: options.maxTokens,
tools: memory.tools,
history: [],
system: this.memorySystem(memory.list),
});
} catch(err) {
if(!store.find(backup.name)) store.cache ? store.cache.add(backup) : store.list.push(backup);
throw err;
} finally {
store.cache?.commit(store.list);
}
}
private async reconcileSimilar(history: LLMMessage[], memories: MemoryStore, options: LLMRequest): Promise<void> {
const store = memoryStore(memories);
const touched = new Set(this.touchedNames(history));
const targets = store.list.filter(m => touched.has(m.name) || [...touched].some(t => duplicateScore(m, store.find(t) || m) >= DUPLICATE_THRESHOLD));
const deleted = new Set<string>();
if(!targets.length) return;
await this.backfillEmbeddings(store);
for(const node of targets) {
if(!store.find(node.name) || deleted.has(node.name) || PROTECTED_MEMORIES.includes(node.name)) continue;
const closest = this.closestDuplicate(node, store);
if(!closest) continue;
const [drop, home] = pickMerge(node, closest, touched);
if(deleted.has(drop.name) || PROTECTED_MEMORIES.includes(drop.name)) continue;
deleted.add(drop.name);
await this.rehomeDeleted(drop, home, memories, options);
await this.backfillEmbeddings(store);
}
}
async recollect(query: string, memory: MemoryStore, limit = 15): Promise<Memory[]> {
const store = memoryStore(memory);
if(!store.list.length || !query?.trim()) return [];
const [chunk] = await this.llm.embedding(query, {maxTokens: 8000, overlapTokens: 0});
return !chunk?.embedding ? [] : store.search(chunk.embedding, limit).map(ref => store.find(ref.name)).filter((m: Memory | undefined): m is Memory => !!m);
}
get tools(): {read: (memory: MemoryStore) => AiTool[]} {
return {read: (memory: MemoryStore) => memoryTools(this.llm, memory).tools};
}
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest = {},): Promise<Memory[]> {
const store = memoryStore(memories);
const previous = this.memorized.get(history);
let start = 0;
if(previous) {
const index = history.indexOf(previous);
if(index >= 0) start = index + 1;
}
const turns = history.slice(start).filter((h: any) => h.role === 'user' || h.role === 'assistant');
const conversation = turns.map((h: any) => `[${h.role}]: ${h.content}`).join('\n\n').trim();
if(!conversation) return store.list;
const memory = memoryTools(this.llm, memories);
const memoryHistory: LLMMessage[] = [];
await this.llm.ask(conversation, {
model: options.memoryModel || options.model,
temperature: 0.2,
maxTokens: options.maxTokens,
tools: memory.tools,
history: memoryHistory,
system: this.memorySystem(memory.list),
});
await this.reconcileSimilar(memoryHistory, memories, options);
const lastTurn = turns.at(-1);
if(lastTurn) this.memorized.set(history, lastTurn);
return store.list;
}
}
+192 -108
View File
@@ -1,84 +1,99 @@
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';
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 {
let client = this.clients.get(token);
if(!client) {
client = new openAI(clean({baseURL: this.host, apiKey: token || undefined}));
this.clients.set(token, client);
}
return client;
}
private toWireContent(content: any): any {
if(!Array.isArray(content)) return content;
return content.map(c => c.type === 'image'
? {type: 'image_url', image_url: {url: `data:${c.mime};base64,${c.data}`}}
: {type: 'text', text: c.text});
}
/** Convert standard history -> OpenAI wire format */
private toWire(history: LLMMessage[], system?: string): any[] {
const wire: any[] = [];
if(system) wire.push({role: 'system', content: system});
for(let i = 0; i < history.length; i++) { for(let i = 0; i < history.length; i++) {
const h = history[i]; const h = history[i];
if(h.role === 'assistant' && h.tool_calls) {
const tools = h.tool_calls.map((tc: any) => ({ if(h.role !== 'tool') {
role: 'tool', wire.push({role: h.role, content: this.toWireContent(h.content)});
id: tc.id, continue;
name: tc.function.name,
args: JSONAttemptParse(tc.function.arguments, {}),
timestamp: h.timestamp
}));
history.splice(i, 1, ...tools);
i += tools.length - 1;
} else if(h.role === 'tool' && h.content) {
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);
i--;
}
if(!history[i]?.timestamp) history[i].timestamp = Date.now();
}
return history;
} }
private fromStandard(history: LLMMessage[]): any[] { const calls: any[] = [];
return history.reduce((result, h) => { const results: any[] = [];
if(h.role === 'tool') {
result.push({ while(i < history.length && history[i].role === 'tool') {
const tool: any = history[i];
calls.push({
id: tool.id,
type: 'function',
function: {
name: tool.name,
arguments: JSON.stringify(tool.args || {})
}
});
results.push({
role: 'tool',
tool_call_id: tool.id,
content: tool.error || tool.content || ''
});
i++;
}
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: calls
refusal: null,
annotations: []
}, {
role: 'tool',
tool_call_id: h.id,
content: h.error || h.content
}); });
} else {
const {timestamp, ...rest} = h; wire.push(...results);
result.push(rest); i--;
}
return result;
}, [] as any[]);
} }
ask(message: string, options: LLMRequest = {}): AbortablePromise<string> { return wire;
}
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) => {
if(options.system) { if(!options.history) options.history = [];
if(options.history?.[0]?.role != 'system') options.history?.splice(0, 0, {role: 'system', content: options.system, timestamp: Date.now()}); const history = options.history;
else options.history[0].content = options.system; if(message) history.push({role: 'user', content: message, timestamp: Date.now()});
}
let history = this.fromStandard([...options.history || [], {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_tokens: options.max_tokens || this.ai.options.llm?.max_tokens || 4096, max_completion_tokens: options.maxTokens ?? this.ai.options.llm?.maxTokens,
temperature: options.temperature || this.ai.options.llm?.temperature || 0.7, temperature: options.temperature ?? this.ai.options.llm?.temperature,
tools: tools.map(t => ({ tools: tools.map(t => ({
type: 'function', type: 'function',
function: { function: {
@@ -86,84 +101,153 @@ export class OpenAi extends LLMProvider {
description: t.description, description: t.description,
parameters: { parameters: {
type: 'object', type: 'object',
properties: t.args ? objectMap(t.args, (key, value) => ({...value, required: undefined})) : {}, properties: t.args
required: t.args ? Object.entries(t.args).filter(t => t[1].required).map(t => t[0]) : [] ? objectMap(t.args, (key, value) => ({...value, required: undefined}))
: {},
required: t.args
? Object.entries(t.args).filter(t => t[1].required).map(t => t[0])
: []
} }
} }
})) }))
}; };
let resp: any, isFirstMessage = true; if(options.schema) {
const schema = convertSchema(options.schema);
requestParams.response_format = {
type: 'json_schema',
json_schema: {name: 'response', strict: true, schema}
};
}
if(options.stream) requestParams.stream_options = {include_usage: true};
try {
let terminal = false;
let iteration = 0;
do { do {
resp = await this.client.chat.completions.create(requestParams).catch(err => { iteration++;
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`; requestParams.messages = this.toWire(history.filter(h => h.role !== 'system'), options.system);
const callStart = Date.now();
const resp: any = await this.tokenPool.run(token =>
this.getClient(token).chat.completions.create(requestParams)
).catch(err => {
err.message += `\n\nMessages:\n${JSON.stringify(requestParams.messages, null, 2)}`;
throw err; throw err;
}); });
let usage: any;
let finishReason: string | undefined;
let msg: any = {content: '', tool_calls: []};
let streamedChars = 0;
if(options.stream) { if(options.stream) {
if(!isFirstMessage) options.stream({text: '\n\n'}); let streamCompleted = false;
else isFirstMessage = false; try {
resp.choices = [{message: {role: 'assistant', content: '', tool_calls: []}}];
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.usage) usage = chunk.usage;
resp.choices[0].message.content += chunk.choices[0].delta.content;
options.stream({text: chunk.choices[0].delta.content}); const choice = chunk.choices?.[0];
if(choice?.finish_reason) finishReason = choice.finish_reason;
if(choice?.delta?.content) {
msg.content += choice.delta.content;
streamedChars += choice.delta.content.length;
options.stream({text: choice.delta.content});
}
if(choice?.delta?.tool_calls) {
for(const deltaTC of choice.delta.tool_calls) {
const index = deltaTC.index ?? msg.tool_calls.length;
let existing = msg.tool_calls.find((tc: any) => tc.index === index);
if(!existing) {
existing = {index, id: '', function: {name: '', arguments: ''}};
msg.tool_calls.push(existing);
} }
if(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);
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 {
resp.choices[0].message.tool_calls.push({
index: deltaTC.index,
id: deltaTC.id || '',
type: deltaTC.type || 'function',
function: {
name: deltaTC.function?.name || '',
arguments: deltaTC.function?.arguments || ''
}
});
}
}
} }
} }
} }
if(resp.error) throw new Error(resp.error); streamCompleted = true;
const toolCalls = resp.choices[0].message.tool_calls || []; } catch(err) {
if(!controller.signal.aborted) throw err;
}
if(streamCompleted && !finishReason) finishReason = msg.tool_calls.length ? 'tool_calls' : 'stop';
} else {
usage = resp.usage;
finishReason = resp.choices[0].finish_reason;
msg = resp.choices[0].message;
}
const duration = Date.now() - callStart;
const tps = usage?.completion_tokens && duration > 0 ? usage.completion_tokens / (duration / 1000) : 0;
if(finishReason === 'length' && !controller.signal.aborted) {
if(msg.content?.trim()) history.push({role: 'assistant', content: msg.content.trim(), timestamp: Date.now(), duration, tps});
throw new Error(`[OpenAI] Response hit token limit before completing`);
}
if(!finishReason && !controller.signal.aborted) {
throw new Error('[OpenAI] Completion ended without a usable response');
}
const toolCalls = msg.tool_calls || [];
if(toolCalls.length && !controller.signal.aborted) { if(toolCalls.length && !controller.signal.aborted) {
history.push(resp.choices[0].message); 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 = {
if(!tool) return {role: 'tool', tool_call_id: toolCall.id, content: '{"error": "Tool not found"}'}; role: 'tool',
id: tc.id,
name: tc.function.name,
args: JSONAttemptParse(tc.function.arguments, {}),
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.function.name));
if(options.stream) options.stream({tool: tc.function.name});
if(!tool) return entry.error = 'Tool not found';
try { try {
const args = JSONAttemptParse(toolCall.function.arguments, {}); const toolStream = options.stream && ((chunk: any) => {
const result = await tool.fn(args, options.stream, this.ai); if(chunk.done) return;
return {role: 'tool', tool_call_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result}; options.stream!(chunk);
} catch (err: any) { });
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'})};
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
} catch(err: any) {
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 (!controller.signal.aborted && resp.choices?.[0]?.message?.tool_calls?.length); } while(!terminal && !controller.signal.aborted);
history.push({role: 'assistant', content: resp.choices[0].message.content.trim() || ''});
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); const turnStart = history.map(h => h.role).lastIndexOf('user');
res(history.at(-1)?.content); const finalContent = history.slice(turnStart + 1).reduce((str, h) => h.role === 'assistant' ? str + (h.content || '') : str, '').trim();
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
} catch(err) {
rej(err);
}
}), {abort: () => controller.abort()}); }), {abort: () => controller.abort()});
} }
} }
+1 -1
View File
@@ -1,5 +1,5 @@
import {AbortablePromise} from './ai.ts'; import {AbortablePromise} from './ai.ts';
import {LLMMessage, LLMRequest} from './llm.ts'; import {LLMRequest} from './llm.ts';
export abstract class LLMProvider { export abstract class LLMProvider {
abstract ask(message: string, options: LLMRequest): AbortablePromise<string>; abstract ask(message: string, options: LLMRequest): AbortablePromise<string>;
+65
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);
}
}
+596 -128
View File
@@ -1,6 +1,6 @@
import * as cheerio from 'cheerio'; import * as cheerio from 'cheerio';
import {$Sync} from '@ztimson/node-utils'; import {$Sync} from '@ztimson/node-utils';
import {ASet, consoleInterceptor, Http, fn as Fn, decodeHtml} from '@ztimson/utils'; import {ASet, consoleInterceptor, Http, fn as Fn, decodeHtml, objectMap} from '@ztimson/utils';
import * as os from 'node:os'; import * as os from 'node:os';
import {Ai} from './ai.ts'; import {Ai} from './ai.ts';
import {LLMRequest} from './llm.ts'; import {LLMRequest} from './llm.ts';
@@ -41,28 +41,83 @@ export type AiTool = {
/** Tool arguments */ /** Tool arguments */
args?: AiToolArg, args?: AiToolArg,
/** Callback function */ /** Callback function */
fn: (args: any, stream: LLMRequest['stream'], ai: Ai) => any | Promise<any>, fn: (args: any, stream: LLMRequest['stream'], ai: Ai, toolId?: string) => any | Promise<any>,
}; };
export const CliTool: AiTool = { export function convertSchema(schema: any): any {
if(!schema) return null;
const convertProp = (prop: any): any => {
const converted: any = {
type: prop.type || 'string',
};
if(prop.description) converted.description = prop.description;
if(prop.default !== undefined) converted.default = prop.default;
if(prop.enum) converted.enum = prop.enum;
if(prop.pattern) converted.pattern = prop.pattern;
// Handle array items
if(prop.type === 'array' && prop.items) {
converted.items = convertProp(prop.items);
}
// Handle object properties
if(prop.type === 'object' && prop.items) {
converted.properties = objectMap(prop.items, (key, value) => convertProp(value));
const required = Object.entries(prop.items).filter(([_, v]: any) => v.required).map(([k]) => k);
if(required.length) converted.required = required;
converted.additionalProperties = false;
}
// Handle min/max based on type
if(prop.min !== undefined) {
if(prop.type === 'string' || prop.type === 'array') converted.minLength = prop.min;
else converted.minimum = prop.min;
}
if(prop.max !== undefined) {
if(prop.type === 'string' || prop.type === 'array') converted.maxLength = prop.max;
else converted.maximum = prop.max;
}
return converted;
};
return {
type: 'object',
properties: objectMap(schema, (key, value) => convertProp(value)),
required: Object.entries(schema).filter(([_, v]: any) => v.required).map(([k]) => k),
additionalProperties: false
};
}
export const ExecCliTool: AiTool = {
name: 'cli', name: 'cli',
description: 'Use the command line interface, returns any output', description: 'Use the command line interface, returns any output',
args: {command: {type: 'string', description: 'Command to run', required: true}}, args: {command: {type: 'string', description: 'Command to run', required: true}},
fn: (args: {command: string}) => $Sync`${args.command}` fn: (args: {command: string}) => $Sync`${args.command}`
} }
export const DateTimeTool: AiTool = { export const ExecJSTool: AiTool = {
name: 'get_datetime', name: 'exec_javascript',
description: 'Get local date / time', description: 'Execute commonjs javascript',
args: {}, args: {
fn: async () => new Date().toString() code: {type: 'string', description: 'CommonJS javascript', required: true}
},
fn: async (args: {code: string}) => {
const c = consoleInterceptor(null);
const resp = await Fn<any>({console: c}, args.code, true).catch((err: any) => c.output.error.push(err));
return {...c.output, return: resp, stdout: undefined, stderr: undefined};
}
} }
export const DateTimeUTCTool: AiTool = { export const ExecPythonTool: AiTool = {
name: 'get_datetime_utc', name: 'exec_python',
description: 'Get current UTC date / time', description: 'Execute commonjs javascript',
args: {}, args: {
fn: async () => new Date().toUTCString() code: {type: 'string', description: 'CommonJS javascript', required: true}
},
fn: async (args: {code: string}) => ({result: $Sync`python -c "${args.code}"`})
} }
export const ExecTool: AiTool = { export const ExecTool: AiTool = {
@@ -76,11 +131,11 @@ export const ExecTool: AiTool = {
try { try {
switch(args.language) { switch(args.language) {
case 'cli': case 'cli':
return await CliTool.fn({command: args.code}, stream, ai); return await ExecCliTool.fn({command: args.code}, stream, ai);
case 'node': case 'node':
return await JSTool.fn({code: args.code}, stream, ai); return await ExecJSTool.fn({code: args.code}, stream, ai);
case 'python': case 'python':
return await PythonTool.fn({code: args.code}, stream, ai); return await ExecPythonTool.fn({code: args.code}, stream, ai);
default: default:
throw new Error(`Unsupported language: ${args.language}`); throw new Error(`Unsupported language: ${args.language}`);
} }
@@ -90,8 +145,483 @@ export const ExecTool: AiTool = {
} }
} }
export const FetchTool: AiTool = { export const FsDeleteTool = (whitelist: null | string[] = null): AiTool => {
name: 'fetch', return {
name: 'fs_delete',
description: 'Delete a file or directory',
args: {
path: {type: 'string', description: 'Path to file or directory', required: true},
recursive: {type: 'boolean', description: 'Delete all children', required: false}
},
fn: async ({path, recursive = false}) => {
const {existsSync, rmSync} = await import('fs');
const normalizePath = p => p.replace(/\\/g, '/');
path = normalizePath(path);
if(whitelist && !whitelist.some(p => path.startsWith(p))) return {error: 'Permission denied'};
if(!existsSync(path)) return {error: 'Path does not exist'};
rmSync(path, {recursive, force: true});
return {success: true, path};
}
}
}
export const FsMoveTool = (whitelist: null | string[] = null): AiTool => {
return {
name: 'fs_move',
description: 'Move or rename a file or directory',
args: {
source: {type: 'string', description: 'Path to source file or directory', required: true},
destination: {type: 'string', description: 'Path to destination file or directory', required: true}
},
fn: async ({source, destination}) => {
const {existsSync, renameSync} = await import('fs');
const normalizePath = p => p.replace(/\\/g, '/');
source = normalizePath(source);
destination = normalizePath(destination);
if(whitelist && !whitelist.some(p => source.startsWith(p) && destination.startsWith(p))) return {error: 'Permission denied'};
if(!existsSync(source)) return {error: 'Source path does not exist'};
if(existsSync(destination)) return {error: 'Destination path already exists'};
renameSync(source, destination);
return {success: true, source, destination};
}
}
}
export const FsReadTool = (whitelist: null | string[] = null): AiTool => {
return {
name: 'fs_read',
description: 'Read the contents of a provided path. Works with files and directories',
args: {path: {type: 'string', description: 'Path to file or directory', required: true}},
fn: async ({path}) => {
const {existsSync, lstatSync, readdirSync, readFileSync} = await import('fs');
const {join} = await import('path');
const normalizePath = p => p.replace(/\\/g, '/');
path = normalizePath(path);
if(whitelist && !whitelist.some(p => path.startsWith(p))) return {error: 'Permission denied'};
if(!existsSync(path)) return {error: 'Path does not exist'};
const stats = lstatSync(path);
if(stats.isDirectory()) {
const children = readdirSync(path).map(name => {
const childPath = normalizePath(join(path, name));
const childStats = lstatSync(childPath);
return {name, type: childStats.isDirectory() ? 'directory' : 'file', size: childStats.size};
});
return {type: 'directory', children};
}
const content = readFileSync(path, 'utf-8');
return {type: 'file', content};
}
}
}
export const FsSearchTool = (whitelist: null | string[] = null): AiTool => {
return {
name: 'fs_search',
description: 'Scan a directory for matching glob patterns (e.g. "**/*.js", "src/**/*.test.ts")',
args: {
pattern: {type: 'string', description: 'Glob pattern to match against paths', required: true},
root: {type: 'string', description: 'Directory to search from', required: false, default: '.'}
},
fn: async ({pattern, root = '.'}) => {
const {existsSync, lstatSync, readdirSync} = await import('fs');
const {join, relative} = await import('path');
const normalizePath = p => p.replace(/\\/g, '/');
root = normalizePath(root);
if(!existsSync(root)) return {error: 'Root path does not exist'};
if(!lstatSync(root).isDirectory()) return {error: 'Root path is not a directory'};
if(whitelist && !whitelist.some(p => root.startsWith(p))) return {error: 'Permission denied'};
const globToRegex = (glob) => {
let re = '';
for(let i = 0; i < glob.length; i++) {
const c = glob[i];
if(c === '*') {
if(glob[i + 1] === '*') {
const isSlash = glob[i + 2] === '/';
re += '.*';
i += isSlash ? 2 : 1;
} else {
re += '[^/]*';
}
} else if(c === '?') {
re += '[^/]';
} else if('.+^$(){}|[]\\'.includes(c)) {
re += '\\' + c;
} else {
re += c;
}
}
return new RegExp('^' + re + '$');
};
const regex = globToRegex(pattern);
const results: any = [];
const walk = (dir) => {
for(const name of readdirSync(dir)) {
const fullPath = normalizePath(join(dir, name));
const stats = lstatSync(fullPath);
const relPath = normalizePath(relative(root, fullPath));
if(regex.test(relPath)) {
results.push({path: relPath, type: stats.isDirectory() ? 'directory' : 'file', size: stats.size});
}
if(stats.isDirectory()) walk(fullPath);
}
};
walk(root);
return results;
}
}
}
export const FsWriteTool = (whitelist: null | string[] = null): AiTool => {
return {
name: 'fs_write',
description: 'Create a directory, write content to a file or preform a find & replace',
args: {
path: {type: 'string', description: 'Path to file or directory', required: true},
content: {type: 'string', description: 'Content to write or replace (Omit to create a directory)'},
find: {type: 'string', description: 'Text or regex pattern to match (regex must match pattern: "/pattern/g")'}
},
fn: async ({path, content, find}) => {
const {existsSync, mkdirSync, readFileSync, writeFileSync} = await import('fs');
const {dirname} = await import('path');
const normalizePath = p => p.replace(/\\/g, '/');
path = normalizePath(path);
if(whitelist && !whitelist.some(p => path.startsWith(p))) return {error: 'Permission denied'};
if(content === undefined) {
mkdirSync(path, {recursive: true});
return {success: true, type: 'directory', path};
}
const dir = normalizePath(dirname(path));
if(!existsSync(dir)) mkdirSync(dir, {recursive: true});
if(find && existsSync(path)) {
const existing = readFileSync(path, 'utf-8');
const regexMatch = find.match(/^\/(.+)\/([gimuy]*)$/);
const pattern = regexMatch ? new RegExp(regexMatch[1], regexMatch[2]) : find;
if(!existing.match(pattern)) return {error: 'Find pattern not found in file'};
const updated = existing.replace(pattern, content);
writeFileSync(path, updated, 'utf-8');
return {success: true, type: 'file', path, replaced: true, content: updated};
}
writeFileSync(path, content, 'utf-8');
return {success: true, type: 'file', path, content};
}
}
}
export const GetPathsTool: AiTool = {
name: 'get_paths',
description: 'Get the current working directory, and paths to the users home directory',
fn: async () => {
return {
home: os.homedir(),
cwd: process.cwd()
};
}
}
export const GetDatetimeTool: AiTool = {
name: 'get_datetime',
description: 'Get local/UTC timestamp',
args: {
timezone: {type: 'string', description: 'Which timezone to return, defaults to local', enum: ['local', 'utc'], default: 'local'}
},
fn: ({timezone}) => new Date()[timezone === 'local' ? 'toString' : 'toUTCString']()
}
export const GetDevice: AiTool = {
name: 'get_device',
description: 'Get comprehensive system information including hostname, specs, load, storage, and network status',
args: {},
fn: async () => {
const platform = os.platform();
const hostname = os.hostname();
// CPU Info
const cpus = os.cpus();
const cpuModel = cpus[0].model;
const cpuCores = cpus.length;
// Memory Info
const totalMem: any = (os.totalmem() / 1024 / 1024 / 1024).toFixed(2);
const freeMem: any = (os.freemem() / 1024 / 1024 / 1024).toFixed(2);
const usedMem: any = (totalMem - freeMem).toFixed(2);
const memUsage: any = ((usedMem / totalMem) * 100).toFixed(1);
// Load Average (not available on Windows)
const loadAvg = platform === 'win32' ? ['N/A', 'N/A', 'N/A'] : os.loadavg().map(l => l.toFixed(2));
// Storage Usage
let storage = {};
if(platform === 'win32') {
const ps = $Sync`powershell "Get-PSDrive C | Select-Object Used,Free | ConvertTo-Json"`.trim();
const drive = JSON.parse(ps);
const used: any = (drive.Used / 1024 / 1024 / 1024).toFixed(2);
const free: any = (drive.Free / 1024 / 1024 / 1024).toFixed(2);
const total: any = (parseFloat(used) + parseFloat(free)).toFixed(2);
const usage: any = ((used / total) * 100).toFixed(1);
storage = {
filesystem: 'C:',
size: `${total} GB`,
used: `${used} GB`,
available: `${free} GB`,
usage: `${usage}%`
};
} else {
const df = $Sync`df -h / | tail -1`.trim();
const s = df.split(/\s+/);
storage = {
filesystem: s[0],
size: s[1],
used: s[2],
available: s[3],
usage: s[4]
};
}
// Network Status
const interfaces = os.networkInterfaces();
const activeIfaces = Object.entries(interfaces)
.filter(([name]) => name !== 'lo' && !name.includes('Loopback'))
.map(([name, addrs]) => {
const ipv4 = addrs?.find(a => a.family === 'IPv4');
return ipv4 ? {name, ip: ipv4.address} : null;
})
.filter(Boolean);
// Internet connectivity check
let internet = false;
try {
if(platform === 'win32') {
$Sync`powershell "Test-Connection -ComputerName 8.8.8.8 -Count 1 -Quiet"`;
} else {
$Sync`ping -c 1 -W 2 8.8.8.8 > /dev/null 2>&1`;
}
internet = true;
} catch {}
// Uptime
const uptime = os.uptime();
const days = Math.floor(uptime / 86400);
const hours = Math.floor((uptime % 86400) / 3600);
const minutes = Math.floor((uptime % 3600) / 60);
return {
hostname,
cpu: {
model: cpuModel,
cores: cpuCores
},
memory: {
total: `${totalMem} GB`,
used: `${usedMem} GB`,
free: `${freeMem} GB`,
usage: `${memUsage}%`
},
load: {
'1min': loadAvg[0],
'5min': loadAvg[1],
'15min': loadAvg[2]
},
storage,
network: {
interfaces: activeIfaces,
internet: internet ? 'connected' : 'disconnected'
},
uptime: `${days}d ${hours}h ${minutes}m`,
platform: `${os.type()} ${os.release()}`
};
}
}
export const GetWikipediaTool: AiTool = {
name: 'get_wikipedia',
description: 'Search Wikipedia for matching articles',
args: {
query: {type: 'string', description: 'Search term or article title', required: true},
mode: {type: 'string', description: 'search - look for articles, summary - intro of first found article (default), full - complete first found article', enum: ['search', 'summary', 'full'], default: 'summary'},
ua: {type: 'string', description: 'User Agent'},
},
fn: async ({query, mode, ua}) => {
class WikipediaClient {
useragent = 'Mozilla/5.0 (Windows NT 10.0; Win64; x64)';
constructor(useragent: string) {
this.useragent = useragent;
}
async get(url) {
const resp = await fetch(url, {headers: {'User-Agent': this.useragent}});
return resp.json();
}
api(params) {
const qs = new URLSearchParams({...params, format: 'json', utf8: '1'}).toString();
return this.get(`https://en.wikipedia.org/w/api.php?${qs}`);
}
clean(text) {
const cutoffs = ['== See also ==', '== References ==', '== Bibliography ==', '== External links =='];
for (const marker of cutoffs) {
const idx = text.indexOf(marker);
if (idx !== -1) text = text.slice(0, idx);
}
return text
.replace(/^={4}\s*(.+?)\s*={4}$/gm, '#### $1')
.replace(/^={3}\s*(.+?)\s*={3}$/gm, '### $1')
.replace(/^={2}\s*(.+?)\s*={2}$/gm, '## $1')
.replace(/\n{3,}/g, '\n\n')
.replace(/ {2,}/g, ' ')
.replace(/\[\d+]/g, '')
.trim();
}
async searchTitles(query: string, limit = 6) {
const data = await this.api({action: 'query', list: 'search', srsearch: query, srlimit: limit, srprop: 'snippet'});
return data.query?.search || [];
}
async fetchExtract(title: string, introOnly = false) {
const params: any = {action: 'query', prop: 'extracts', titles: title, explaintext: 1, redirects: 1};
if(introOnly) params.exintro = 1;
const data = await this.api(params);
const page: any = Object.values(data.query?.pages || {})[0];
return this.clean(page?.extract || '');
}
pageUrl(title: string) {
return `https://en.wikipedia.org/wiki/${encodeURIComponent(title.replace(/ /g, '_'))}`;
}
stripHtml(text: string) {
return text.replace(/<[^>]+>/g, '');
}
async lookup(query: string, detail = 'summary') {
const results = await this.searchTitles(query, 6);
if(!results.length) return `❌ No Wikipedia articles found for "${query}"`;
const title = results[0].title;
const url = this.pageUrl(title);
const introOnly = detail !== 'full';
const content = await this.fetchExtract(title, introOnly);
return `## ${title}\n🔗 ${url}\n\n${content}`;
}
async search(query: string) {
const results = await this.searchTitles(query, 8);
if(!results.length) return `❌ No results for "${query}"`;
const lines = [`### Search results for "${query}"\n`];
for(let i = 0; i < results.length; i++) {
const r = results[i];
const snippet = this.stripHtml(r.snippet || '').trim();
lines.push(`**${i + 1}. ${r.title}**\n${snippet}\n${this.pageUrl(r.title)}`);
}
return lines.join('\n\n');
}
}
const wiki = new WikipediaClient(ua);
if(mode === 'search') return wiki.search(query);
return wiki.lookup(query, mode || 'summary');
}
};
export const GeoCodeTool: AiTool = {
name: 'geo_code',
description: 'Converts coordinates to address OR vice versa',
args: {
query: {type: 'string', description: 'Search query - coordinates (lat,lon) or address string', required: true},
},
fn: async ({query}) => {
const coordinates = /(-?\d+(?:\.\d+)?).*?,.*?(-?\d+(?:\.\d+)?)/.exec(query);
if(coordinates) { // Geolocate
const url = `https://nominatim.openstreetmap.org/reverse?format=json&lat=${encodeURIComponent(coordinates[1])}&lon=${encodeURIComponent(coordinates[2])}`;
const response = await fetch(url, {headers: {'User-Agent': 'OpenSight/1.0', 'Accept-Language': 'en'}});
const data = await response.json();
if(data.display_name) return {address: data.display_name, mode: 'geolocate'};
} else { // Geocode
const url = `https://nominatim.openstreetmap.org/search?format=json&q=${encodeURIComponent(query)}`;
const response = await fetch(url, {headers: {'User-Agent': 'OpenSight/1.0'}});
const data = await response.json();
if(data[0]) return {latitude: parseFloat(data[0].lat), longitude: parseFloat(data[0].lon), mode: 'geocode'};
}
return {error: 'Not found'};
},
}
export const GeoWeatherTool: AiTool = {
name: 'geo_weather',
description: 'Gets weather and air quality info for a location and time',
args: {
query: {type: 'string', description: 'Location - address or place name', required: true},
day: {type: 'string', description: 'Date to retrieve (YYYY-MM-DD), defaults to today'},
},
fn: async ({query, day}) => {
day = day || new Date().toISOString().slice(0, 10);
const geoUrl = `https://nominatim.openstreetmap.org/search?format=json&q=${encodeURIComponent(query)}`;
const geoResponse = await fetch(geoUrl, {headers: {'User-Agent': 'OpenSight/1.0'}});
const geoData = await geoResponse.json();
if(!geoData[0]) return {error: 'Location not found'};
const lat = parseFloat(geoData[0].lat);
const lon = parseFloat(geoData[0].lon);
const weatherUrl = `https://api.open-meteo.com/v1/forecast?latitude=${lat}&longitude=${lon}&start_date=${day}&end_date=${day}&daily=weathercode,temperature_2m_max,temperature_2m_min,apparent_temperature_max,apparent_temperature_min,precipitation_sum,precipitation_probability_max,windspeed_10m_max,winddirection_10m_dominant,uv_index_max,sunrise,sunset&timezone=auto`;
const airUrl = `https://air-quality-api.open-meteo.com/v1/air-quality?latitude=${lat}&longitude=${lon}&start_date=${day}&end_date=${day}&hourly=us_aqi,european_aqi,pm10,pm2_5&timezone=auto`;
const [weatherResponse, airResponse] = await Promise.all([fetch(weatherUrl), fetch(airUrl)]);
const weatherData = await weatherResponse.json();
const airData = await airResponse.json();
const avg = arr => (arr && arr.length) ? arr.reduce((a, b) => a + b, 0) / arr.length : null;
return {
location: geoData[0].display_name,
latitude: lat,
longitude: lon,
elevation: weatherData.elevation,
date: day,
weatherCode: weatherData.daily?.weathercode?.[0],
tempMax: weatherData.daily?.temperature_2m_max?.[0],
tempMin: weatherData.daily?.temperature_2m_min?.[0],
feelsLikeMax: weatherData.daily?.apparent_temperature_max?.[0],
feelsLikeMin: weatherData.daily?.apparent_temperature_min?.[0],
precipitation: weatherData.daily?.precipitation_sum?.[0],
precipitationChance: weatherData.daily?.precipitation_probability_max?.[0],
windSpeedMax: weatherData.daily?.windspeed_10m_max?.[0],
windDirection: weatherData.daily?.winddirection_10m_dominant?.[0],
uvIndexMax: weatherData.daily?.uv_index_max?.[0],
sunrise: weatherData.daily?.sunrise?.[0],
sunset: weatherData.daily?.sunset?.[0],
usAqi: avg(airData.hourly?.us_aqi),
europeanAqi: avg(airData.hourly?.european_aqi),
pm10: avg(airData.hourly?.pm10),
pm2_5: avg(airData.hourly?.pm2_5),
};
},
}
export const WebFetchTool: AiTool = {
name: 'web_fetch',
description: 'Make HTTP request to URL', description: 'Make HTTP request to URL',
args: { args: {
url: {type: 'string', description: 'URL to fetch', required: true}, url: {type: 'string', description: 'URL to fetch', required: true},
@@ -107,30 +637,59 @@ export const FetchTool: AiTool = {
}) => new Http({url: args.url, headers: args.headers}).request({method: args.method || 'GET', body: args.body}) }) => new Http({url: args.url, headers: args.headers}).request({method: args.method || 'GET', body: args.body})
} }
export const JSTool: AiTool = { export const WebFlareSolverTool = (host: string) => {
name: 'exec_javascript', return {
description: 'Execute commonjs javascript', name: 'web_flaresolverr',
description: 'Use a flaresolverr proxy to bypass cloudflare bot detection',
args: { args: {
code: {type: 'string', description: 'CommonJS javascript', required: true} url: {type: 'string', description: 'URL to fetch', required: true},
cmd: {type: 'string', description: 'Flaresolverr cmd', enum: ['request.get', 'request.post'], default: 'request.get'},
maxTimeout: {type: 'number', description: 'Fetch time limit', default: 60_000},
postData: {type: 'object', description: 'Data to send during request.post requests'},
}, },
fn: async (args: {code: string}) => { fn: async ({url, cmd, maxTimeout, postData}) => {
const c = consoleInterceptor(null); function toFormUrlEncoded(obj, prefix = '') {
const resp = await Fn<any>({console: c}, args.code, true).catch((err: any) => c.output.error.push(err)); const pairs: any = [];
return {...c.output, return: resp, stdout: undefined, stderr: undefined}; for (const key in obj) {
if (!obj.hasOwnProperty(key)) continue;
const value = obj[key];
const encodedKey = prefix
? `${prefix}[${encodeURIComponent(key)}]`
: encodeURIComponent(key);
if (value === null || value === undefined) {
pairs.push(`${encodedKey}=`);
} else if (typeof value === 'object' && !Array.isArray(value)) {
pairs.push(toFormUrlEncoded(value, encodedKey));
} else if (Array.isArray(value)) {
value.forEach(item => {
pairs.push(`${encodedKey}[]=${encodeURIComponent(item)}`);
});
} else {
pairs.push(`${encodedKey}=${encodeURIComponent(value)}`);
}
}
return pairs.join('&');
}
const res = await fetch(host + '/v1', {
method: 'POST',
headers: {'Content-Type': 'application/json'},
body: JSON.stringify({cmd, url, maxTimeout, postData: postData ? toFormUrlEncoded(postData) : undefined}),
});
if(!res.ok) throw new Error(`FlareSolverr HTTP error: ${res.status} ${res.statusText}`);
const data = await res.json();
if(data.status !== 'ok') throw new Error(`FlareSolverr error: ${data.message ?? data.status}`);
return data.solution.response;
}
} }
} }
export const PythonTool: AiTool = { export const WebReadTool: AiTool = {
name: 'exec_javascript', name: 'web_read',
description: 'Execute commonjs javascript',
args: {
code: {type: 'string', description: 'CommonJS javascript', required: true}
},
fn: async (args: {code: string}) => ({result: $Sync`python -c "${args.code}"`})
}
export const ReadWebpageTool: AiTool = {
name: 'read_webpage',
description: 'Extract clean content from webpages, or convert media/documents to accessible formats', description: 'Extract clean content from webpages, or convert media/documents to accessible formats',
args: { args: {
url: {type: 'string', description: 'URL to read', required: true}, url: {type: 'string', description: 'URL to read', required: true},
@@ -258,94 +817,3 @@ export const WebSearchTool: AiTool = {
return results; return results;
} }
} }
class WikipediaClient {
private async get(url: string): Promise<any> {
const resp = await fetch(url, {headers: {'User-Agent': UA}});
return resp.json();
}
private api(params: Record<string, any>): Promise<any> {
const qs = new URLSearchParams({...params, format: 'json', utf8: '1'}).toString();
return this.get(`https://en.wikipedia.org/w/api.php?${qs}`);
}
private clean(text: string): string {
return text.replace(/\n{3,}/g, '\n\n').replace(/ {2,}/g, ' ').replace(/\[\d+\]/g, '').trim();
}
private truncate(text: string, max: number): string {
if(text.length <= max) return text;
const cut = text.slice(0, max);
const lastPara = cut.lastIndexOf('\n\n');
return lastPara > max * 0.7 ? cut.slice(0, lastPara) : cut;
}
private async searchTitles(query: string, limit = 6): Promise<any[]> {
const data = await this.api({action: 'query', list: 'search', srsearch: query, srlimit: limit, srprop: 'snippet'});
return data.query?.search || [];
}
private async fetchExtract(title: string, intro = false): Promise<string> {
const params: any = {action: 'query', prop: 'extracts', titles: title, explaintext: 1, redirects: 1};
if(intro) params.exintro = 1;
const data = await this.api(params);
const page = Object.values(data.query?.pages || {})[0] as any;
return this.clean(page?.extract || '');
}
private pageUrl(title: string): string {
return `https://en.wikipedia.org/wiki/${encodeURIComponent(title.replace(/ /g, '_'))}`;
}
private stripHtml(text: string): string {
return text.replace(/<[^>]+>/g, '');
}
async lookup(query: string, detail: 'intro' | 'full' = 'intro'): Promise<string> {
const results = await this.searchTitles(query, 6);
if(!results.length) return `❌ No Wikipedia articles found for "${query}"`;
const title = results[0].title;
const url = this.pageUrl(title);
const content = await this.fetchExtract(title, detail === 'intro');
const text = this.truncate(content, detail === 'intro' ? 2000 : 8000);
return `## ${title}\n🔗 ${url}\n\n${text}`;
}
async search(query: string): Promise<string> {
const results = await this.searchTitles(query, 8);
if(!results.length) return `❌ No results for "${query}"`;
const lines = [`### Search results for "${query}"\n`];
for(let i = 0; i < results.length; i++) {
const r = results[i];
const snippet = this.truncate(this.stripHtml(r.snippet || ''), 150);
lines.push(`**${i + 1}. ${r.title}**\n${snippet}\n${this.pageUrl(r.title)}`);
}
return lines.join('\n\n');
}
}
export const WikipediaLookupTool: AiTool = {
name: 'wikipedia_lookup',
description: 'Get Wikipedia article content',
args: {
query: {type: 'string', description: 'Topic or article title', required: true},
detail: {type: 'string', description: 'Content level: "intro" (summary, default) or "full" (complete article)', enum: ['intro', 'full'], default: 'intro'}
},
fn: async (args: {query: string; detail?: 'intro' | 'full'}) => {
const wiki = new WikipediaClient();
return wiki.lookup(args.query, args.detail || 'intro');
}
};
export const WikipediaSearchTool: AiTool = {
name: 'wikipedia_search',
description: 'Search Wikipedia for matching articles',
args: {
query: {type: 'string', description: 'Search terms', required: true}
},
fn: async (args: {query: string}) => {
const wiki = new WikipediaClient();
return wiki.search(args.query);
}
};
+82
View File
@@ -0,0 +1,82 @@
import {Memory} from './memory/memory.ts';
export function cosineDistance(a: number[], b: number[]): number {
let dot = 0, normA = 0, normB = 0;
for(let i = 0; i < a.length; i++) {
dot += a[i] * b[i];
normA += a[i] * a[i];
normB += b[i] * b[i];
}
const denom = Math.sqrt(normA) * Math.sqrt(normB);
return denom === 0 ? 1 : 1 - dot / denom;
}
export async function embedMemoryFields(node: Memory, llm: any): Promise<void> {
const body = stripHeader(node.content);
const [titleE] = await llm.embedding(node.name.split('/').pop() || node.name);
const [descE] = await llm.embedding(node.description || '');
const bodyChunks = body ? await llm.embedding(body) : [];
if(titleE) node.titleEmbedding = titleE.embedding;
if(descE) node.embedding = descE.embedding;
node.bodyEmbeddings = bodyChunks.map((c: any) => c.embedding).filter(Boolean);
}
export function euclideanDistance(a: number[], b: number[]): number {
let sum = 0;
for(let i = 0; i < a.length; i++) {
const d = a[i] - b[i];
sum += d * d;
}
return Math.sqrt(sum);
}
export function getWeekStart(date: Date = new Date()): string {
const d = new Date(Date.UTC(date.getFullYear(), date.getMonth(), date.getDate()));
const day = d.getUTCDay();
const diff = day === 0 ? -6 : 1 - day;
d.setUTCDate(d.getUTCDate() + diff);
return d.toISOString().slice(0, 10);
}
export function journalDescription(journalName?: string): string {
const start = journalName?.split('/').pop() || getWeekStart();
const d = new Date(`${start}T00:00:00Z`);
d.setUTCDate(d.getUTCDate() + 6);
const end = d.toISOString().slice(0, 10);
return `Log from ${start} - ${end}`;
}
function parseFrontmatter(content: string): {fm: Map<string, string>, body: string} {
const match = content.match(/^---\n([\s\S]*?)\n---\n?([\s\S]*)$/);
if(!match) return {fm: new Map(), body: content};
const fm = new Map<string, string>();
for(const line of match[1].split('\n')) {
const i = line.indexOf(':');
if(i === -1) continue;
const key = line.slice(0, i).trim();
const raw = line.slice(i + 1).trim();
let value = raw;
try { value = JSON.parse(raw); } catch { }
fm.set(key, value);
}
return {fm, body: match[2]};
}
export function writeFrontmatter(fm: Map<string, string>, body: string): string {
const lines = [...fm.entries()].map(([k, v]) =>
`${k}: ${JSON.stringify(String(v).replace(/\s+/g, ' ').trim())}`);
return `---\n${lines.join('\n')}\n---\n\n${body.trimStart()}`;
}
export function stripHeader(content: string): string {
return content.replace(/^---[\s\S]*?\n---\n?/, '').trimStart();
}
export function updateMemory(node: Memory, body: string): string {
const {fm} = parseFrontmatter(node.content);
fm.set('name', node.name);
fm.set('description', (node.name.startsWith('Journal/') ? journalDescription(node.name) : node.description)
|| 'Persistent memory document');
fm.set('modified', new Date().toISOString());
return writeFrontmatter(fm, stripHeader(body));
}
+24 -5
View File
@@ -12,12 +12,31 @@ export class Vision {
*/ */
ocr(path: string): AbortablePromise<string | null> { ocr(path: string): AbortablePromise<string | null> {
let worker: any; let worker: any;
const p = new Promise<string | null>(async res => { let reject: (err: any) => void;
const handler = (err: Error) => {
if(err.stack?.includes('tesseract.js')) {
process.off('uncaughtException', handler);
reject?.(err);
return;
}
throw err;
};
process.on('uncaughtException', handler);
const p = (async () => {
worker = await createWorker(this.ai.options.ocr || 'eng', 2, {cachePath: this.ai.options.path}); worker = await createWorker(this.ai.options.ocr || 'eng', 2, {cachePath: this.ai.options.path});
const {data} = await worker.recognize(path); return await new Promise<string | null>((res, rej) => {
await worker.terminate(); reject = rej;
res(data.text.trim() || null); worker.recognize(path)
}).finally(() => worker?.terminate()); .then(({data}: any) => res(data.text.trim() || null))
.catch(rej);
});
})().finally(() => {
process.off('uncaughtException', handler);
worker?.terminate();
});
return Object.assign(p, {abort: () => worker?.terminate()}); return Object.assign(p, {abort: () => worker?.terminate()});
} }
} }
+4 -1
View File
@@ -4,7 +4,10 @@
"target": "ESNext", "target": "ESNext",
"useDefineForClassFields": true, "useDefineForClassFields": true,
"module": "ESNext", "module": "ESNext",
"lib": ["ESNext"], "lib": [
"ESNext",
"dom"
],
"skipLibCheck": true, "skipLibCheck": true,
/* Bundler mode */ /* Bundler mode */