Compare commits

...

60 Commits
0.8.7 ... 1.6.4

Author SHA1 Message Date
0a6f1e4d62 Refined memory management prompts
All checks were successful
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
08a351e028 Better memory management
All checks were successful
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
85c01d3ef1 Added official file support
All checks were successful
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
5826573d5c Added official file support
All checks were successful
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
797a40a566 Added official file support
All checks were successful
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
7308927a3c max token rename
All checks were successful
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
04f038ba65 Memory prompt refinement
All checks were successful
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
d42f58d710 Memory refinement
All checks were successful
Publish Library / Build NPM Project (push) Successful in 54s
Publish Library / Tag Version (push) Successful in 11s
2026-08-05 12:22:13 -04:00
878a8794ee Rebuild graph edges on changes
All checks were successful
Publish Library / Build NPM Project (push) Successful in 46s
Publish Library / Tag Version (push) Successful in 19s
2026-08-04 17:05:58 -04:00
3f1289d993 Small agent tweaks
All checks were successful
Publish Library / Build NPM Project (push) Successful in 49s
Publish Library / Tag Version (push) Successful in 9s
2026-08-04 14:33:28 -04:00
077f75cdd9 Fixed delegate agent history... again
All checks were successful
Publish Library / Build NPM Project (push) Successful in 48s
Publish Library / Tag Version (push) Successful in 13s
2026-08-04 13:58:47 -04:00
566d84fd7a Added memory graph traversal helpers
All checks were successful
Publish Library / Build NPM Project (push) Successful in 43s
Publish Library / Tag Version (push) Successful in 14s
2026-08-04 12:58:39 -04:00
4230b534fc bump 1.4.0
All checks were successful
Publish Library / Build NPM Project (push) Successful in 1m17s
Publish Library / Tag Version (push) Successful in 14s
2026-08-04 12:44:41 -04:00
119f8472f2 token pools
Some checks failed
Publish Library / Tag Version (push) Has been cancelled
Publish Library / Build NPM Project (push) Has been cancelled
2026-08-04 12:44:21 -04:00
9c04e58c63 Pass deligate subagents full history, improved memory managment 2026-08-04 12:24:23 -04:00
7fbb42c26a improved subagent instructions 2026-08-04 12:03:31 -04:00
be08db8e2c Attach tps to response promise
All checks were successful
Publish Library / Build NPM Project (push) Successful in 42s
Publish Library / Tag Version (push) Successful in 9s
2026-08-04 09:48:20 -04:00
497f051c62 bump 1.3.5
All checks were successful
Publish Library / Build NPM Project (push) Successful in 52s
Publish Library / Tag Version (push) Successful in 7s
2026-08-04 09:30:45 -04:00
62fbe73b22 Added tps + duration to AI history
All checks were successful
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
d53b1c6328 Removed <tool> blocks from responses
All checks were successful
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
89619e211e Fixed message history and response
All checks were successful
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
afc6653364 fixed openai system calls in history breaking anthropic calls
All checks were successful
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
68e72445a2 Keep recent memories in context
All checks were successful
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
1aa6cdf329 Agent/subagent support
All checks were successful
Publish Library / Build NPM Project (push) Successful in 45s
Publish Library / Tag Version (push) Successful in 15s
2026-08-01 18:28:16 -04:00
d022a5ef4d Improved levenshtein fuzzy match
All checks were successful
Publish Library / Build NPM Project (push) Successful in 52s
Publish Library / Tag Version (push) Successful in 17s
2026-08-01 12:00:26 -04:00
a1d438a20a Tools can now emit "done" event and end chat early gracefully
All checks were successful
Publish Library / Build NPM Project (push) Successful in 1m0s
Publish Library / Tag Version (push) Successful in 9s
2026-07-31 17:49:06 -04:00
52a9e3aaa4 Fixed history poisoning on empty tool response
All checks were successful
Publish Library / Build NPM Project (push) Successful in 51s
Publish Library / Tag Version (push) Successful in 13s
2026-07-30 22:12:49 -04:00
a7aec4ee29 Improved memory prompt slightly
All checks were successful
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
dda2d4c2a3 Bump 1.2.8
All checks were successful
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
58e0e488e4 Added Geo, FS and flarescraperr tools
Some checks failed
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
8dfcd06752 More memory fixes
All checks were successful
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
14f6cdd313 Personal file memory organization instructions
All checks were successful
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
73d6ee0f2a Personal file memory organization instructions
All checks were successful
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
bee4085666 updatememory awaits full result
Some checks failed
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
3b5c71de7c Improved memory management
All checks were successful
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
8229e02a52 Improved memory management
All checks were successful
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
a6fb8ae828 New memory system
All checks were successful
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
d1230bcaad Updated wiki tool
All checks were successful
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
2d49c9aa80 Removed redundant llama protocol (Use openai)
All checks were successful
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
9a39f00f94 Diarization fix
All checks were successful
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
436757daad Added new json output support
Some checks failed
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
69b3297bb3 Proper error handling for OCR
All checks were successful
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
710c6ce52c Proper error handling for OCR
Some checks failed
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
4ac3036000 Proper error handling for OCR
All checks were successful
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
3121d542d4 OCR
All checks were successful
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
51ab8f2538 Memory / history fixes
All checks were successful
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
7dd3307a07 Update LLM models at runtime
All checks were successful
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
209d3b120b Export memory types
All checks were successful
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
0b1c25dfda Added MCP, Hybrid Memories and Skill support
All checks were successful
Publish Library / Build NPM Project (push) Successful in 56s
Publish Library / Tag Version (push) Successful in 16s
2026-06-06 22:02:19 -04:00
af6522ad88 Bump 0.9.0
All checks were successful
Publish Library / Build NPM Project (push) Successful in 26s
Publish Library / Tag Version (push) Successful in 9s
2026-03-29 23:01:30 -04:00
ee7b85301b * Fixed llm response object (double encoding)
All checks were successful
Publish Library / Build NPM Project (push) Successful in 25s
Publish Library / Tag Version (push) Successful in 12s
+ added wikitools
+ Improved webpage reading tool
2026-03-29 23:00:40 -04:00
d2e711fbf2 Added wikipedia tools
All checks were successful
Publish Library / Build NPM Project (push) Successful in 1m5s
Publish Library / Tag Version (push) Successful in 11s
2026-03-29 21:50:26 -04:00
596e99daa7 Use word count for summary (more predictable)
All checks were successful
Publish Library / Build NPM Project (push) Successful in 55s
Publish Library / Tag Version (push) Successful in 33s
2026-03-26 13:10:46 -04:00
eda4eed87d Added JSON / Summary LLM safeguard
All checks were successful
Publish Library / Build NPM Project (push) Successful in 41s
Publish Library / Tag Version (push) Successful in 21s
2026-03-26 12:50:52 -04:00
7f88c2d1d0 Added JSON / Summary LLM safeguard
All checks were successful
Publish Library / Build NPM Project (push) Successful in 1m17s
Publish Library / Tag Version (push) Successful in 13s
2026-03-26 12:33:50 -04:00
5eae84f6cf Added JSON / Summary LLM safeguard
All checks were successful
Publish Library / Build NPM Project (push) Successful in 1m1s
Publish Library / Tag Version (push) Successful in 14s
2026-03-26 12:24:20 -04:00
52a3e73484 Improved read_webpage tool
All checks were successful
Publish Library / Build NPM Project (push) Successful in 57s
Publish Library / Tag Version (push) Successful in 14s
2026-03-21 14:34:24 -04:00
ccb1bdf043 Added Non-UTC version of date/time tool
All checks were successful
Publish Library / Build NPM Project (push) Successful in 54s
Publish Library / Tag Version (push) Successful in 6s
2026-03-13 18:55:38 -04:00
b814ea8b28 Improved memory recall results
All checks were successful
Publish Library / Build NPM Project (push) Successful in 42s
Publish Library / Tag Version (push) Successful in 10s
2026-03-03 00:26:00 -05:00
06dda88dbc Removed log statements
All checks were successful
Publish Library / Build NPM Project (push) Successful in 36s
Publish Library / Tag Version (push) Successful in 10s
2026-03-02 14:00:58 -05:00
18 changed files with 4712 additions and 2050 deletions

120
README.md
View File

@@ -103,7 +103,125 @@ A TypeScript library that provides a unified interface for working with multiple
## Documentation ## Documentation
[Available Here](https://ai-utils.docs.zakscode.com/) ### Setup
```javascript
const ai = new Ai({
path: '/ai-models',
// Setup audio
whisper: '/path/to/binary', // Required for ASR
hfToken: '...', // Required for diarization
asr: 'ggml-base.en.bin', // Override default ASR model
// Setup LLM
embedder: 'bge-small-en-v1.5', // Override default embedder model
llm: {
system: 'You are a helpful assistant.',
compress: {max: 90_000, min: 50_000}, // Compress chat history to min tokens when max is reached
temperature: 0.8,
maxTokens: 100_000,
memoryModel: 'gpt-4o', // Cheap model for managing memories in background, defaults to current model
models: {
'claude-3-5-sonnet': {proto: 'anthropic', token: process.env.ANTHROPIC_TOKEN},
'gpt-4o': {proto: 'openai', token: process.env.OPENAI_TOKEN},
'llama3': {proto: 'ollama', host: 'http://localhost:11434'},
},
mcp: [
{name: 'files', url: 'https://mcp.example.com', token: process.env.MCP_TOKEN}
],
skills: [
{name: 'Tone of voice', description: 'Brand writing guidelines', content: '# Tone of Voice\n\nAlways be concise and friendly...'}
],
tools: [{
name: 'Marco?',
description: 'Where is marco polo?',
args: {
shout: {type: 'boolean', default: 'Shout into the void?', description: false, required: false}
},
fn: (args: any, stream: LLMRequest['stream'], ai: Ai) => {
const {shout} = args;
return shout ? 'Polo!' : 'Polo';
}
}],
},
// Setup Vision
ocr: 'eng' // Override default OCR model
});
```
### Audio
```javascript
// Crate audio transcript
const text = await ai.audio.asr('./path/to/audio.mp3');
console.log(text);
// Break transcript into speakers
const text = await ai.audio.asr('./path/to/audio.mp3', {diarization: true});
console.log(text);
// Break transcript into named speakers
const text = await ai.audio.asr('./path/to/audio.mp3', {diarization: 'llm'});
console.log(text);
```
### Language
```javascript
const history = [], memory = [];
// Wait for entire response
const text = await ai.language.ask('My favorite color is blue, whats yours?', {history, memory});
console.log(text);
// Stream response
const chunks = '';
await ai.language.ask('Write me a poem', {
history, memory,
stream: chunk => chunks += chunk,
});
console.log(chunks);
// Manually compile history into memories at end of conversation
// Happens automatically when coverstaions are compressed
await ai.language.memorize(history, memory);
// Summarize text
const summary = await ai.language.summarize(longText, 200);
// Code response (no conversation or extra BS)
const code = await ai.language.code('Write a fibonacci function');
// Structured JSON response
const data = await ai.language.json('Extract the name and age', `{
"name": "string",
"age": "number"
}`, {system: 'Extract from user input'});
```
#### Premade LLM Tools:
- `cli`: Run a shell command, returns its output
- `get_datetime`: Returns local date/time
- `get_datetime_utc`: Returns current UTC date/time
- `exec`: Execute code in cli, node, or python
- `fetch`: Make HTTP requests (GET/POST/PUT/DELETE)
- `exec_javascript`: Execute CommonJS JavaScript
- `exec_python`: Execute Python via python -c
- `read_webpage`: Scrape & clean content from a URL, handles HTML, JSON, CSV, media, PDFs etc.
- `web_search`: Anonymous DuckDuckGo search, returns a list of URLs
- `wikipedia_lookup`: Fetch a Wikipedia article (intro or full)
- `wikipedia_search`: Search Wikipedia and return matching articles
- `get_weather`: Fetch current weather + forecast for a location (just built!)
### Vision
```javascript
// Extract text from image
const text = await ai.vision.ocr('./path/to/image.png');
console.log(text);
```
## License ## License

3535
package-lock.json generated

File diff suppressed because it is too large Load Diff

View File

@@ -1,6 +1,6 @@
{ {
"name": "@ztimson/ai-utils", "name": "@ztimson/ai-utils",
"version": "0.8.7", "version": "1.6.4",
"description": "AI Utility library", "description": "AI Utility library",
"author": "Zak Timson", "author": "Zak Timson",
"license": "MIT", "license": "MIT",
@@ -25,21 +25,22 @@
"watch": "npx vite build --watch" "watch": "npx vite build --watch"
}, },
"dependencies": { "dependencies": {
"@anthropic-ai/sdk": "^0.78.0", "@anthropic-ai/sdk": "^0.102.0",
"@huggingface/transformers": "^4.2.0",
"@tensorflow/tfjs": "^4.22.0", "@tensorflow/tfjs": "^4.22.0",
"@xenova/transformers": "^2.17.2",
"@ztimson/node-utils": "^1.0.7", "@ztimson/node-utils": "^1.0.7",
"@ztimson/utils": "^0.28.13", "@ztimson/utils": "^0.29.4",
"cheerio": "^1.2.0", "cheerio": "^1.2.0",
"openai": "^6.22.0", "openai": "^6.42.0",
"pdf-parse": "^2.4.5",
"tesseract.js": "^7.0.0" "tesseract.js": "^7.0.0"
}, },
"devDependencies": { "devDependencies": {
"@types/node": "^24.8.1", "@types/node": "^24.13.1",
"typedoc": "^0.26.7", "typedoc": "^0.26.7",
"typescript": "^5.3.3", "typescript": "^5.6.3",
"vite": "^7.2.7", "vite": "^8.0.16",
"vite-plugin-dts": "^4.5.3" "vite-plugin-dts": "^5.0.2"
}, },
"files": [ "files": [
"dist" "dist"

View File

@@ -1,5 +1,5 @@
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';
@@ -8,7 +8,7 @@ export type AbortablePromise<T> = Promise<T> & {
}; };
export type AiOptions = { export type AiOptions = {
/** Token to pull models from hugging face */ /** Token to pull diarization models from hugging face */
hfToken?: string; hfToken?: string;
/** Path to models */ /** Path to models */
path?: string; path?: string;
@@ -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;

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: JSONSanitize(result)}; if(chunk.done) { terminal = true; return; }
options.stream!(chunk);
});
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
} catch(err: any) { } catch(err: any) {
return {type: 'tool_result', tool_use_id: toolCall.id, is_error: true, content: err?.message || err?.toString() || 'Unknown'}; entry.error = err?.message || err?.toString() || 'Unknown';
} }
})); }));
history.push({role: 'user', content: results}); } 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()});
} }
} }

View File

@@ -2,7 +2,6 @@ import {execSync, spawn} from 'node:child_process';
import {mkdtempSync} from 'node:fs'; import {mkdtempSync} from 'node:fs';
import fs from 'node:fs/promises'; import fs from 'node:fs/promises';
import {tmpdir} from 'node:os'; import {tmpdir} from 'node:os';
import * as path from 'node:path';
import Path, {join} from 'node:path'; import Path, {join} from 'node:path';
import {AbortablePromise, Ai} from './ai.ts'; import {AbortablePromise, Ai} from './ai.ts';
@@ -142,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;
} }
@@ -155,7 +161,7 @@ print(json.dumps(segments))
const p = new Promise<any>((resolve, reject) => { const p = new Promise<any>((resolve, reject) => {
this.downloadAsrModel(opts.model).then(m => { this.downloadAsrModel(opts.model).then(m => {
if(opts.diarization) { if(opts.diarization) {
let output = path.join(path.dirname(file), 'transcript'); let output = join(Path.dirname(file), 'transcript');
proc = spawn(<string>this.ai.options.whisper, proc = spawn(<string>this.ai.options.whisper,
['-m', m, '-f', file, '-np', '-ml', '1', '-oj', '-of', output], ['-m', m, '-f', file, '-np', '-ml', '1', '-oj', '-of', output],
{stdio: ['ignore', 'ignore', 'pipe']} {stdio: ['ignore', 'ignore', 'pipe']}
@@ -226,11 +232,11 @@ print(json.dumps(segments))
return <any>Object.assign(p, {abort}); return <any>Object.assign(p, {abort});
} }
asr(file: string, options: { model?: string; diarization?: boolean | 'llm' } = {}): AbortablePromise<string | null> { asr(path: string, options: { model?: string; diarization?: boolean | 'llm' } = {}): AbortablePromise<string | null> {
if(!this.ai.options.whisper) throw new Error('Whisper not configured'); if(!this.ai.options.whisper) throw new Error('Whisper not configured');
const tmp = join(mkdtempSync(join(tmpdir(), 'audio-')), 'converted.wav'); const tmp = join(mkdtempSync(join(tmpdir(), 'audio-')), 'converted.wav');
execSync(`ffmpeg -i "${file}" -ar 16000 -ac 1 -f wav "${tmp}"`, { stdio: 'ignore' }); execSync(`ffmpeg -i "${path}" -ar 16000 -ac 1 -f wav "${tmp}"`, { stdio: 'ignore' });
const clean = () => fs.rm(Path.dirname(tmp), {recursive: true, force: true}).catch(() => {}); const clean = () => fs.rm(Path.dirname(tmp), {recursive: true, force: true}).catch(() => {});
if(!options.diarization) return this.runAsr(tmp, {model: options.model}); if(!options.diarization) return this.runAsr(tmp, {model: options.model});

View File

@@ -1,13 +1,13 @@
import { pipeline } from '@xenova/transformers'; import { pipeline } from '@huggingface/transformers';
const [modelDir, model] = process.argv.slice(2); const [modelDir, model] = process.argv.slice(2);
let text = ''; let text = '';
process.stdin.on('data', chunk => text += chunk); process.stdin.on('data', chunk => text += chunk);
process.stdin.on('end', async () => { process.stdin.on('end', async () => {
const embedder = await pipeline('feature-extraction', 'Xenova/' + model, {quantized: true, cache_dir: modelDir}); const embedder = await pipeline('feature-extraction', 'Xenova/' + model, {cache_dir: modelDir});
const output = await embedder(text, { pooling: 'mean', normalize: true }); const output = await embedder(text, { pooling: 'mean', normalize: true });
const embedding = Array.from(output.data); const embedding = Array.from(output.data);
console.log(JSON.stringify({embedding})); process.stdout.write(JSON.stringify({embedding}));
process.exit(); process.exit();
}); });

85
src/helpers.ts Normal file
View File

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

View File

@@ -1,8 +1,11 @@
export * from './ai'; export * from './ai';
export * from './antrhopic'; export * from './antrhopic';
export * from './audio'; export * from './audio';
export * from './helpers';
export * from './llm'; export * from './llm';
export * from './memory';
export * from './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';

335
src/kd-tree.ts Normal file
View File

@@ -0,0 +1,335 @@
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;
}
// ─── Distance helpers ─────────────────────────────────────────────────────────
function euclidean(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);
}
function cosine(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; // distance = 1 - similarity
}
/**
* 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
* - 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 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" ? cosine : euclidean;
if (points && points.length > 0) {
this.validateAll(points);
this.root = this.buildBalanced([...points], 0);
this._size = points.length;
}
}
/** Total number of points stored in the tree. */
get size(): number { return this._size; }
// ── 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++;
}
// ── KNN search ─────────────────────────────────────────────────────────────
/**
* Find the k nearest 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 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 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 points as a balanced tree.
* Useful after many individual insertions to restore O(log n) query time.
*/
rebalance(): void {
const points = this.toArray();
this.root = points.length ? this.buildBalanced(points, 0) : null;
}
// ── 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;
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);
// Only explore the far side if it could contain a closer point.
// For cosine distance we can't prune by axis gap alone, so always explore.
const shouldExplore =
this.distanceFn === cosine
? 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;
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 === cosine ? 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;
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);
}
}

View File

@@ -1,24 +1,63 @@
import {JSONAttemptParse} from '@ztimson/utils'; 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 {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, MemoryCache, MemoryManager, MemoryOptions, stripHeader} 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';
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 Agent = {
name: string;
description?: string;
model?: string | null;
temperature?: number;
system: string;
delegate?: boolean;
skills?: Skill[] | null;
tools?: AiTool[] | null;
mcp?: McpServer[] | null;
agents?: string[] | 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,25 +73,21 @@ export type LLMMessage = {
error?: undefined | string; error?: undefined | string;
/** Timestamp */ /** Timestamp */
timestamp?: number; timestamp?: number;
} /** Response duration in ms */
duration?: number;
/** Background information the AI will be fed */ /** Tokens per second */
export type LLMMemory = { tps?: number;
/** What entity is this fact about */
owner: string;
/** The information that will be remembered */
fact: string;
/** Owner and fact embedding vector */
embeddings: [number[], 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 */
@@ -62,17 +97,49 @@ export type LLMRequest = {
/** Stream response */ /** Stream response */
stream?: (chunk: {text?: string, tool?: string, done?: true}) => any; stream?: (chunk: {text?: string, tool?: string, done?: true}) => any;
/** Compress old messages in the chat to free up context */ /** Compress old messages in the chat to free up context */
compress?: { compress?: {max: number; min: number};
/** Trigger chat compression once context exceeds the token count */ /** User's memory documents - RAG injected automatically each turn */
max: number; memory?: Memory[] | MemoryCache | MemoryOptions;
/** Compress chat until context size smaller than */ /** Model to use for memory operations */
min: number memoryModel?: string;
}, /** Skill documents the AI can browse and read on demand */
/** Background information the AI will be fed */ skills?: Skill[];
memory?: LLMMemory[], /** MCP servers to connect and expose as tools */
mcp?: McpServer[];
/** Subagents exposed as delegatable/wrapped tools */
agents?: Agent[];
/** Attach files to request */
files?: LLMFile[];
/** @internal recursion guard for nested agent delegation */
_agentDepth?: number;
}
export type McpServer = {
/** MCP server name for humans */
name: string;
/** Host URL */
host: string;
/** Server access token */
token?: string;
}
export type Skill = {
/** Name of skill for humans */
name: string;
/** Description LLM will use to decide to learn a skill */
description: string;
/** Skill instructions */
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;
defaultModel!: string; defaultModel!: string;
models: {[model: string]: LLMProvider} = {}; models: {[model: string]: LLMProvider} = {};
@@ -81,21 +148,246 @@ 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);
}
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;
} }
/** /**
* Chat with LLM * Extract text from a PDF. Pages with no text layer (scanned/image-only) are handled as either:
* @param {string} message Question * - Rendered to images and returned alongside the text so the (vision-capable) model can read them directly
* @param {LLMRequest} options Configuration options and chat history * - OCR'd via Tesseract when the doc is too large to reasonably pass as images
* @returns {{abort: () => void, response: Promise<string>}} Function to abort response and chat history
*/ */
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(agents: Agent[] = [], allAgents: Agent[], history: LLMMessage[], aborts: (() => void)[], depth = 0, delegateState: {resp: string | null}): AiTool[] {
return agents.map(a => {
const toolName = `${a.delegate ? '' : 'sub'}agent_${snakeCase(a.name)}`;
return {
name: toolName,
description: `${a.delegate ? 'Delegate to ' : ''}Subagent: ${a.description || a.name}`,
args: clean<any>({
context: !a.delegate ? {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true} : undefined,
instructions: {type: 'string', description: 'Detailed instructions for subagent to complete', required: true},
}),
fn: async (args: any, stream: any, ai: any, id?: string) => {
if(depth >= MAX_AGENT_DEPTH) return 'Max agent delegation depth exceeded';
const nested = (a.agents || [])
.map(name => allAgents.find(x => x.name === name))
.filter((x): x is Agent => !!x && x.name !== a.name);
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: nested,
_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[]}> {
if(!servers?.length) return {prompt: '', tools: []};
const allTools: AiTool[] = [];
await Promise.all(servers.map(async server => {
const res = await fetch(`${server.host}/tools`, {headers: server.token ? {Authorization: `Bearer ${server.token}`} : {}});
const mcp: any = await res.json();
if(!mcp?.tools) return;
for(const t of mcp.tools) {
const args: Record<string, any> = {};
if(t.inputSchema?.properties) {
for(const [key, val] of Object.entries<any>(t.inputSchema.properties)) {
args[key] = {type: val.type || 'string', description: val.description || '', required: t.inputSchema.required?.includes(key)};
}
}
allTools.push({
name: `${server.name}_${t.name}`,
description: t.description || '',
args,
fn: async (a: any) => {
const r = await fetch(`${server.host}/tools/call`, {
method: 'POST',
headers: {'Content-Type': 'application/json', ...(server.token ? {Authorization: `Bearer ${server.token}`} : {})},
body: JSON.stringify({name: t.name, arguments: a})
});
const data: any = await r.json();
return data?.content?.[0]?.text ?? JSON.stringify(data);
}
});
}
}));
const list = allTools.map(t => `- ${t.name}: ${t.description}`).join('\n');
return {
prompt: `## MCP\nYou have access to the following MCP tools:\n${list}`,
tools: allTools
};
}
private setupSkills(skills: Skill[] = []): {prompt: string, tools: AiTool[]} {
if(!skills?.length) return {prompt: '', tools: []};
const list = skills.map(s => `- ${s.name}: ${s.description}`).join('\n');
return {
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: [{
name: 'skill_read',
description: 'Read the full content of a skill/knowledge document',
args: {
name: {type: 'string', description: 'Exact skill name', required: true}
},
fn: (args: any) => {
const skill = skills.find(s => s.name === args.name);
if(!skill) return `Skill not found. Available:\n${list}`;
return `# ${skill.name}\n${skill.content}`;
}
}]
}
}
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: [],
@@ -103,86 +395,145 @@ 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;
if(!options.history) options.history = []; const nestedAborts: (() => void)[] = [];
// If memories were passed, find any relevant ones and add a tool for ADHOC lookups const abort = () => {
if(options.memory) { aborted = true;
const search = async (query?: string | null, subject?: string | null, limit = 10) => { request?.abort?.();
const [o, q] = await Promise.all([ nestedAborts.forEach(a => a());
subject ? this.embedding(subject) : Promise.resolve(null), };
query ? this.embedding(query) : Promise.resolve(null),
]); let promise: any;
return (options.memory || []).map(m => { const requestStart = Date.now();
const score = (o ? this.cosineSimilarity(m.embeddings[0], o[0].embedding) : 0)
+ (q ? this.cosineSimilarity(m.embeddings[1], q[0].embedding) : 0); promise = (async () => {
return {...m, score}; let tools: AiTool[] = options.tools || this.ai.options.llm?.tools || [];
}).toSorted((a: any, b: any) => a.score - b.score).slice(0, limit); const prompts: string[] = [];
let history = options.history || [];
const files = options.files || [];
if(message || files.length) history.push({role: 'user', content: message || '', timestamp: Date.now()});
// MCP
const mcp = options.mcp || this.ai.options?.llm?.mcp;
if(mcp?.length) {
const m = await this.setupMcp(mcp);
prompts.unshift(m.prompt);
tools.push(...m.tools);
} }
options.system += '\nYou have RAG memory and will be given the top_k closest memories regarding the users query. Save anything new you have learned worth remembering from the user message using the remember tool and feel free to recall memories manually.\n'; // Skills
const relevant = await search(message); const skills = options.skills || this.ai.options?.llm?.skills;
if(relevant.length) options.history.push({role: 'tool', name: 'recall', id: 'auto_recall_' + Math.random().toString(), args: {}, content: 'Things I remembered:\n' + relevant.map(m => `${m.owner}: ${m.fact}`).join('\n')}); if(skills?.length) {
options.tools = [{ const s = this.setupSkills(skills);
name: 'recall', prompts.unshift(s.prompt);
description: 'Recall the closest memories you have regarding a query using RAG', tools.push(...s.tools);
args: {
subject: {type: 'string', description: 'Find information by a subject topic, can be used with or without query argument'},
query: {type: 'string', description: 'Search memory based on a query, can be used with or without subject argument'},
topK: {type: 'number', description: 'Result limit, default 5'},
},
fn: (args) => {
if(!args.subject && !args.query) throw new Error('Either a subject or query argument is required');
return search(args.query, args.subject, args.topK);
}
}, {
name: 'remember',
description: 'Store important facts user shares for future recall',
args: {
owner: {type: 'string', description: 'Subject/person this fact is about'},
fact: {type: 'string', description: 'The information to remember'}
},
fn: async (args) => {
if(!options.memory) return;
const e = await Promise.all([
this.embedding(args.owner),
this.embedding(`${args.owner}: ${args.fact}`)
]);
const newMem = {owner: args.owner, fact: args.fact, embeddings: <any>[e[0][0].embedding, e[1][0].embedding]};
options.memory.splice(0, options.memory.length, ...[
...options.memory.filter(m => {
return !(this.cosineSimilarity(newMem.embeddings[0], m.embeddings[0]) >= 0.9 && this.cosineSimilarity(newMem.embeddings[1], m.embeddings[1]) >= 0.8);
}),
newMem
]);
return 'Remembered!';
}
}, ...options.tools || []];
} }
// Ask // Agents
const resp = await this.models[m].ask(message, options); const agents = options.agents || this.ai.options?.llm?.agents;
const delegateState: {resp: string | null} = {resp: null};
if(agents?.length) tools.push(...this.setupAgent(agents, agents, history, nestedAborts, options._agentDepth || 0, delegateState));
// Remove any memory calls from history // Memory
if(options.memory) options.history.splice(0, options.history.length, ...options.history.filter(h => h.role != 'tool' || (h.name != 'recall' && h.name != 'remember'))); const mem = MemoryManager.normalize(options.memory);
if(mem) {
const mems = mem.memory instanceof MemoryCache ? mem.memory.memories : mem.memory;
if(mems.length) {
if(mem.inject) {
const pool = 15; // candidates considered, cheap since only refs are listed
const budget = mem.maxTokens ?? 2000; // actual content injected
const relevant = await this.memoryManager.recollect(message, mem.memory, pool);
// Compress message history let used = 0;
if(options.compress) { const preloaded: typeof relevant = [];
const compressed = await this.ai.language.compressHistory(options.history, options.compress.max, options.compress.min, options); const listed: typeof relevant = [];
options.history.splice(0, options.history.length, ...compressed); 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);
} }
return res(resp); prompts.unshift(`## Memory
}), {abort}); 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` : ''}
${preloaded.length ? `### Prefetched Memories (Most relevant first):
${preloaded.map(r => `Memory: ${r.name}
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));
}
} }
async code(message: string, options?: LLMRequest): Promise<any> { if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
const resp = await this.ask(message, {...options, system: [
options?.system, const lastMsg = history[history.length - 1];
'Return your response in a code block' if(files.length && lastMsg?.role === 'user') lastMsg.files = files;
].filter(t => !!t).join(('\n'))}); const restores: {msg: LLMMessage, content: any}[] = [];
const codeBlock = /```(?:.+)?\s*([\s\S]*?)```/.exec(resp); for(const msg of history) {
return codeBlock ? codeBlock[1].trim() : null; 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;
}
const toolTimings = new Map<string, {duration: number, tps: number}>();
tools = this.wrapToolTiming(tools, toolTimings);
if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
prompts.unshift(options.system || this.ai.options.llm?.system || '');
request = this.models[m].ask('', {...options, tools, system: prompts.filter(Boolean).join('\n\n')});
let resp = await request;
// Strip the file injection shim
restores.forEach(({msg, content}) => msg.content = content);
// Capture meta (duration / tps)
for(const h of history) {
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});
} }
/** /**
@@ -271,7 +622,7 @@ class LLM {
* @param {maxTokens?: number, overlapTokens?: number} opts Options for embedding such as chunk sizes * @param {maxTokens?: number, overlapTokens?: number} opts Options for embedding such as chunk sizes
* @returns {Promise<Awaited<{index: number, embedding: number[], text: string, tokens: number}>[]>} Chunked embeddings * @returns {Promise<Awaited<{index: number, embedding: number[], text: string, tokens: number}>[]>} Chunked embeddings
*/ */
embedding(target: object | string, opts: {maxTokens?: number, overlapTokens?: number} = {}): AbortablePromise<any[]> { embedding(target: object | string, opts: {maxTokens?: number, overlapTokens?: number} = {}): AbortablePromise<{index: number, embedding: number[], text: string, tokens: number}[]> {
let {maxTokens = 500, overlapTokens = 50} = opts; let {maxTokens = 500, overlapTokens = 50} = opts;
let aborted = false; let aborted = false;
const abort = () => { aborted = true; }; const abort = () => { aborted = true; };
@@ -279,7 +630,6 @@ class LLM {
const embed = (text: string): Promise<number[]> => { const embed = (text: string): Promise<number[]> => {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
if(aborted) return reject(new Error('Aborted')); if(aborted) return reject(new Error('Aborted'));
const args: string[] = [ const args: string[] = [
join(dirname(fileURLToPath(import.meta.url)), 'embedder.js'), join(dirname(fileURLToPath(import.meta.url)), 'embedder.js'),
<string>this.ai.options.path, <string>this.ai.options.path,
@@ -288,7 +638,6 @@ class LLM {
const proc = spawn('node', args, {stdio: ['pipe', 'pipe', 'ignore']}); const proc = spawn('node', args, {stdio: ['pipe', 'pipe', 'ignore']});
proc.stdin.write(text); proc.stdin.write(text);
proc.stdin.end(); proc.stdin.end();
let output = ''; let output = '';
proc.stdout.on('data', (data: Buffer) => output += data.toString()); proc.stdout.on('data', (data: Buffer) => output += data.toString());
proc.on('close', (code: number) => { proc.on('close', (code: number) => {
@@ -298,7 +647,7 @@ class LLM {
const result = JSON.parse(output); const result = JSON.parse(output);
resolve(result.embedding); resolve(result.embedding);
} catch(err) { } catch(err) {
reject(new Error('Failed to parse embedding output')); reject(err);
} }
} else { } else {
reject(new Error(`Embedder process exited with code ${code}`)); reject(new Error(`Embedder process exited with code ${code}`));
@@ -318,7 +667,7 @@ class LLM {
} }
return results; return results;
})(); })();
return Object.assign(p, { abort }); return <any>Object.assign(p, {abort});
} }
/** /**
@@ -337,41 +686,98 @@ 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[]> {
const code = await this.code(text, {...options, system: [ return this.memoryManager.memorize(history, memories, {model: this.defaultModel, ...options});
options?.system,
`Only respond using JSON matching this schema:\n\`\`\`json\n${schema}\n\`\`\``
].filter(t => !!t).join('\n')});
return code ? JSONAttemptParse(code, {}) : null;
} }
/** /**
* Create a summary of some text * Create a summary of some text
* @param {string} text Text to summarize * @param {string} text Text to summarize
* @param {number} tokens Max number of tokens * @param {number} length Max number of words
* @param options LLM request options * @param options LLM request options
* @returns {Promise<string>} Summary * @returns {Promise<string>} Summary
*/ */
summarize(text: string, tokens: number = 500, options?: LLMRequest): Promise<string | null> { async summarize(text: string, length: number = 500, options?: LLMRequest): Promise<string | null> {
return this.ask(text, {system: `Generate the shortest summary possible <= ${tokens} tokens. Output nothing else`, temperature: 0.3, ...options}); let system = `Your job is to summarize the users message using tool calls. Call the \`submit\` tool at least once with the shortest summary possible that's <= ${length} words. The tool call will respond with the token count. Responses are ignored`;
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 summary',
args: {summary: {type: 'string', description: 'Text summarization', required: true}},
fn: (args) => {
if(!args.summary) return 'No summary provided';
const count = args.summary.split(' ').length;
if(count > length) return `Too long: ${length} words`;
done = true;
resolve(args.summary || null);
return `Saved: ${length} words`;
}
}, ...(options?.tools || [])],
});
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] ?? '';
} }
} }

644
src/memory.ts Normal file
View File

@@ -0,0 +1,644 @@
import {MemoryNode, rebuildGraph} from './helpers.ts';
import {LLMRequest, LLMMessage} from './llm.ts';
import {AiTool} from './tools.ts';
import {KDPoint, KDTree} from './kd-tree.ts';
import {escapeRegex} from '@ztimson/utils';
const MERGE_THRESHOLD = 0.88;
const PENDING_HEADING = '## Pending';
const GENERIC_TEMPLATE = `# {{Title}}
## Summary
## Details
## Related`;
export type Memory = {
name: string;
description: string;
content: string;
embedding: number[];
links: string[];
backlinks: string[];
}
type MemoryRef = {
name: string;
description: string;
}
type FactBucket = {
subject: string;
facts: string[];
}
type FactAgentResult = {
buckets: FactBucket[];
journal: string;
}
function dedupeFacts(facts: string[]): string[] {
const seen = new Map<string, string>();
for (const f of facts) {
const clean = f.trim();
if (clean) seen.set(clean.toLowerCase(), clean);
}
return [...seen.values()];
}
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;
}
function cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] {
return memories
.filter(m => m.embedding?.length)
.map(m => ({ref: {name: m.name, description: m.description}, distance: cosineDistance(query, m.embedding)}))
.sort((a, b) => a.distance - b.distance)
.slice(0, limit)
.map(s => s.ref);
}
export function stripHeader(content: string): string {
return content.replace(/^---[\s\S]*?\n---\n?/, '').trimStart();
}
export class MemoryCache {
private tree!: KDTree<MemoryRef>;
public memories: Memory[];
public nodes: MemoryNode[] = [];
get length() { return this.memories.length; }
constructor(memories: Memory[]) {
this.memories = memories;
this.rebuild();
}
private buildTree(): KDTree<MemoryRef> {
const embedded = this.memories.filter(m => m.embedding?.length);
if (!embedded.length) return new KDTree<MemoryRef>(0);
const dims = embedded[0].embedding.length;
const points: KDPoint<MemoryRef>[] = embedded.map(m => ({
vector: m.embedding,
payload: {name: m.name, description: m.description},
}));
return new KDTree<MemoryRef>(dims, 'cosine', points);
}
search(query: number[], limit: number): MemoryRef[] {
if (!this.tree || this.tree.dims === 0) return [];
return this.tree.knn(query, limit).map(r => r.point.payload);
}
add(memory: Memory): void {
this.memories.push(memory);
this.rebuild();
}
update(memory: Memory): void {
const existing = this.memories.find(m => m.name === memory.name);
if (existing) Object.assign(existing, memory);
else this.memories.push(memory);
this.rebuild();
}
remove(name: string): void {
const idx = this.memories.findIndex(m => m.name === name);
if (idx !== -1) {
this.memories.splice(idx, 1);
this.rebuild();
}
}
rebuild(): void {
this.nodes = rebuildGraph(this.memories);
this.tree = this.buildTree();
}
}
class MemoryAccessor {
readonly list: Memory[];
private readonly cache: MemoryCache | null;
constructor(memories: Memory[] | MemoryCache) {
this.cache = memories instanceof MemoryCache ? memories : null;
this.list = this.cache ? this.cache.memories : <Memory[]>memories;
}
find(name: string): Memory | undefined {
return this.list.find(m => m.name === name);
}
commit(): MemoryNode[] {
if (this.cache) {
this.cache.rebuild();
return this.cache.nodes;
}
return rebuildGraph(this.list);
}
ghosts(): string[] {
const nodes = this.cache ? this.cache.nodes : rebuildGraph(this.list);
return nodes.filter(n => n.missing).map(n => n.name);
}
search(vector: number[], limit: number): MemoryRef[] {
return this.cache ? this.cache.search(vector, limit) : cosineSearch(vector, this.list, limit);
}
forget(name: string): boolean {
const idx = this.list.findIndex(m => m.name === name);
if (idx === -1) return false;
this.list.splice(idx, 1);
this.commit();
return true;
}
async backfillEmbeddings(llm: any): Promise<number> {
const missing = this.list.filter(m => !m.embedding?.length);
if (!missing.length) return 0;
await Promise.all(missing.map(async node => {
const [e] = await llm.embedding(`${node.description}\n\n${stripHeader(node.content)}`.trim());
if (e) node.embedding = e.embedding;
}));
this.commit();
return missing.length;
}
}
export type MemoryOptions = {
/** Memory object */
memory: Memory[] | MemoryCache;
/** Inject N memories into the system prompt */
inject?: boolean;
/** expose recall tool to LLM */
tool?: boolean;
/** Update memory on compression */
update?: boolean;
/** Max context size of memories to inject to each call (removed immediately after use) */
maxTokens?: number;
}
export class MemoryManager {
private recentlyTouched = new Map<string, number>();
private queues = new Map<string, {
dirty: boolean,
request: {abort?: () => void} | null,
task: Promise<void>,
}>();
tools = {
forget: (memories: Memory[] | MemoryCache): AiTool => ({
name: 'memory_forget',
description: 'Permanently delete a memory document and clean up all references to it',
args: {
name: {type: 'string', description: 'Exact memory name to forget', required: true}
},
fn: (args: any) => {
const result = this.forget(args.name, memories);
return result ? `Forgotten: ${args.name}` : `Not found: ${args.name}`;
},
}),
read: (memories: Memory[] | MemoryCache): AiTool => ({
name: 'memory_recall',
description: 'Read the full content of a memory document',
args: {
name: {type: 'string', description: 'Exact memory name', required: true},
},
fn: (args: any) => {
const mem = this.access(memories).find(args.name);
if (!mem) return 'Document not found';
this.touch(mem.name);
return mem.content;
},
}),
search: (memories: Memory[] | MemoryCache): AiTool => ({
name: 'memory_search',
description: 'Use embeddings to find the MOST relevant memories, even if NOT relevant',
args: {
query: {type: 'string', description: 'What to look for in the memories', required: true},
limit: {type: 'number', description: 'Number of memories to return', default: 1},
},
fn: async ({query, limit}) => {
const mem = await this.recollect(query, memories, limit)
return mem.map(m => `Memory: ${m.name}
Description: ${m.description}
Links: ${[...m.links, ...m.backlinks].join(', ')}
\`\`\`
${m.content}
\`\`\``).join('\n\n');
},
}),
};
constructor(private llm: any) {}
static normalize(m?: Memory[] | MemoryCache | MemoryOptions) {
if (!m) return null;
const raw = m instanceof MemoryCache || Array.isArray(m);
return raw ? {memory: <Memory[] | MemoryCache>m, inject: true, tool: true, update: true} : {inject: true, tool: true, update: true, ...m};
}
private access(memories: Memory[] | MemoryCache): MemoryAccessor {
return new MemoryAccessor(memories);
}
private stage(node: Memory, block: string): void {
this.ensureDoc(node);
const body = stripHeader(node.content);
const idx = body.indexOf(PENDING_HEADING);
const newBody = idx === -1
? `${body.trimEnd()}\n\n${PENDING_HEADING}\n${block}\n`
: `${body.slice(0, idx + PENDING_HEADING.length)}\n${block}${body.slice(idx + PENDING_HEADING.length)}`;
node.content = this.touchHeader(node, newBody);
}
private ensureDoc(node: Memory): void {
if (node.content) return;
const title = node.name.split('/').pop() ?? node.name;
node.content = this.touchHeader(node, `# ${title}\n`);
}
private sanitizeDescription(text: string): string {
return (text ?? '').replace(/\s+/g, ' ').trim().slice(0, 240);
}
private relink(memories: Memory[], from: string, to: string): void {
const pattern = new RegExp(`\\[\\[${escapeRegex(from)}\\]\\]`, 'g');
for (const m of memories) if (pattern.test(m.content)) m.content = m.content.replace(pattern, `[[${to}]]`);
}
private async factAgent(conversation: string, store: MemoryAccessor, options: LLMRequest): Promise<FactAgentResult> {
const ghosts = store.ghosts();
const response = await this.llm.ask(conversation, {
model: options.model,
temperature: 0.2,
system: `You are a fact extractor for Obsidian-style knowledge vaults. Analyze the conversation and produce:
1. Journal recap (single paragraph)
- "Captains Log" style record keeping
- What was discussed/worked on, decisions, user's events/state/mood, general context
- Leave empty only for trivial/empty exchanges/small talk
2. Fact buckets
- ONLY facts the USER explicitly stated about themselves, their work, projects, or decisions made during this conversation
- NEVER extract greetings, pleasantries, or anything the assistant itself said
- Extract the final/end state, not deltas
Path assignment rules:
- Reuse existing node names whenever possible
- Documents should be grouped and named by the root subject
- Person → People/Name
- Project → Projects/Name
- Concept → Concepts/Name
- A bug report, its investigation, should be nested and attached to the same root subject node
- Tickets/one-off tasks → file under the project/name/component they belong to
- Only create a new top-level node when the fact belongs to a genuinely new subject (person/project/concept)\`
Available nodes:
${this.listNodes(store.list).map(n => `- ${n.name}: ${n.description}`).join('\n') || 'None yet.'}
${ghosts.length ? `${ghosts.map(g => `- ${g}: (Ghost)`).join('\n')}` : ''}`,
schema: {
journal: {type: 'string', description: 'Short day-to-day recap; empty if nothing happened.', required: false},
buckets: {type: 'array', description: 'Groups of facts to remember; empty array if nothing worth storing.', items: {
type: 'object', items: {
subject: {type: 'string', description: 'Exact node name or new path (e.g. "People/Sarah", "Projects/Oxide")', required: true},
facts: {type: 'array', description: 'Facts to store here', items: {type: 'string'}},
},
},
},
},
});
const buckets = new Map<string, string[]>();
for (const bucket of response.buckets ?? []) {
const subject = bucket.subject.trim();
const facts = buckets.get(subject) ?? [];
facts.push(...dedupeFacts(bucket.facts));
buckets.set(subject, facts);
}
return {
buckets: buckets.entries().toArray().map(([subject, facts]) => ({subject, facts})),
journal: (response.journal ?? '').trim(),
};
}
private getWeekMonday(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);
}
private listNodes(memories: Memory[]): MemoryRef[] {
return memories.map(m => ({name: m.name, description: m.description}));
}
/** Find the nearest node above the similarity threshold and fold the smaller/less-connected one into
* the other. Journals are exempt — they're partitioned by date, not topic, and merging across weeks
* would wreck the timeline. Returns 'merged' if `node` absorbed another (caller should re-run the doc
* agent), 'absorbed' if `node` itself got folded away (caller should stop touching it), or null. */
private async checkMerge(node: Memory, memories: Memory[] | MemoryCache, options: LLMRequest, threshold = MERGE_THRESHOLD): Promise<Memory | null> {
if (!node.embedding?.length || node.name.startsWith('Journal/')) return null;
const store = this.access(memories);
let closest: Memory | null = null, closestDist = Infinity;
for (const other of store.list) {
if (other.name === node.name || other.name.startsWith('Journal/') || !other.embedding?.length) continue;
const d = cosineDistance(node.embedding, other.embedding);
if (d < closestDist) { closestDist = d; closest = other; }
}
if (!closest || closestDist > threshold) return null;
const result = await this.mergeAgent(node, closest, options);
const merged: Memory = {name: result.name, description: this.sanitizeDescription(result.description), content: '', embedding: [], links: [], backlinks: []};
merged.content = this.touchHeader(merged, result.content);
const [e] = await this.llm.embedding(`${merged.description}\n\n${result.content}`.trim());
if (e) merged.embedding = e.embedding;
this.relink(store.list, node.name, merged.name);
this.relink(store.list, closest.name, merged.name);
this.queues.get(closest.name)?.request?.abort?.();
this.queues.delete(closest.name);
store.forget(node.name);
store.forget(closest.name);
store.list.push(merged);
store.commit();
return merged;
}
private reconcile(node: Memory, memories: Memory[] | MemoryCache, options: LLMRequest): Promise<void> {
const key = node.name;
const existing = this.queues.get(key);
if (existing) {
existing.dirty = true;
existing.request?.abort?.();
return existing.task;
}
const entry = {dirty: false, request: null, task: Promise.resolve()};
this.queues.set(key, entry);
const store = this.access(memories);
entry.task = (async () => {
let current = node;
do {
entry.dirty = false;
await this.docAgent(current, store.list, options, entry);
const merged = await this.checkMerge(current, memories, options);
if (merged) { current = merged; entry.dirty = true; }
} while (entry.dirty);
})().finally(() => {
this.queues.delete(key);
store.commit();
});
return entry.task;
}
private async docAgent(node: Memory, memories: Memory[], options: LLMRequest, entry: {request: {abort?: () => void} | null}): Promise<void> {
if(!memories.includes(node)) return;
const currentBody = stripHeader(node.content);
let update;
try {
for (let i = 0; i < 2 && !update?.content; i++) {
const request = this.llm.ask(currentBody, {
model: options.model,
temperature: 0.3,
schema: {
description: {type: 'string', description: 'One factual sentence describing the document\'s ENTIRE SUBJECT MATTER — for use as a search/merge fingerprint', required: true},
content: {type: 'string', description: 'Rewritten document body in markdown, without the frontmatter block', required: true},
},
system: `You are a knowledge base editor maintaining one Obsidian-style document.
If it has a "## Pending" section, fold all new material into the appropriate part, resolve overlap, then remove the section entirely. If no section, just tidy per the rules below.
Use this loose structure, adapting headings to what the content needs:
# Title
## Summary
## Details
## Related
Rules:
- Contradictions: newer facts always win — delete outdated statements entirely
- Journals (Journal/...): keep entries as a chronological timeline; clean up grammar within entries but never delete history
- Use Obsidian markdown: # headings, **bold**, bullet/numbered lists, tables for 2D data
- Link specific entities and concepts with [[WikiLink]] (e.g., [[Projects/KiwixServer]]); skip generics
- Keep concise, factual, human-readable
- NO frontmatter, filler, preamble, or AI commentary
Available nodes to link to (don't duplicate their content):
${this.listNodes(memories).filter(n => n.name !== node.name).map(n => n.name).join(', ') || 'none'}
Current document:
\`\`\`markdown
${currentBody}
\`\`\``,
});
entry.request = request;
update = await request;
}
} catch (err: any) {
if (err?.name === 'AbortError') return;
throw err;
} finally {
entry.request = null;
}
if (!update?.content) return;
node.description = node.name !== 'People/User' ? this.sanitizeDescription(update.description) : 'All information about the current user';
node.content = this.touchHeader(node, update.content);
const [e] = await this.llm.embedding(`${node.description}\n\n${update.content}`.trim());
if (e) node.embedding = e.embedding;
}
private async mergeAgent(a: Memory, b: Memory, options: LLMRequest): Promise<{name: string, description: string, content: string}> {
return this.llm.ask('', {
model: options.model,
temperature: 0.3,
schema: {
name: {type: 'string', description: 'New path for the merged doc, collection/subject format (e.g. Projects/Oxide) — only reuse an old title if it\'s genuinely the best fit', required: true},
description: {type: 'string', description: 'One factual sentence describing the merged document\'s subject matter', required: true},
content: {type: 'string', description: 'Fully reconciled body in markdown, without frontmatter', required: true},
},
system: `You are a knowledge base editor merging two overlapping Obsidian documents into one. Newer facts win on contradiction.
Structure loosely:
# Title
## Summary
## Details
## Related
Combine both documents, resolve duplication and contradictions.
Document A ("${a.name}"):
\`\`\`markdown
${stripHeader(a.content)}
\`\`\`
Document B ("${b.name}"):
\`\`\`markdown
${stripHeader(b.content)}
\`\`\``,
});
}
private 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 { /* legacy unquoted value, keep raw */ }
fm.set(key, value);
}
return {fm, body: match[2]};
}
private touchHeader(node: Memory, body: string): string {
const {fm} = this.parseFrontmatter(node.content);
fm.set('name', node.name);
fm.set('description', node.description || '');
fm.set('modified', new Date().toISOString());
return this.writeFrontmatter(fm, body);
}
private 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()}`;
}
decay() {
for (const [name, ttl] of this.recentlyTouched) {
if (ttl <= 1) this.recentlyTouched.delete(name);
else this.recentlyTouched.set(name, ttl - 1);
}
}
touch(name: string, ttl = 2) {
this.recentlyTouched.set(name, ttl);
}
forget(name: string, memories: Memory[] | MemoryCache): boolean {
return this.access(memories).forget(name);
}
async recollect(query: string, memories: Memory[] | MemoryCache, limit = 5, graphDepth = 1): Promise<Memory[]> {
const store = this.access(memories);
if (!store.list.length) return [];
await store.backfillEmbeddings(this.llm);
const [e] = await this.llm.embedding(query);
if (!e) return [];
const vectorResults = store.search(e.embedding, limit);
const found = new Set<string>(vectorResults.map(r => r.name));
if (graphDepth > 0) {
let frontier = [...found];
for (let depth = 0; depth < graphDepth && frontier.length; depth++) {
const next: string[] = [];
for (const name of frontier) {
const node = store.find(name);
if (!node) continue;
for (const link of node.links) {
if (!found.has(link) && store.find(link)) {
found.add(link);
next.push(link);
}
}
}
frontier = next;
}
}
const vectorOrder = vectorResults.map(r => r.name);
const graphExpansions = [...found].filter(n => !vectorOrder.includes(n));
return [...vectorOrder, ...graphExpansions].map(n => store.find(n)!).filter(Boolean);
}
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise<Memory[]> {
const conversation = history
.filter(h => h.role === 'user' || h.role === 'assistant')
.map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim();
if (!conversation) return [];
const uid = `${Date.now()}_${Math.random().toString(36).slice(2)}`;
const pending = {role: 'tool', name: 'memory_process', id: uid, content: conversation} as unknown as LLMMessage;
history.push(pending);
const store = this.access(memories);
const {buckets, journal} = await this.factAgent(conversation, store, options);
const touched: Memory[] = [];
if (journal) {
const journalName = `Journal/${this.getWeekMonday()}`;
let jnode = store.find(journalName);
if (!jnode) {
jnode = {name: journalName, description: '', content: '', embedding: [], links: [], backlinks: []};
store.list.push(jnode);
}
this.stage(jnode, `### ${new Date().toISOString().slice(0, 10)}\n${journal}`);
touched.push(jnode);
}
for (const {subject, facts} of buckets) {
let node = store.find(subject);
if (!node) {
node = {name: subject, description: '', content: '', embedding: [], links: [], backlinks: []};
store.list.push(node);
}
this.stage(node, facts.map(f => `- ${f}`).join('\n'));
touched.push(node);
}
for (const node of touched) {
const [e] = await this.llm.embedding(`${node.description}\n\n${stripHeader(node.content)}`.trim());
if (e) node.embedding = e.embedding;
this.touch(node.name);
}
if (touched.length) {
store.commit();
(pending as any).content = `Saved to ${touched.map(n => `[[${n.name}]]`).join(', ')}`;
await Promise.all(touched.map(node => this.reconcile(node, memories, options).catch(() => {})));
} else {
(pending as any).content = 'Nothing worth remembering.';
}
(touched as any).uid = uid;
return touched;
}
async reconcileAll(memories: Memory[] | MemoryCache, options: LLMRequest, scope: 'touched' | 'all' = 'touched'): Promise<void> {
const store = this.access(memories);
const targets = scope === 'all' ? store.list : store.list.filter(m => m.content.includes(PENDING_HEADING));
await Promise.all(targets.map(node => this.reconcile(node, memories, options)));
store.commit();
}
}

View File

@@ -1,84 +1,72 @@
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 {
for(let i = 0; i < history.length; i++) { let client = this.clients.get(token);
const h = history[i]; if(!client) {
if(h.role === 'assistant' && h.tool_calls) { client = new openAI(clean({baseURL: this.host, apiKey: token || undefined}));
const tools = h.tool_calls.map((tc: any) => ({ this.clients.set(token, client);
role: 'tool',
id: tc.id,
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); return client;
i--;
}
if(!history[i]?.timestamp) history[i].timestamp = Date.now();
}
return history;
} }
private fromStandard(history: LLMMessage[]): any[] { private toWireContent(content: any): any {
return history.reduce((result, h) => { 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(const h of history) {
if(h.role === 'tool') { if(h.role === 'tool') {
result.push({ wire.push({
role: 'assistant', role: 'assistant',
content: null, content: null,
tool_calls: [{id: h.id, type: 'function', function: {name: h.name, arguments: JSON.stringify(h.args)}}], tool_calls: [{id: h.id, type: 'function', function: {name: h.name, arguments: JSON.stringify(h.args)}}],
refusal: null,
annotations: []
}, { }, {
role: 'tool', role: 'tool',
tool_call_id: h.id, tool_call_id: h.id,
content: h.error || h.content content: h.error || h.content || '',
}); });
} else { } else {
const {timestamp, ...rest} = h; wire.push({role: h.role, content: this.toWireContent(h.content)});
result.push(rest);
} }
return result; }
}, [] as any[]); 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, 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 || undefined,
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 => ({
type: 'function', type: 'function',
function: { function: {
@@ -93,77 +81,96 @@ export class OpenAi extends LLMProvider {
})) }))
}; };
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;
do { do {
resp = await this.client.chat.completions.create(requestParams).catch(err => { requestParams.messages = this.toWire(history.filter(h => h.role !== 'system'), options.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).chat.completions.create(requestParams)).catch(err => {
err.message += `\n\nMessages:\n${JSON.stringify(requestParams.messages, null, 2)}`;
throw err; throw err;
}); });
let usage: any, msg: any = {content: '', tool_calls: []};
if(options.stream) { if(options.stream) {
if(!isFirstMessage) options.stream({text: '\n\n'});
else isFirstMessage = false;
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; if(chunk.choices[0]?.delta?.content) {
msg.content += chunk.choices[0].delta.content;
options.stream({text: chunk.choices[0].delta.content}); options.stream({text: chunk.choices[0].delta.content});
} }
if(chunk.choices[0]?.delta?.tool_calls) {
if(chunk.choices[0].delta.tool_calls) {
for(const deltaTC of chunk.choices[0].delta.tool_calls) { for(const deltaTC of chunk.choices[0].delta.tool_calls) {
const existing = resp.choices[0].message.tool_calls.find(tc => tc.index === deltaTC.index); const existing = msg.tool_calls.find((tc: any) => tc.index === deltaTC.index);
if(existing) { if(existing) {
if(deltaTC.id) existing.id = deltaTC.id; if(deltaTC.id) existing.id = deltaTC.id;
if(deltaTC.type) existing.type = deltaTC.type; if(deltaTC.function?.name) existing.function.name = deltaTC.function.name;
if(deltaTC.function) { if(deltaTC.function?.arguments) existing.function.arguments += deltaTC.function.arguments;
if(!existing.function) existing.function = {};
if(deltaTC.function.name) existing.function.name = deltaTC.function.name;
if(deltaTC.function.arguments) existing.function.arguments = (existing.function.arguments || '') + deltaTC.function.arguments;
}
} else { } else {
resp.choices[0].message.tool_calls.push({ msg.tool_calls.push({
index: deltaTC.index, index: deltaTC.index,
id: deltaTC.id || '', id: deltaTC.id || '',
type: deltaTC.type || 'function', function: {name: deltaTC.function?.name || '', arguments: deltaTC.function?.arguments || ''}
function: {
name: deltaTC.function?.name || '',
arguments: deltaTC.function?.arguments || ''
}
}); });
} }
} }
} }
} }
} else {
usage = resp.usage;
msg = resp.choices[0].message;
} }
const duration = Date.now() - callStart;
const tps = usage?.completion_tokens && duration > 0 ? usage.completion_tokens / (duration / 1000) : 0;
const toolCalls = resp.choices[0].message.tool_calls || []; 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 = {role: 'tool', id: tc.id, name: tc.function.name, args: JSONAttemptParse(tc.function.arguments, {}), content: undefined, timestamp: Date.now()};
if(!tool) return {role: 'tool', tool_call_id: toolCall.id, content: '{"error": "Tool not found"}'}; history.push(entry);
return {tc, entry};
});
await Promise.all(entries.map(async ({tc, entry}: any) => {
const tool = tools.find(findByProp('name', tc.function.name));
if(options.stream) options.stream({tool: tc.function.name});
if(!tool) { entry.error = 'Tool not found'; return; }
try { try {
const args = JSONAttemptParse(toolCall.function.arguments, {}); const toolStream = options.stream && ((chunk: any) => {
const result = await tool.fn(args, options.stream, this.ai); if(chunk.done) { terminal = true; return; }
console.log(result); options.stream!(chunk);
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize(result)}; });
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
} catch(err: any) { } catch(err: any) {
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'})}; 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 || ''});
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()});
} }
} }

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
src/token-pool.ts Normal file
View File

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

View File

@@ -1,10 +1,12 @@
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} 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';
const UA = 'Mozilla/5.0 (Windows NT 10.0; Win64; x64)';
const getShell = () => { const getShell = () => {
if(os.platform() == 'win32') return 'cmd'; if(os.platform() == 'win32') return 'cmd';
return $Sync`echo $SHELL`?.split('/').pop() || 'bash'; return $Sync`echo $SHELL`?.split('/').pop() || 'bash';
@@ -39,21 +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 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}) => {
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 ExecPythonTool: AiTool = {
name: 'exec_python',
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 ExecTool: AiTool = { export const ExecTool: AiTool = {
@@ -67,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}`);
} }
@@ -81,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},
@@ -98,62 +637,160 @@ 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}) => {
console.log('executing js') function toFormUrlEncoded(obj, prefix = '') {
const c = consoleInterceptor(null); const pairs: any = [];
const resp = await Fn<any>({console: c}, args.code, true).catch((err: any) => c.output.error.push(err)); for (const key in obj) {
return {...c.output, return: resp, stdout: undefined, stderr: undefined}; 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)}`);
} }
} }
export const PythonTool: AiTool = { return pairs.join('&');
name: 'exec_javascript',
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 = { const res = await fetch(host + '/v1', {
name: 'read_webpage', method: 'POST',
description: 'Extract clean, structured content from a webpage. Use after web_search to read specific URLs', headers: {'Content-Type': 'application/json'},
args: { body: JSON.stringify({cmd, url, maxTimeout, postData: postData ? toFormUrlEncoded(postData) : undefined}),
url: {type: 'string', description: 'URL to extract content from', required: true}, });
focus: {type: 'string', description: 'Optional: What aspect to focus on (e.g., "pricing", "features", "contact info")'}
},
fn: async (args: {url: string; focus?: string}) => {
const html = await fetch(args.url, {headers: {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64)"}})
.then(r => r.text()).catch(err => {throw new Error(`Failed to fetch: ${err.message}`)});
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 WebReadTool: AiTool = {
name: 'web_read',
description: 'Extract clean content from webpages, or convert media/documents to accessible formats',
args: {
url: {type: 'string', description: 'URL to read', required: true},
mimeRegex: {type: 'string', description: 'Optional regex to filter MIME types (e.g., "^image/", "text/")'}
},
fn: async (args: {url: string; mimeRegex?: string}) => {
const ua = 'AiTools-Webpage/1.0';
const maxSize = 10 * 1024 * 1024;
const response = await fetch(args.url, {
headers: {
'User-Agent': ua,
'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,image/webp,*/*;q=0.8',
'Accept-Language': 'en-US,en;q=0.5'
},
redirect: 'follow'
}).catch(err => {throw new Error(`Failed to fetch: ${err.message}`)});
const contentType = response.headers.get('content-type') || '';
const mimeType = contentType.split(';')[0].trim().toLowerCase();
if(args.mimeRegex && !new RegExp(args.mimeRegex, 'i').test(mimeType)) {
return `❌ MIME type rejected: ${mimeType} (filter: ${args.mimeRegex})`;
}
if(mimeType.match(/^(image|audio|video)\//)) {
const buffer = await response.arrayBuffer();
if(buffer.byteLength > maxSize) {
return `❌ File too large: ${(buffer.byteLength / 1024 / 1024).toFixed(1)}MB (max 10MB)\nType: ${mimeType}`;
}
const base64 = Buffer.from(buffer).toString('base64');
return `## Media File\n**Type:** ${mimeType}\n**Size:** ${(buffer.byteLength / 1024).toFixed(1)}KB\n**Data URL:** \`data:${mimeType};base64,${base64.slice(0, 100)}...\``;
}
if(mimeType.match(/^text\/(plain|csv|xml)/) || args.url.match(/\.(txt|csv|xml|md|yaml|yml)$/i)) {
const text = await response.text();
const truncated = text.length > 50000 ? text.slice(0, 50000) : text;
return `## Text File\n**Type:** ${mimeType}\n**URL:** ${args.url}\n\n${truncated}`;
}
if(mimeType.match(/application\/(json|xml|csv)/)) {
const text = await response.text();
const truncated = text.length > 50000 ? text.slice(0, 50000) : text;
return `## Structured Data\n**Type:** ${mimeType}\n**URL:** ${args.url}\n\n\`\`\`\n${truncated}\n\`\`\``;
}
if(mimeType === 'application/pdf' || (mimeType.startsWith('application/') && !mimeType.includes('html'))) {
const buffer = await response.arrayBuffer();
if(buffer.byteLength > maxSize) {
return `❌ File too large: ${(buffer.byteLength / 1024 / 1024).toFixed(1)}MB (max 10MB)\nType: ${mimeType}`;
}
const base64 = Buffer.from(buffer).toString('base64');
return `## Binary File\n**Type:** ${mimeType}\n**Size:** ${(buffer.byteLength / 1024).toFixed(1)}KB\n**Data URL:** \`data:${mimeType};base64,${base64.slice(0, 100)}...\``;
}
// HTML
const html = await response.text();
const $ = cheerio.load(html); const $ = cheerio.load(html);
$('script, style, nav, footer, header, aside, iframe, noscript, [role="navigation"], [role="banner"], .ad, .ads, .cookie, .popup').remove(); $('script, style, nav, footer, header, aside, iframe, noscript, svg').remove();
const metadata = { $('[role="navigation"], [role="banner"], [role="complementary"]').remove();
title: $('meta[property="og:title"]').attr('content') || $('title').text() || '', $('[aria-hidden="true"], [hidden], .visually-hidden, .sr-only, .screen-reader-text').remove();
description: $('meta[name="description"]').attr('content') || $('meta[property="og:description"]').attr('content') || '', $('.ad, .ads, .advertisement, .cookie, .popup, .modal, .sidebar, .related, .comments, .social-share').remove();
}; $('button, [class*="share"], [class*="follow"], [class*="social"]').remove();
const title = $('meta[property="og:title"]').attr('content') || $('title').text().trim() || '';
const description = $('meta[name="description"]').attr('content') || $('meta[property="og:description"]').attr('content') || '';
const author = $('meta[name="author"]').attr('content') || '';
let content = ''; let content = '';
const contentSelectors = ['article', 'main', '[role="main"]', '.content', '.post', '.entry', 'body']; const selectors = ['article', 'main', '[role="main"]', '.content', '.post-content', '.entry-content', '.article-content'];
for (const selector of contentSelectors) { for(const sel of selectors) {
const el = $(selector).first(); const el = $(sel).first();
if(el.length && el.text().trim().length > 200) { if(el.length && el.text().trim().length > 200) {
content = el.text(); const paragraphs: string[] = [];
el.find('p').each((_, p) => {
const text = $(p).text().trim();
if(text.length > 80) paragraphs.push(text);
});
if(paragraphs.length > 2) {
content = paragraphs.join('\n\n');
break; break;
} }
} }
if (!content) content = $('body').text(); }
content = content.replace(/\s+/g, ' ').trim().slice(0, 8000);
return {url: args.url, title: metadata.title.trim(), description: metadata.description.trim(), content, focus: args.focus}; if(!content) {
const paragraphs: string[] = [];
$('body p').each((_, p) => {
const text = $(p).text().trim();
if(text.length > 80) paragraphs.push(text);
});
content = paragraphs.slice(0, 30).join('\n\n');
} }
// Decode escaped newlines and clean
const parts = [`## ${title || 'Webpage'}`];
if(description) parts.push(`_${description}_`);
if(author) parts.push(`👤 ${author}`);
parts.push(`🔗 ${args.url}\n`);
parts.push(content);
return decodeHtml(parts.join('\n\n').replaceAll(/\n{3,}/g, '\n\n'));
} }
};
export const WebSearchTool: AiTool = { export const WebSearchTool: AiTool = {
name: 'web_search', name: 'web_search',
@@ -167,7 +804,7 @@ export const WebSearchTool: AiTool = {
length: number; length: number;
}) => { }) => {
const html = await fetch(`https://html.duckduckgo.com/html/?q=${encodeURIComponent(args.query)}`, { const html = await fetch(`https://html.duckduckgo.com/html/?q=${encodeURIComponent(args.query)}`, {
headers: {"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64)", "Accept-Language": "en-US,en;q=0.9"} headers: {"User-Agent": UA, "Accept-Language": "en-US,en;q=0.9"}
}).then(resp => resp.text()); }).then(resp => resp.text());
let match, regex = /<a .*?href="(.+?)".+?<\/a>/g; let match, regex = /<a .*?href="(.+?)".+?<\/a>/g;
const results = new ASet<string>(); const results = new ASet<string>();

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

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 */