Compare commits

..

16 Commits
1.3.2 ... 1.5.0

Author SHA1 Message Date
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
11 changed files with 885 additions and 845 deletions

View File

@@ -119,7 +119,7 @@ const ai = new Ai({
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,
max_tokens: 100_000,
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},
@@ -186,7 +186,7 @@ console.log(chunks);
// Manually compile history into memories at end of conversation
// Happens automatically when coverstaions are compressed
await ai.language.updateMemory(history, memory);
await ai.language.memorize(history, memory);
// Summarize text
const summary = await ai.language.summarize(longText, 200);

View File

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

View File

@@ -1,65 +1,56 @@
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 {LLMMessage, LLMRequest} from './llm.ts';
import {LLMProvider} from './provider.ts';
import {TokenPool} from './token-pool.ts';
import {convertSchema} from './tools.ts';
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();
this.client = new anthropic({apiKey: apiToken});
this.tokenPool = new TokenPool(...makeArray(apiToken).filter(Boolean));
}
private toStandard(history: any[]): LLMMessage[] {
const timestamp = Date.now();
const messages: LLMMessage[] = [];
for(let h of history) {
if(typeof h.content == 'string') {
messages.push(<any>{timestamp, ...h});
} else {
const textContent = h.content?.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
if(textContent) messages.push({role: h.role, content: textContent, timestamp: timestamp});
h.content.forEach((c: any) => {
if(c.type == 'tool_use') {
messages.push({role: 'tool', id: c.id, name: c.name, args: c.input, timestamp: c.timestamp, 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;
private getClient(token: string): anthropic {
let client = this.clients.get(token);
if(!client) {
client = new anthropic({apiKey: token});
this.clients.set(token, client);
}
});
}
}
return messages;
return client;
}
private fromStandard(history: LLMMessage[]): any[] {
for(let i = 0; i < history.length; i++) {
if(history[i].role == 'tool') {
const h: any = history[i];
history.splice(i, 1,
/** 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: 'user', content: [{type: 'tool_result', tool_use_id: h.id, is_error: !!h.error, content: h.error || h.content}]}
)
i++;
{role: 'user', content: [{type: 'tool_result', tool_use_id: h.id, is_error: !!h.error, content: h.error || h.content || ''}]}
);
} else {
wire.push({role: h.role, content: h.content});
}
}
return history;
return wire;
}
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
const controller = new AbortController();
return Object.assign(new Promise<any>(async (res) => {
let history = this.fromStandard([
...(options.history || []).filter(h => h.role !== 'system'),
{role: 'user', content: message, timestamp: Date.now()}
]);
return Object.assign(new Promise<any>(async (res, rej) => {
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 requestParams: any = {
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 || '',
temperature: options.temperature || this.ai.options.llm?.temperature || undefined,
tools: tools.map(t => ({
@@ -69,94 +60,97 @@ export class Anthropic extends LLMProvider {
type: 'object',
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]) : []
},
fn: undefined
}
})),
messages: history,
stream: !!options.stream,
};
// Add structured output support
if(options.schema) {
requestParams.output_config = {
format: {
type: 'json_schema',
schema: convertSchema(options.schema)
}
};
requestParams.output_config = {format: {type: 'json_schema', schema: convertSchema(options.schema)}};
}
let resp: any, hasStreamedText = false, terminal = false;
try {
let terminal = false;
do {
requestParams.messages = history.map(({timestamp, ...m}) => m);
resp = await this.client.messages.create(requestParams).catch(err => {
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
requestParams.messages = this.toWire(history.filter(h => h.role !== 'system'));
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;
});
// Streaming mode
let usage: any, content: any[] = [];
if(options.stream) {
if(hasStreamedText) options.stream({text: '\n\n'});
resp.content = [];
for await (const chunk of resp) {
if(controller.signal.aborted) break;
if(chunk.type === 'content_block_start') {
if(chunk.content_block.type === 'text') {
resp.content.push({type: 'text', text: ''});
} 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>''});
}
if(chunk.content_block.type === 'text') 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.type === 'content_block_delta') {
if(chunk.delta.type === 'text_delta') {
const text = chunk.delta.text;
resp.content.at(-1).text += text;
if(text) { hasStreamedText = true; options.stream({text}); }
content.at(-1).text += chunk.delta.text;
options.stream({text: chunk.delta.text});
} 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') {
const last = resp.content.at(-1);
if(last?.input != null) last.input = last.input ? JSONAttemptParse(last.input, {}) : {};
const last = content.at(-1);
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') {
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 = resp.content.filter((c: any) => c.type === 'tool_use');
const toolCalls = content.filter((c: any) => c.type === 'tool_use');
if(toolCalls.length && !controller.signal.aborted) {
history.push({role: 'assistant', content: resp.content, timestamp: Date.now()});
const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
const tool = tools.find(findByProp('name', toolCall.name));
if(options.stream) options.stream({tool: toolCall.name});
if(!tool) return {tool_use_id: toolCall.id, is_error: true, content: 'Tool not found'};
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});
const entries = toolCalls.map((tc: any) => {
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 {
// Wrap stream so a tool's `done` ends turn gracefully
const toolStream = options.stream && ((chunk: any) => {
if(chunk.done) { terminal = true; return; }
options.stream!(chunk);
});
const result = await tool.fn(toolCall.input, toolStream, this.ai);
return {type: 'tool_result', tool_use_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result};
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
} catch(err: any) {
return {type: 'tool_result', tool_use_id: toolCall.id, is_error: true, content: err?.message || err?.toString() || 'Unknown'};
entry.error = err?.message || err?.toString() || 'Unknown';
}
}));
history.push({role: 'user', content: results, timestamp: Date.now()});
requestParams.messages = history;
} else {
terminal = true;
const text = content.filter((c: any) => c.type === 'text').map((c: any) => c.text).join('\n\n').trim();
if(text) history.push({role: 'assistant', content: text, timestamp: Date.now(), duration, tps});
}
} while (!terminal && !controller.signal.aborted && resp.content.some((c: any) => c.type === 'tool_use'));
} while(!terminal && !controller.signal.aborted);
if(!terminal) {
const textContent = resp.content.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
history.push({role: 'assistant', content: textContent, timestamp: Date.now()});
}
history = this.toStandard(history);
if(options.stream) options.stream({done: true});
if(options.history) options.history.splice(0, options.history.length, ...history);
const finalContent = 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()});
}
}

85
src/helpers.ts Normal file
View File

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

View File

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

View File

@@ -103,9 +103,10 @@ class BoundedMaxHeap<T> {
export class KDTree<T = unknown> {
private root: KDNode<T> | null = null;
private _size = 0;
private readonly dims: number;
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".

View File

@@ -1,4 +1,4 @@
import {snakeCase} from '@ztimson/utils';
import {clean, makeUnique, snakeCase} from '@ztimson/utils';
import {AbortablePromise, Ai} from './ai.ts';
import {Anthropic} from './antrhopic.ts';
import {OpenAi} from './open-ai.ts';
@@ -7,10 +7,12 @@ import {AiTool, AiToolArg} from './tools.ts';
import {fileURLToPath} from 'url';
import {dirname, join} from 'path';
import {spawn} from 'node:child_process';
import {Memory, MemoryCache, MemoryManager, MemoryOptions} from './memory.ts';
import {Memory, MemoryCache, MemoryManager, MemoryOptions, stripHeader} from './memory.ts';
export type AnthropicConfig = {proto: 'anthropic', token: string};
export type OpenAiConfig = {proto: 'openai', host?: string, token: string};
const MAX_AGENT_DEPTH = 5;
export type AnthropicConfig = {proto: 'anthropic', token: string | string[]};
export type OpenAiConfig = {proto: 'openai', host?: string, token: string | string[]};
export type Agent = {
name: string;
@@ -22,7 +24,6 @@ export type Agent = {
skills?: Skill[] | null;
tools?: AiTool[] | null;
mcp?: McpServer[] | null;
/** Explicit whitelist of agents this agent may delegate to. Default: none - must opt-in, self is always excluded */
agents?: string[] | null;
}
@@ -33,6 +34,10 @@ export type LLMMessage = {
content: string | any;
/** Timestamp */
timestamp?: number;
/** Response duration in ms */
duration?: number;
/** Tokens per second */
tps?: number;
} | {
/** Tool call */
role: 'tool';
@@ -48,6 +53,10 @@ export type LLMMessage = {
error?: undefined | string;
/** Timestamp */
timestamp?: number;
/** Response duration in ms */
duration?: number;
/** Tokens per second */
tps?: number;
}
export type LLMRequest = {
@@ -58,7 +67,7 @@ export type LLMRequest = {
/** Message history */
history?: LLMMessage[];
/** Max tokens for request */
max_tokens?: number;
maxTokens?: number;
/** 0 = Rigid Logic, 1 = Balanced, 2 = Hyper Creative **/
temperature?: number;
/** Available tools */
@@ -101,8 +110,6 @@ export type Skill = {
content: string;
}
const MAX_AGENT_DEPTH = 5;
class LLM {
private memoryManager!: MemoryManager;
@@ -119,40 +126,35 @@ class LLM {
this.memoryManager = new MemoryManager(this);
}
/**
* Wrap agents as tools. Nested delegation is opt-in only (empty by default, like
* tools/skills/mcp) and an agent can never call itself even if explicitly whitelisted.
* Delegate results are queued in `pending` and spliced into history by `ask()` after
* the provider's own end-of-turn history sync has already run.
*/
private setupAgent(agents: Agent[] = [], allAgents: Agent[], pending: Map<string, {resp: string, subHistory: LLMMessage[]}[]>, aborts: (() => void)[], depth = 0): AiTool[] {
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: {
context: {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true},
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) => {
}),
fn: async (args: any, stream: any, ai: any, id?: string) => {
if(depth >= MAX_AGENT_DEPTH) return 'Max agent delegation depth exceeded';
const subHistory: LLMMessage[] = [];
// Opt-in only, self always excluded regardless of whitelist
const nested = (a.agents || [])
.map(name => allAgents.find(x => x.name === name))
.filter((x): x is Agent => !!x && x.name !== a.name);
const request = this.ask(`${args.instructions}${args.context ? `\n\n<context>${args.context}</context>` : ''}`, {
system: `You are a specialized subagent. ${a.delegate ? 'Your output streams directly to the user for the remainder of this turn.' : 'You are wrapped in a tool call that will be analysis by an LLM'}
As a subagent, focus on executing your task completely using available tools and returning only the final result - no commentary, questions, or dialogue.
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: subHistory,
history: a.delegate ? history : [],
mcp: a.mcp || undefined,
skills: a.skills || undefined,
tools: a.tools || undefined,
@@ -163,8 +165,7 @@ ${a.system}`,
const resp = await request;
if(a.delegate) {
if(!pending.has(toolName)) pending.set(toolName, []);
pending.get(toolName)!.push({resp, subHistory});
delegateState.resp = resp;
return '';
}
return resp;
@@ -206,7 +207,7 @@ ${a.system}`,
const list = allTools.map(t => `- ${t.name}: ${t.description}`).join('\n');
return {
prompt: `You have access to the following MCP tools:\n${list}`,
prompt: `## MCP\nYou have access to the following MCP tools:\n${list}`,
tools: allTools
};
}
@@ -215,7 +216,7 @@ ${a.system}`,
if(!skills?.length) return {prompt: '', tools: []};
const list = skills.map(s => `- ${s.name}: ${s.description}`).join('\n');
return {
prompt: `You have access to the following skill documents, use \`read_skill\` to access them:\n${list}`,
prompt: `## Skills\nYou have access to the following skill documents, whenever there is overlap between a question and a skill file, use \`skill_read\` to get instructions and background knowledge:\n${list}`,
tools: [{
name: 'skill_read',
description: 'Read the full content of a skill/knowledge document',
@@ -231,6 +232,20 @@ ${a.system}`,
}
}
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> {
options = <any>{
system: '',
@@ -250,10 +265,14 @@ ${a.system}`,
nestedAborts.forEach(a => a());
};
const promise = (async () => {
let promise: any;
const requestStart = Date.now();
promise = (async () => {
let tools: AiTool[] = options.tools || this.ai.options.llm?.tools || [];
const prompts: string[] = [];
let history = options.history || [];
if(message) history.push({role: 'user', content: message, timestamp: Date.now()});
// MCP
const mcp = options.mcp || this.ai.options?.llm?.mcp;
@@ -273,8 +292,8 @@ ${a.system}`,
// Agents
const agents = options.agents || this.ai.options?.llm?.agents;
const pendingDelegates = new Map<string, {resp: string, subHistory: LLMMessage[]}[]>();
if(agents?.length) tools.push(...this.setupAgent(agents, agents, pendingDelegates, nestedAborts, options._agentDepth || 0));
const delegateState: {resp: string | null} = {resp: null};
if(agents?.length) tools.push(...this.setupAgent(agents, agents, history, nestedAborts, options._agentDepth || 0, delegateState));
// Memory
const mem = MemoryManager.normalize(options.memory);
@@ -297,18 +316,26 @@ ${a.system}`,
} else listed.push(r);
}
prompts.unshift(`You have access to the following memory files:
${mems.map(m => `- ${m.name}: ${m.description}`).join('\n')}
${preloaded.length ? `
Relevant memories have been preloaded:
${preloaded.map(r => `
**${r.name}**
${r.description}
${r.content}
`).join('\n---\n')}
` : ''}${listed.length ? `
Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.name).join(', ')}
` : ''}`.trim());
prompts.unshift(`## Memory
You have a background memory process which has prefetched relevant information${mem.update ? ' and will create new memories from this conversation' : ''} for you
Assume it is perfect and never mention this process to anyone ever
Always use your memories to craft a personalized response, they contain links / [[wiki links]] which you use navigate between them
${mem.tool ? `You can access memory files via the \`memory_search\` and \`memory_recall\` tools
When you need information about the user, \`memory_recall\` \`People/User\` before asking (fetch if not included bellow)
When you need information not provided, attempt 1-3 \`memory_search\` calls with distinct queries before asking` : ''}
${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));
}
@@ -316,53 +343,42 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam
if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
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(message, {...options, tools, system: prompts.filter(Boolean).join('\n\n')});
request = this.models[m].ask('', {...options, tools, system: prompts.filter(Boolean).join('\n\n')});
let resp = await request;
// Spice delegated agents response into history
let lastDelegateResp: string | null = null;
if(pendingDelegates.size) {
for(let i = 0; i < history.length; i++) {
const h = history[i];
if(h.role !== 'tool' || h.content !== '') continue;
const queue = pendingDelegates.get(h.name);
if(!queue?.length) continue;
const {resp: delegateResp, subHistory} = queue.shift()!;
const insert: LLMMessage[] = [...subHistory.filter(sh => sh.role === 'tool'), {role: 'assistant', content: delegateResp, timestamp: Date.now()}];
history.splice(i + 1, 0, ...insert);
lastDelegateResp = delegateResp;
i += insert.length;
}
// 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 the orchestrator added no commentary of its own, its answer IS the delegate's answer
if(typeof resp === 'string' && !resp.trim() && lastDelegateResp !== null) resp = lastDelegateResp;
if(typeof resp === 'string' && !resp.trim() && delegateState.resp !== null) resp = delegateState.resp;
// Trim memory injections from history
if(mem?.tool) history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'memory_recall'));
// Auto-memorize before compressing
if(options.compress && this.estimateTokens(history) >= options.compress.max) {
if(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});
}
/**
* Digest full conversation history into memory documents.
* Call on session end to persist the conversation.
*/
async updateMemory(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest = {}): Promise<Memory[]> {
return this.memoryManager.memorize(history, memories, {model: this.defaultModel, ...options});
}
/**
* Compress chat history to reduce context size
* @param {LLMMessage[]} history Chatlog that will be compressed
@@ -542,6 +558,14 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam
};
}
/**
* Digest full conversation history into memory documents.
* Call on session end to persist the conversation.
*/
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest = {}): Promise<Memory[]> {
return this.memoryManager.memorize(history, memories, {model: this.defaultModel, ...options});
}
/**
* Create a summary of some text
* @param {string} text Text to summarize

File diff suppressed because it is too large Load Diff

View File

@@ -1,86 +1,64 @@
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 {LLMMessage, LLMRequest} from './llm.ts';
import {LLMProvider} from './provider.ts';
import {TokenPool} from './token-pool.ts';
import {convertSchema} from './tools.ts';
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();
this.client = new openAI(clean({
baseURL: host,
apiKey: token || (host ? 'ignored' : undefined)
}));
const tokens = makeArray(token).filter(Boolean);
this.tokenPool = new TokenPool(...(tokens.length ? tokens : [host ? 'ignored' : '']));
}
private toStandard(history: any[]): LLMMessage[] {
for(let i = 0; i < history.length; i++) {
const h = history[i];
if(h.role === 'assistant' && h.tool_calls) {
const tools = h.tool_calls.map((tc: any) => ({
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') {
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 || '';
private getClient(token: string): openAI {
let client = this.clients.get(token);
if(!client) {
client = new openAI(clean({baseURL: this.host, apiKey: token || undefined}));
this.clients.set(token, client);
}
history.splice(i, 1);
i--;
}
if(!history[i]?.timestamp) history[i].timestamp = Date.now();
}
return history;
return client;
}
private fromStandard(history: LLMMessage[]): any[] {
return history.reduce((result, h) => {
/** 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') {
result.push({
wire.push({
role: 'assistant',
content: null,
tool_calls: [{id: h.id, type: 'function', function: {name: h.name, arguments: JSON.stringify(h.args)}}],
refusal: null,
annotations: [],
timestamp: h.timestamp,
}, {
role: 'tool',
tool_call_id: h.id,
content: h.error || h.content,
timestamp: h.timestamp,
content: h.error || h.content || '',
});
} else {
result.push(h);
wire.push({role: h.role, content: h.content});
}
return result;
}, [] as any[]);
}
return wire;
}
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
const controller = new AbortController();
return Object.assign(new Promise<any>(async (res, rej) => {
const base = (options.history || []).filter(h => h.role !== 'system');
let history = this.fromStandard([
...(options.system ? [{role: <any>'system', content: options.system, timestamp: Date.now()}] : []),
...base,
{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 requestParams: any = {
model: options.model || this.model,
messages: history,
stream: !!options.stream,
max_completion_tokens: options.max_tokens || this.ai.options.llm?.max_tokens || undefined,
max_completion_tokens: options.maxTokens || this.ai.options.llm?.maxTokens || undefined,
temperature: options.temperature || this.ai.options.llm?.temperature || undefined,
tools: tools.map(t => ({
type: 'function',
@@ -98,96 +76,94 @@ export class OpenAi extends LLMProvider {
if(options.schema) {
const schema = convertSchema(options.schema);
requestParams.response_format = {
type: 'json_schema',
json_schema: {
name: 'response',
strict: true,
schema
}
};
requestParams.response_format = {type: 'json_schema', json_schema: {name: 'response', strict: true, schema}};
}
if(options.stream) requestParams.stream_options = {include_usage: true};
let resp: any, hasStreamedText = false, terminal = false;
try {
let terminal = false;
do {
requestParams.messages = history.map(({timestamp, ...m}) => m);
resp = await this.client.chat.completions.create(requestParams).catch(err => {
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
requestParams.messages = this.toWire(history.filter(h => h.role !== 'system'), options.system);
const callStart = Date.now();
const resp: any = await this.tokenPool.run(token => this.getClient(token).chat.completions.create(requestParams)).catch(err => {
err.message += `\n\nMessages:\n${JSON.stringify(requestParams.messages, null, 2)}`;
throw err;
});
let usage: any, msg: any = {content: '', tool_calls: []};
if(options.stream) {
if(hasStreamedText) options.stream({text: '\n\n'});
resp.choices = [{message: {role: 'assistant', content: '', tool_calls: [], timestamp: Date.now()}}];
for await (const chunk of resp) {
if(controller.signal.aborted) break;
if(chunk.choices[0].delta.content) {
const text = chunk.choices[0].delta.content;
resp.choices[0].message.content += text;
if(text) { hasStreamedText = true; options.stream({text}); }
if(chunk.usage) usage = chunk.usage;
if(chunk.choices[0]?.delta?.content) {
msg.content += 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) {
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(deltaTC.id) existing.id = deltaTC.id;
if(deltaTC.type) existing.type = deltaTC.type;
if(deltaTC.function) {
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;
}
if(deltaTC.function?.name) existing.function.name = deltaTC.function.name;
if(deltaTC.function?.arguments) existing.function.arguments += deltaTC.function.arguments;
} else {
resp.choices[0].message.tool_calls.push({
msg.tool_calls.push({
index: deltaTC.index,
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;
if(resp.error) throw new Error(resp.error);
const toolCalls = resp.choices[0].message.tool_calls || [];
const toolCalls = msg.tool_calls || [];
if(toolCalls.length && !controller.signal.aborted) {
history.push(resp.choices[0].message);
const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
const tool = tools?.find(findByProp('name', toolCall.function.name));
if(options.stream) options.stream({tool: toolCall.function.name});
if(!tool) return {role: 'tool', tool_call_id: toolCall.id, content: '{"error": "Tool not found"}', timestamp: Date.now()};
if(msg.content?.trim()) history.push({role: 'assistant', content: msg.content.trim(), timestamp: Date.now(), duration, tps});
const entries = toolCalls.map((tc: any) => {
const entry: any = {role: 'tool', id: tc.id, name: tc.function.name, args: JSONAttemptParse(tc.function.arguments, {}), content: undefined, timestamp: Date.now()};
history.push(entry);
return {tc, entry};
});
await Promise.all(entries.map(async ({tc, entry}: any) => {
const tool = tools.find(findByProp('name', tc.function.name));
if(options.stream) options.stream({tool: tc.function.name});
if(!tool) { entry.error = 'Tool not found'; return; }
try {
const args = JSONAttemptParse(toolCall.function.arguments, {});
// Wrap stream so a tool's `done` ends turn gracefully
const toolStream = options.stream && ((chunk: any) => {
if(chunk.done) { terminal = true; return; }
options.stream!(chunk);
});
const result = await tool.fn(args, toolStream, this.ai);
return {role: 'tool', tool_call_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result, timestamp: Date.now()};
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
} catch(err: any) {
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'}), timestamp: Date.now()};
entry.error = err?.message || err?.toString() || 'Unknown';
}
}));
history.push(...results);
requestParams.messages = history;
} else {
terminal = true;
const text = (msg.content || '').trim();
if(text) history.push({role: 'assistant', content: text, timestamp: Date.now(), duration, tps});
}
} while (!terminal && !controller.signal.aborted && resp.choices?.[0]?.message?.tool_calls?.length);
} while(!terminal && !controller.signal.aborted);
if(!terminal) {
const textContent = resp.choices[0].message.content?.trim() || '';
history.push({role: 'assistant', content: textContent, timestamp: Date.now()});
}
history = this.toStandard(history);
if(options.history) options.history.splice(0, options.history.length, ...history.filter(h => h.role !== 'system'));
if(options.stream) options.stream({done: true});
const finalContent = 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()});
}
}

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

@@ -41,7 +41,7 @@ export type AiTool = {
/** Tool arguments */
args?: AiToolArg,
/** 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 function convertSchema(schema: any): any {