Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 497f051c62 | |||
| 62fbe73b22 | |||
| d53b1c6328 | |||
| 89619e211e | |||
| afc6653364 |
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@ztimson/ai-utils",
|
"name": "@ztimson/ai-utils",
|
||||||
"version": "1.3.1",
|
"version": "1.3.5",
|
||||||
"description": "AI Utility library",
|
"description": "AI Utility library",
|
||||||
"author": "Zak Timson",
|
"author": "Zak Timson",
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
|
|||||||
@@ -21,10 +21,10 @@ export class Anthropic extends LLMProvider {
|
|||||||
messages.push(<any>{timestamp, ...h});
|
messages.push(<any>{timestamp, ...h});
|
||||||
} else {
|
} else {
|
||||||
const textContent = h.content?.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
|
const textContent = h.content?.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
|
||||||
if(textContent) messages.push({role: h.role, content: textContent, timestamp: timestamp});
|
if(textContent) messages.push({role: h.role, content: textContent, timestamp: timestamp, duration: h.duration, tps: h.tps});
|
||||||
h.content.forEach((c: any) => {
|
h.content.forEach((c: any) => {
|
||||||
if(c.type == 'tool_use') {
|
if(c.type == 'tool_use') {
|
||||||
messages.push({role: 'tool', id: c.id, name: c.name, args: c.input, timestamp: c.timestamp, content: undefined});
|
messages.push({role: 'tool', id: c.id, name: c.name, args: c.input, timestamp: h.timestamp, content: undefined, duration: h.duration, tps: h.tps});
|
||||||
} else if(c.type == 'tool_result') {
|
} else if(c.type == 'tool_result') {
|
||||||
const m: any = messages.findLast(m => (<any>m).id == c.tool_use_id);
|
const m: any = messages.findLast(m => (<any>m).id == c.tool_use_id);
|
||||||
if(m) m[c.is_error ? 'error' : 'content'] = c.content;
|
if(m) m[c.is_error ? 'error' : 'content'] = c.content;
|
||||||
@@ -52,7 +52,10 @@ export class Anthropic extends LLMProvider {
|
|||||||
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
|
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
|
||||||
const controller = new AbortController();
|
const controller = new AbortController();
|
||||||
return Object.assign(new Promise<any>(async (res) => {
|
return Object.assign(new Promise<any>(async (res) => {
|
||||||
let history = this.fromStandard([...options.history || [], {role: 'user', content: message, timestamp: Date.now()}]);
|
let history = this.fromStandard([
|
||||||
|
...(options.history || []).filter(h => h.role !== 'system'),
|
||||||
|
{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,
|
||||||
@@ -83,18 +86,17 @@ export class Anthropic extends LLMProvider {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
let resp: any, isFirstMessage = true, terminal = false;
|
let resp: any, terminal = false, duration = 0, tps = 0;
|
||||||
do {
|
do {
|
||||||
requestParams.messages = history.map(({timestamp, ...m}) => m);
|
requestParams.messages = history.map(({timestamp, ...m}) => m);
|
||||||
|
const callStart = Date.now();
|
||||||
resp = await this.client.messages.create(requestParams).catch(err => {
|
resp = await this.client.messages.create(requestParams).catch(err => {
|
||||||
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
|
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
|
||||||
throw err;
|
throw err;
|
||||||
});
|
});
|
||||||
|
|
||||||
// Streaming mode
|
let usage: any;
|
||||||
if(options.stream) {
|
if(options.stream) {
|
||||||
if(!isFirstMessage) options.stream({text: '\n\n'});
|
|
||||||
else isFirstMessage = false;
|
|
||||||
resp.content = [];
|
resp.content = [];
|
||||||
for await (const chunk of resp) {
|
for await (const chunk of resp) {
|
||||||
if(controller.signal.aborted) break;
|
if(controller.signal.aborted) break;
|
||||||
@@ -115,27 +117,31 @@ export class Anthropic extends LLMProvider {
|
|||||||
} else if(chunk.type === 'content_block_stop') {
|
} else if(chunk.type === 'content_block_stop') {
|
||||||
const last = resp.content.at(-1);
|
const last = resp.content.at(-1);
|
||||||
if(last?.input != null) last.input = last.input ? JSONAttemptParse(last.input, {}) : {};
|
if(last?.input != null) last.input = last.input ? JSONAttemptParse(last.input, {}) : {};
|
||||||
|
} else if(chunk.type === 'message_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;
|
||||||
}
|
}
|
||||||
|
duration = Date.now() - callStart;
|
||||||
|
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 = 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, timestamp: Date.now()});
|
history.push({role: 'assistant', content: resp.content, timestamp: Date.now(), duration, tps});
|
||||||
const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
|
const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
|
||||||
const tool = tools.find(findByProp('name', toolCall.name));
|
const tool = tools.find(findByProp('name', toolCall.name));
|
||||||
if(options.stream) options.stream({tool: toolCall.name});
|
if(options.stream) options.stream({tool: toolCall.name});
|
||||||
if(!tool) return {tool_use_id: toolCall.id, is_error: true, content: 'Tool not found'};
|
if(!tool) return {tool_use_id: toolCall.id, is_error: true, content: 'Tool not found'};
|
||||||
try {
|
try {
|
||||||
// Wrap stream so a tool's `done` ends turn gracefully
|
|
||||||
const toolStream = options.stream && ((chunk: any) => {
|
const toolStream = options.stream && ((chunk: any) => {
|
||||||
if(chunk.done) { terminal = true; return; }
|
if(chunk.done) { terminal = true; return; }
|
||||||
options.stream!(chunk);
|
options.stream!(chunk);
|
||||||
});
|
});
|
||||||
const result = await tool.fn(toolCall.input, toolStream, this.ai);
|
const result = await tool.fn(toolCall.input, toolStream, this.ai, toolCall.id);
|
||||||
return {type: 'tool_result', tool_use_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result};
|
return {type: 'tool_result', tool_use_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result};
|
||||||
} catch (err: any) {
|
} catch (err: any) {
|
||||||
return {type: 'tool_result', tool_use_id: toolCall.id, is_error: true, content: err?.message || err?.toString() || 'Unknown'};
|
return {type: 'tool_result', tool_use_id: toolCall.id, is_error: true, content: err?.message || err?.toString() || 'Unknown'};
|
||||||
@@ -148,12 +154,19 @@ export class Anthropic extends LLMProvider {
|
|||||||
|
|
||||||
if(!terminal) {
|
if(!terminal) {
|
||||||
const textContent = resp.content.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
|
const textContent = resp.content.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n');
|
||||||
history.push({role: 'assistant', content: textContent, timestamp: Date.now()});
|
history.push({role: 'assistant', content: textContent.trim(), timestamp: Date.now(), duration, tps});
|
||||||
}
|
}
|
||||||
|
|
||||||
history = this.toStandard(history);
|
history = this.toStandard(history);
|
||||||
if(options.stream) options.stream({done: true});
|
|
||||||
if(options.history) options.history.splice(0, options.history.length, ...history);
|
if(options.history) options.history.splice(0, options.history.length, ...history);
|
||||||
const finalContent = history.at(-1)?.content;
|
if(options.stream) options.stream({done: true});
|
||||||
|
|
||||||
|
const turnStart = history.map(h => h.role).lastIndexOf('user');
|
||||||
|
const finalContent = history.slice(turnStart + 1).reduce((str, h) => {
|
||||||
|
if(h.role === 'assistant') return str + (h.content || '');
|
||||||
|
return str;
|
||||||
|
}, '').trim();
|
||||||
|
|
||||||
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
|
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
|
||||||
}), {abort: () => controller.abort()});
|
}), {abort: () => controller.abort()});
|
||||||
}
|
}
|
||||||
|
|||||||
63
src/llm.ts
63
src/llm.ts
@@ -22,7 +22,6 @@ export type Agent = {
|
|||||||
skills?: Skill[] | null;
|
skills?: Skill[] | null;
|
||||||
tools?: AiTool[] | null;
|
tools?: AiTool[] | null;
|
||||||
mcp?: McpServer[] | 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;
|
agents?: string[] | null;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -48,6 +47,10 @@ export type LLMMessage = {
|
|||||||
error?: undefined | string;
|
error?: undefined | string;
|
||||||
/** Timestamp */
|
/** Timestamp */
|
||||||
timestamp?: number;
|
timestamp?: number;
|
||||||
|
/** Response duration in ms */
|
||||||
|
duration?: number;
|
||||||
|
/** Tokens per second */
|
||||||
|
tps?: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
export type LLMRequest = {
|
export type LLMRequest = {
|
||||||
@@ -119,13 +122,7 @@ class LLM {
|
|||||||
this.memoryManager = new MemoryManager(this);
|
this.memoryManager = new MemoryManager(this);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
private setupAgent(agents: Agent[] = [], allAgents: Agent[], pending: Map<string, any>, aborts: (() => void)[], depth = 0): AiTool[] {
|
||||||
* Wrap agents as tools. Nested delegation is opt-in only (empty by default, like
|
|
||||||
* tools/skills/mcp) and an agent can never call itself even if explicitly whitelisted.
|
|
||||||
* Delegate results are queued in `pending` and spliced into history by `ask()` after
|
|
||||||
* the provider's own end-of-turn history sync has already run.
|
|
||||||
*/
|
|
||||||
private setupAgent(agents: Agent[] = [], allAgents: Agent[], pending: Map<string, {resp: string, subHistory: LLMMessage[]}[]>, aborts: (() => void)[], depth = 0): AiTool[] {
|
|
||||||
return agents.map(a => {
|
return agents.map(a => {
|
||||||
const toolName = `${a.delegate ? '' : 'sub'}agent_${snakeCase(a.name)}`;
|
const toolName = `${a.delegate ? '' : 'sub'}agent_${snakeCase(a.name)}`;
|
||||||
return {
|
return {
|
||||||
@@ -135,15 +132,16 @@ class LLM {
|
|||||||
context: {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true},
|
context: {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true},
|
||||||
instructions: {type: 'string', description: 'Detailed instructions for subagent to complete', required: true},
|
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';
|
if(depth >= MAX_AGENT_DEPTH) return 'Max agent delegation depth exceeded';
|
||||||
|
|
||||||
const subHistory: LLMMessage[] = [];
|
const subHistory: LLMMessage[] = [];
|
||||||
|
|
||||||
// Opt-in only, self always excluded regardless of whitelist
|
// Opt-in only, self always excluded regardless of whitelist
|
||||||
const nested = (a.agents || [])
|
const nested = (a.agents || [])
|
||||||
.map(name => allAgents.find(x => x.name === name))
|
.map(name => allAgents.find(x => x.name === name))
|
||||||
.filter((x): x is Agent => !!x && x.name !== a.name);
|
.filter((x): x is Agent => !!x && x.name !== a.name);
|
||||||
|
|
||||||
|
const start = Date.now();
|
||||||
const request = this.ask(`${args.instructions}${args.context ? `\n\n<context>${args.context}</context>` : ''}`, {
|
const request = this.ask(`${args.instructions}${args.context ? `\n\n<context>${args.context}</context>` : ''}`, {
|
||||||
system: `You are a specialized subagent. ${a.delegate ? 'Your output streams directly to the user for the remainder of this turn.' : 'You are wrapped in a tool call that will be analysis by an LLM'}
|
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.
|
As a subagent, focus on executing your task completely using available tools and returning only the final result - no commentary, questions, or dialogue.
|
||||||
@@ -161,10 +159,14 @@ ${a.system}`,
|
|||||||
} as any);
|
} as any);
|
||||||
aborts.push(request.abort);
|
aborts.push(request.abort);
|
||||||
const resp = await request;
|
const resp = await request;
|
||||||
|
const duration = Date.now() - start;
|
||||||
|
const assistantTurns = subHistory.filter((h: any) => h.role === 'assistant' && h.duration);
|
||||||
|
const genTime = assistantTurns.reduce((s, h: any) => s + h.duration, 0);
|
||||||
|
const genTokens = assistantTurns.reduce((s, h: any) => s + (h.tps || 0) * (h.duration / 1000), 0);
|
||||||
|
const tps = genTime > 0 ? genTokens / (genTime / 1000) : 0;
|
||||||
|
|
||||||
if(a.delegate) {
|
if(a.delegate) {
|
||||||
if(!pending.has(toolName)) pending.set(toolName, []);
|
pending.set(<string>id, {resp, subHistory, duration, tps});
|
||||||
pending.get(toolName)!.push({resp, subHistory});
|
|
||||||
return '';
|
return '';
|
||||||
}
|
}
|
||||||
return resp;
|
return resp;
|
||||||
@@ -231,6 +233,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> {
|
ask(message: string, options: LLMRequest = {}): AbortablePromise<string> {
|
||||||
options = <any>{
|
options = <any>{
|
||||||
system: '',
|
system: '',
|
||||||
@@ -273,7 +289,7 @@ ${a.system}`,
|
|||||||
|
|
||||||
// Agents
|
// Agents
|
||||||
const agents = options.agents || this.ai.options?.llm?.agents;
|
const agents = options.agents || this.ai.options?.llm?.agents;
|
||||||
const pendingDelegates = new Map<string, {resp: string, subHistory: LLMMessage[]}[]>();
|
const pendingDelegates = new Map<string, any>();
|
||||||
if(agents?.length) tools.push(...this.setupAgent(agents, agents, pendingDelegates, nestedAborts, options._agentDepth || 0));
|
if(agents?.length) tools.push(...this.setupAgent(agents, agents, pendingDelegates, nestedAborts, options._agentDepth || 0));
|
||||||
|
|
||||||
// Memory
|
// Memory
|
||||||
@@ -316,20 +332,29 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam
|
|||||||
|
|
||||||
if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
|
if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
|
||||||
|
|
||||||
|
// Time each tool call's real execution so its history entry gets its own duration/tps
|
||||||
|
const toolTimings = new Map<string, {duration: number, tps: number}>();
|
||||||
|
tools = this.wrapToolTiming(tools, toolTimings);
|
||||||
|
|
||||||
prompts.unshift(options.system || this.ai.options.llm?.system || '');
|
prompts.unshift(options.system || this.ai.options.llm?.system || '');
|
||||||
request = this.models[m].ask(message, {...options, tools, system: prompts.filter(Boolean).join('\n\n')});
|
request = this.models[m].ask(message, {...options, tools, system: prompts.filter(Boolean).join('\n\n')});
|
||||||
let resp = await request;
|
let resp = await request;
|
||||||
|
|
||||||
|
// Providers stamp duration/tps on assistant entries themselves (from real API usage).
|
||||||
|
// Overwrite tool entries with actual tool-execution timing instead of the LLM call timing.
|
||||||
|
for(const h of history) {
|
||||||
|
if(h.role === 'tool' && toolTimings.has(h.id)) Object.assign(h, toolTimings.get(h.id));
|
||||||
|
}
|
||||||
|
|
||||||
// Spice delegated agents response into history
|
// Spice delegated agents response into history
|
||||||
let lastDelegateResp: string | null = null;
|
let lastDelegateResp: string | null = null;
|
||||||
if(pendingDelegates.size) {
|
if(pendingDelegates.size) {
|
||||||
for(let i = 0; i < history.length; i++) {
|
for(let i = 0; i < history.length; i++) {
|
||||||
const h = history[i];
|
const h: any = history[i];
|
||||||
if(h.role !== 'tool' || h.content !== '') continue;
|
if(h.role !== 'tool' || !pendingDelegates.has(h.id)) continue;
|
||||||
const queue = pendingDelegates.get(h.name);
|
const {resp: delegateResp, subHistory, duration, tps} = pendingDelegates.get(h.id)!;
|
||||||
if(!queue?.length) continue;
|
pendingDelegates.delete(h.id);
|
||||||
const {resp: delegateResp, subHistory} = queue.shift()!;
|
const insert: LLMMessage[] = [...subHistory.filter(sh => sh.role === 'tool'), {role: 'assistant', content: delegateResp, timestamp: Date.now(), duration, tps}];
|
||||||
const insert: LLMMessage[] = [...subHistory.filter(sh => sh.role === 'tool'), {role: 'assistant', content: delegateResp, timestamp: Date.now()}];
|
|
||||||
history.splice(i + 1, 0, ...insert);
|
history.splice(i + 1, 0, ...insert);
|
||||||
lastDelegateResp = delegateResp;
|
lastDelegateResp = delegateResp;
|
||||||
i += insert.length;
|
i += insert.length;
|
||||||
|
|||||||
@@ -2,73 +2,6 @@ import {LLMRequest, LLMMessage} from './llm.ts';
|
|||||||
import {AiTool} from './tools.ts';
|
import {AiTool} from './tools.ts';
|
||||||
import {KDPoint, KDTree} from './kd-tree.ts';
|
import {KDPoint, KDTree} from './kd-tree.ts';
|
||||||
|
|
||||||
export function buildMemoryGraph(memories: Memory[] | MemoryCache): MemoryNode[] {
|
|
||||||
const mems = memories instanceof MemoryCache ? memories.memories : memories;
|
|
||||||
const nameSet = new Set(mems.map(m => m.name));
|
|
||||||
const ghosts = new Set<string>();
|
|
||||||
|
|
||||||
const nodes: MemoryNode[] = mems.map(m => {
|
|
||||||
const {links, backlinks} = extractMetadata(m.content);
|
|
||||||
return {
|
|
||||||
name: m.name,
|
|
||||||
missing: false,
|
|
||||||
links,
|
|
||||||
backlinks,
|
|
||||||
};
|
|
||||||
});
|
|
||||||
|
|
||||||
for (const node of nodes) {
|
|
||||||
for (const link of node.links) {
|
|
||||||
if (!nameSet.has(link)) ghosts.add(link);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return [
|
|
||||||
...nodes,
|
|
||||||
...[...ghosts].map(name => ({
|
|
||||||
name,
|
|
||||||
missing: true,
|
|
||||||
links: [],
|
|
||||||
backlinks: nodes
|
|
||||||
.filter(n => n.links.includes(name))
|
|
||||||
.map(n => n.name),
|
|
||||||
}))
|
|
||||||
];
|
|
||||||
}
|
|
||||||
|
|
||||||
export function renderMemoryGraph(nodes) {
|
|
||||||
if (!nodes.length) return 'No memories yet.';
|
|
||||||
|
|
||||||
const groups = new Map();
|
|
||||||
for (const node of nodes) {
|
|
||||||
const [prefix, ...rest] = node.name.split('/');
|
|
||||||
const group = rest.length ? prefix : 'Root';
|
|
||||||
const label = rest.length ? rest.join('/') : node.name;
|
|
||||||
if (!groups.has(group)) groups.set(group, []);
|
|
||||||
groups.get(group).push({...node, label});
|
|
||||||
}
|
|
||||||
|
|
||||||
const ghostCount = nodes.filter(n => n.missing).length;
|
|
||||||
const lines = [`Memory Graph (${nodes.length} nodes, ${ghostCount} ghost${ghostCount === 1 ? '' : 's'})`, ''];
|
|
||||||
|
|
||||||
for (const group of [...groups.keys()].sort()) {
|
|
||||||
const items = groups.get(group).sort((a, b) => a.label.localeCompare(b.label));
|
|
||||||
lines.push(`${group}/`);
|
|
||||||
items.forEach((n, i) => {
|
|
||||||
const last = i === items.length - 1;
|
|
||||||
const branch = last ? '└─' : '├─';
|
|
||||||
const pad = last ? ' ' : '│ ';
|
|
||||||
const tag = n.missing ? ' (ghost)' : '';
|
|
||||||
lines.push(` ${branch} ${n.label}${tag}`);
|
|
||||||
if (n.links.length) lines.push(` ${pad} → ${n.links.join(', ')}`);
|
|
||||||
if (n.backlinks.length) lines.push(` ${pad} ← ${n.backlinks.join(', ')}`);
|
|
||||||
});
|
|
||||||
lines.push('');
|
|
||||||
}
|
|
||||||
|
|
||||||
return lines.join('\n').trimEnd();
|
|
||||||
}
|
|
||||||
|
|
||||||
export class MemoryCache {
|
export class MemoryCache {
|
||||||
private tree: KDTree<MemoryRef>;
|
private tree: KDTree<MemoryRef>;
|
||||||
public memories: Memory[];
|
public memories: Memory[];
|
||||||
|
|||||||
@@ -20,15 +20,19 @@ export class OpenAi extends LLMProvider {
|
|||||||
for(let i = 0; i < history.length; i++) {
|
for(let i = 0; i < history.length; i++) {
|
||||||
const h = history[i];
|
const h = history[i];
|
||||||
if(h.role === 'assistant' && h.tool_calls) {
|
if(h.role === 'assistant' && h.tool_calls) {
|
||||||
const tools = h.tool_calls.map((tc: any) => ({
|
const items: any[] = [];
|
||||||
|
if(h.content) items.push({role: 'assistant', content: h.content, timestamp: h.timestamp, duration: h.duration, tps: h.tps});
|
||||||
|
items.push(...h.tool_calls.map((tc: any) => ({
|
||||||
role: 'tool',
|
role: 'tool',
|
||||||
id: tc.id,
|
id: tc.id,
|
||||||
name: tc.function.name,
|
name: tc.function.name,
|
||||||
args: JSONAttemptParse(tc.function.arguments, {}),
|
args: JSONAttemptParse(tc.function.arguments, {}),
|
||||||
timestamp: h.timestamp
|
timestamp: h.timestamp,
|
||||||
}));
|
duration: h.duration,
|
||||||
history.splice(i, 1, ...tools);
|
tps: h.tps
|
||||||
i += tools.length - 1;
|
})));
|
||||||
|
history.splice(i, 1, ...items);
|
||||||
|
i += items.length - 1;
|
||||||
} else if(h.role === 'tool') {
|
} else if(h.role === 'tool') {
|
||||||
const record = history.find(h2 => h.tool_call_id == h2.id);
|
const record = history.find(h2 => h.tool_call_id == h2.id);
|
||||||
if(record) {
|
if(record) {
|
||||||
@@ -69,11 +73,12 @@ export class OpenAi extends LLMProvider {
|
|||||||
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
|
ask(message: string, options: LLMRequest = {}): AbortablePromise<string | any> {
|
||||||
const controller = new AbortController();
|
const controller = new AbortController();
|
||||||
return Object.assign(new Promise<any>(async (res, rej) => {
|
return Object.assign(new Promise<any>(async (res, rej) => {
|
||||||
if(options.system) {
|
const base = (options.history || []).filter(h => h.role !== 'system');
|
||||||
if(options.history?.[0]?.role != 'system') options.history?.splice(0, 0, {role: 'system', content: options.system, timestamp: Date.now()});
|
let history = this.fromStandard([
|
||||||
else options.history[0].content = options.system;
|
...(options.system ? [{role: <any>'system', content: options.system, timestamp: Date.now()}] : []),
|
||||||
}
|
...base,
|
||||||
let history = this.fromStandard([...options.history || [], {role: 'user', content: message, timestamp: Date.now()}]);
|
{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,
|
||||||
@@ -107,26 +112,27 @@ export class OpenAi extends LLMProvider {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
let resp: any, isFirstMessage = true, terminal = false;
|
if(options.stream) requestParams.stream_options = {include_usage: true};
|
||||||
|
let resp: any, terminal = false, duration = 0, tps = 0;
|
||||||
do {
|
do {
|
||||||
requestParams.messages = history.map(({timestamp, ...m}) => m);
|
requestParams.messages = history.map(({timestamp, ...m}) => m);
|
||||||
|
const callStart = Date.now();
|
||||||
resp = await this.client.chat.completions.create(requestParams).catch(err => {
|
resp = await this.client.chat.completions.create(requestParams).catch(err => {
|
||||||
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
|
err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`;
|
||||||
throw err;
|
throw err;
|
||||||
});
|
});
|
||||||
|
|
||||||
|
let usage: any;
|
||||||
if(options.stream) {
|
if(options.stream) {
|
||||||
if(!isFirstMessage) options.stream({text: '\n\n'});
|
|
||||||
else isFirstMessage = false;
|
|
||||||
resp.choices = [{message: {role: 'assistant', content: '', tool_calls: [], timestamp: Date.now()}}];
|
resp.choices = [{message: {role: 'assistant', content: '', tool_calls: [], timestamp: Date.now()}}];
|
||||||
for await (const chunk of resp) {
|
for await (const chunk of resp) {
|
||||||
if(controller.signal.aborted) break;
|
if(controller.signal.aborted) break;
|
||||||
if(chunk.choices[0].delta.content) {
|
if(chunk.usage) usage = chunk.usage;
|
||||||
|
if(chunk.choices[0]?.delta?.content) {
|
||||||
resp.choices[0].message.content += chunk.choices[0].delta.content;
|
resp.choices[0].message.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 = resp.choices[0].message.tool_calls.find(tc => tc.index === deltaTC.index);
|
||||||
if(existing) {
|
if(existing) {
|
||||||
@@ -151,24 +157,27 @@ export class OpenAi extends LLMProvider {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
} else {
|
||||||
|
usage = resp.usage;
|
||||||
}
|
}
|
||||||
|
duration = Date.now() - callStart;
|
||||||
|
tps = usage?.completion_tokens && duration > 0 ? usage.completion_tokens / (duration / 1000) : 0;
|
||||||
|
|
||||||
if(resp.error) throw new Error(resp.error);
|
if(resp.error) throw new Error(resp.error);
|
||||||
const toolCalls = resp.choices[0].message.tool_calls || [];
|
const toolCalls = resp.choices[0].message.tool_calls || [];
|
||||||
if(toolCalls.length && !controller.signal.aborted) {
|
if(toolCalls.length && !controller.signal.aborted) {
|
||||||
history.push(resp.choices[0].message);
|
history.push({...resp.choices[0].message, duration, tps});
|
||||||
const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
|
const results = await Promise.all(toolCalls.map(async (toolCall: any) => {
|
||||||
const tool = tools?.find(findByProp('name', toolCall.function.name));
|
const tool = tools?.find(findByProp('name', toolCall.function.name));
|
||||||
if(options.stream) options.stream({tool: toolCall.function.name});
|
if(options.stream) options.stream({tool: toolCall.function.name});
|
||||||
if(!tool) return {role: 'tool', tool_call_id: toolCall.id, content: '{"error": "Tool not found"}', timestamp: Date.now()};
|
if(!tool) return {role: 'tool', tool_call_id: toolCall.id, content: '{"error": "Tool not found"}', timestamp: Date.now()};
|
||||||
try {
|
try {
|
||||||
const args = JSONAttemptParse(toolCall.function.arguments, {});
|
const args = JSONAttemptParse(toolCall.function.arguments, {});
|
||||||
// Wrap stream so a tool's `done` ends turn gracefully
|
|
||||||
const toolStream = options.stream && ((chunk: any) => {
|
const toolStream = options.stream && ((chunk: any) => {
|
||||||
if(chunk.done) { terminal = true; return; }
|
if(chunk.done) { terminal = true; return; }
|
||||||
options.stream!(chunk);
|
options.stream!(chunk);
|
||||||
});
|
});
|
||||||
const result = await tool.fn(args, toolStream, this.ai);
|
const result = await tool.fn(args, toolStream, this.ai, toolCall.id);
|
||||||
return {role: 'tool', tool_call_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result, timestamp: Date.now()};
|
return {role: 'tool', tool_call_id: toolCall.id, content: typeof result == 'object' ? JSONSanitize(result) : result, timestamp: Date.now()};
|
||||||
} catch (err: any) {
|
} catch (err: any) {
|
||||||
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'}), timestamp: Date.now()};
|
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'}), timestamp: Date.now()};
|
||||||
@@ -180,13 +189,20 @@ export class OpenAi extends LLMProvider {
|
|||||||
} while (!terminal && !controller.signal.aborted && resp.choices?.[0]?.message?.tool_calls?.length);
|
} while (!terminal && !controller.signal.aborted && resp.choices?.[0]?.message?.tool_calls?.length);
|
||||||
|
|
||||||
if(!terminal) {
|
if(!terminal) {
|
||||||
const textContent = resp.choices[0].message.content?.trim() || '';
|
const textContent = resp.choices[0].message.content || '';
|
||||||
history.push({role: 'assistant', content: textContent, timestamp: Date.now()});
|
history.push({role: 'assistant', content: textContent.trim(), timestamp: Date.now(), duration, tps});
|
||||||
}
|
}
|
||||||
|
|
||||||
history = this.toStandard(history);
|
history = this.toStandard(history);
|
||||||
|
if(options.history) options.history.splice(0, options.history.length, ...history.filter(h => h.role !== 'system'));
|
||||||
if(options.stream) options.stream({done: true});
|
if(options.stream) options.stream({done: true});
|
||||||
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) => {
|
||||||
|
if(h.role === 'assistant') return str + (h.content || '');
|
||||||
|
return str;
|
||||||
|
}, '').trim();
|
||||||
|
|
||||||
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
|
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
|
||||||
}), {abort: () => controller.abort()});
|
}), {abort: () => controller.abort()});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -41,7 +41,7 @@ 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 function convertSchema(schema: any): any {
|
export function convertSchema(schema: any): any {
|
||||||
|
|||||||
Reference in New Issue
Block a user