diff --git a/src/antrhopic.ts b/src/antrhopic.ts index 4c44377..0079396 100644 --- a/src/antrhopic.ts +++ b/src/antrhopic.ts @@ -21,10 +21,10 @@ export class Anthropic extends LLMProvider { messages.push({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}); + if(textContent) messages.push({role: h.role, content: textContent, timestamp: timestamp, duration: h.duration, tps: h.tps}); h.content.forEach((c: any) => { if(c.type == 'tool_use') { - messages.push({role: 'tool', id: c.id, name: c.name, args: c.input, timestamp: 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') { const m: any = messages.findLast(m => (m).id == c.tool_use_id); if(m) m[c.is_error ? 'error' : 'content'] = c.content; @@ -86,15 +86,16 @@ export class Anthropic extends LLMProvider { }; } - let resp: any, terminal = false; + let resp: any, terminal = false, duration = 0, tps = 0; do { requestParams.messages = history.map(({timestamp, ...m}) => m); + const callStart = Date.now(); resp = await this.client.messages.create(requestParams).catch(err => { err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`; throw err; }); - // Streaming mode + let usage: any; if(options.stream) { resp.content = []; for await (const chunk of resp) { @@ -116,22 +117,26 @@ export class Anthropic extends LLMProvider { } else if(chunk.type === 'content_block_stop') { const last = resp.content.at(-1); 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') { 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'); 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 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'}; 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); @@ -149,8 +154,9 @@ export class Anthropic extends LLMProvider { if(!terminal) { const textContent = resp.content.filter((c: any) => c.type == 'text').map((c: any) => c.text).join('\n\n'); - history.push({role: 'assistant', content: textContent.trim(), timestamp: Date.now()}); + history.push({role: 'assistant', content: textContent.trim(), timestamp: Date.now(), duration, tps}); } + history = this.toStandard(history); if(options.history) options.history.splice(0, options.history.length, ...history); if(options.stream) options.stream({done: true}); diff --git a/src/llm.ts b/src/llm.ts index d3f0987..1ddebc8 100644 --- a/src/llm.ts +++ b/src/llm.ts @@ -47,6 +47,10 @@ export type LLMMessage = { error?: undefined | string; /** Timestamp */ timestamp?: number; + /** Response duration in ms */ + duration?: number; + /** Tokens per second */ + tps?: number; } export type LLMRequest = { @@ -118,12 +122,6 @@ 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, aborts: (() => void)[], depth = 0): AiTool[] { return agents.map(a => { const toolName = `${a.delegate ? '' : 'sub'}agent_${snakeCase(a.name)}`; @@ -143,6 +141,7 @@ class LLM { .map(name => allAgents.find(x => x.name === 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${args.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. @@ -160,9 +159,14 @@ ${a.system}`, } as any); aborts.push(request.abort); 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) { - pending.set(id, {resp, subHistory}); + pending.set(id, {resp, subHistory, duration, tps}); return ''; } return resp; @@ -229,6 +233,20 @@ ${a.system}`, } } + private wrapToolTiming(tools: AiTool[], timings: Map): 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 { options = { system: '', @@ -314,19 +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'}); + // Time each tool call's real execution so its history entry gets its own duration/tps + const toolTimings = new Map(); + tools = this.wrapToolTiming(tools, toolTimings); + prompts.unshift(options.system || this.ai.options.llm?.system || ''); request = this.models[m].ask(message, {...options, tools, system: prompts.filter(Boolean).join('\n\n')}); 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 let lastDelegateResp: string | null = null; if(pendingDelegates.size) { for(let i = 0; i < history.length; i++) { const h: any = history[i]; if(h.role !== 'tool' || !pendingDelegates.has(h.id)) continue; - const {resp: delegateResp, subHistory} = pendingDelegates.get(h.id)!; + const {resp: delegateResp, subHistory, duration, tps} = pendingDelegates.get(h.id)!; pendingDelegates.delete(h.id); - const insert: LLMMessage[] = [...subHistory.filter(sh => sh.role === 'tool'), {role: 'assistant', content: delegateResp, timestamp: Date.now()}]; + const insert: LLMMessage[] = [...subHistory.filter(sh => sh.role === 'tool'), {role: 'assistant', content: delegateResp, timestamp: Date.now(), duration, tps}]; history.splice(i + 1, 0, ...insert); lastDelegateResp = delegateResp; i += insert.length; diff --git a/src/memory.ts b/src/memory.ts index ae96d98..a5b94d5 100644 --- a/src/memory.ts +++ b/src/memory.ts @@ -2,73 +2,6 @@ import {LLMRequest, LLMMessage} from './llm.ts'; import {AiTool} from './tools.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(); - - 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 { private tree: KDTree; public memories: Memory[]; diff --git a/src/open-ai.ts b/src/open-ai.ts index 5c5e5c5..1606641 100644 --- a/src/open-ai.ts +++ b/src/open-ai.ts @@ -21,13 +21,15 @@ export class OpenAi extends LLMProvider { const h = history[i]; if(h.role === 'assistant' && h.tool_calls) { const items: any[] = []; - if(h.content) items.push({role: 'assistant', content: h.content, timestamp: h.timestamp}); + if(h.content) items.push({role: 'assistant', content: h.content, timestamp: h.timestamp, duration: h.duration, tps: h.tps}); items.push(...h.tool_calls.map((tc: any) => ({ role: 'tool', id: tc.id, name: tc.function.name, args: JSONAttemptParse(tc.function.arguments, {}), - timestamp: h.timestamp + timestamp: h.timestamp, + duration: h.duration, + tps: h.tps }))); history.splice(i, 1, ...items); i += items.length - 1; @@ -110,23 +112,27 @@ export class OpenAi extends LLMProvider { }; } - let resp: any, terminal = false; + if(options.stream) requestParams.stream_options = {include_usage: true}; + let resp: any, terminal = false, duration = 0, tps = 0; do { requestParams.messages = history.map(({timestamp, ...m}) => m); + const callStart = Date.now(); resp = await this.client.chat.completions.create(requestParams).catch(err => { err.message += `\n\nMessages:\n${JSON.stringify(history, null, 2)}`; throw err; }); + let usage: any; if(options.stream) { 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) { + if(chunk.usage) usage = chunk.usage; + if(chunk.choices[0]?.delta?.content) { resp.choices[0].message.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); if(existing) { @@ -151,19 +157,22 @@ 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); const toolCalls = resp.choices[0].message.tool_calls || []; 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 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()}; 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); @@ -181,8 +190,9 @@ export class OpenAi extends LLMProvider { if(!terminal) { const textContent = resp.choices[0].message.content || ''; - history.push({role: 'assistant', content: textContent.trim(), timestamp: Date.now()}); + history.push({role: 'assistant', content: textContent.trim(), timestamp: Date.now(), duration, tps}); } + history = this.toStandard(history); if(options.history) options.history.splice(0, options.history.length, ...history.filter(h => h.role !== 'system')); if(options.stream) options.stream({done: true});