diff --git a/package.json b/package.json
index adc8374..752c612 100644
--- a/package.json
+++ b/package.json
@@ -1,6 +1,6 @@
{
"name": "@ztimson/ai-utils",
- "version": "1.3.2",
+ "version": "1.3.3",
"description": "AI Utility library",
"author": "Zak Timson",
"license": "MIT",
diff --git a/src/antrhopic.ts b/src/antrhopic.ts
index 488b41a..57554dd 100644
--- a/src/antrhopic.ts
+++ b/src/antrhopic.ts
@@ -86,7 +86,7 @@ export class Anthropic extends LLMProvider {
};
}
- let resp: any, hasStreamedText = false, terminal = false;
+ let resp: any, terminal = false;
do {
requestParams.messages = history.map(({timestamp, ...m}) => m);
resp = await this.client.messages.create(requestParams).catch(err => {
@@ -96,7 +96,6 @@ export class Anthropic extends LLMProvider {
// Streaming mode
if(options.stream) {
- if(hasStreamedText) options.stream({text: '\n\n'});
resp.content = [];
for await (const chunk of resp) {
if(controller.signal.aborted) break;
@@ -110,7 +109,7 @@ export class Anthropic extends LLMProvider {
if(chunk.delta.type === 'text_delta') {
const text = chunk.delta.text;
resp.content.at(-1).text += text;
- if(text) { hasStreamedText = true; options.stream({text}); }
+ options.stream({text});
} else if(chunk.delta.type === 'input_json_delta') {
resp.content.at(-1).input += chunk.delta.partial_json;
}
@@ -137,7 +136,7 @@ export class Anthropic extends LLMProvider {
if(chunk.done) { terminal = true; return; }
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};
} catch (err: any) {
return {type: 'tool_result', tool_use_id: toolCall.id, is_error: true, content: err?.message || err?.toString() || 'Unknown'};
@@ -150,12 +149,19 @@ 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, timestamp: Date.now()});
+ history.push({role: 'assistant', content: textContent.trim(), 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;
+ 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 || '');
+ if(h.role === 'tool') return str + `${h.name}\n\n`;
+ return str;
+ }, '').trim();
+
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
}), {abort: () => controller.abort()});
}
diff --git a/src/llm.ts b/src/llm.ts
index 015c707..d3f0987 100644
--- a/src/llm.ts
+++ b/src/llm.ts
@@ -22,7 +22,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;
}
@@ -125,7 +124,7 @@ class LLM {
* 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[] {
+ 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)}`;
return {
@@ -135,10 +134,10 @@ class LLM {
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},
},
- 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))
@@ -163,8 +162,7 @@ ${a.system}`,
const resp = await request;
if(a.delegate) {
- if(!pending.has(toolName)) pending.set(toolName, []);
- pending.get(toolName)!.push({resp, subHistory});
+ pending.set(id, {resp, subHistory});
return '';
}
return resp;
@@ -273,7 +271,7 @@ ${a.system}`,
// Agents
const agents = options.agents || this.ai.options?.llm?.agents;
- const pendingDelegates = new Map();
+ const pendingDelegates = new Map();
if(agents?.length) tools.push(...this.setupAgent(agents, agents, pendingDelegates, nestedAborts, options._agentDepth || 0));
// Memory
@@ -324,11 +322,10 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam
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 h: any = history[i];
+ if(h.role !== 'tool' || !pendingDelegates.has(h.id)) continue;
+ const {resp: delegateResp, subHistory} = pendingDelegates.get(h.id)!;
+ pendingDelegates.delete(h.id);
const insert: LLMMessage[] = [...subHistory.filter(sh => sh.role === 'tool'), {role: 'assistant', content: delegateResp, timestamp: Date.now()}];
history.splice(i + 1, 0, ...insert);
lastDelegateResp = delegateResp;
diff --git a/src/open-ai.ts b/src/open-ai.ts
index 25b6934..e133524 100644
--- a/src/open-ai.ts
+++ b/src/open-ai.ts
@@ -20,15 +20,17 @@ export class OpenAi extends LLMProvider {
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) => ({
+ const items: any[] = [];
+ if(h.content) items.push({role: 'assistant', content: h.content, timestamp: h.timestamp});
+ 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
- }));
- history.splice(i, 1, ...tools);
- i += tools.length - 1;
+ })));
+ history.splice(i, 1, ...items);
+ i += items.length - 1;
} else if(h.role === 'tool') {
const record = history.find(h2 => h.tool_call_id == h2.id);
if(record) {
@@ -108,7 +110,7 @@ export class OpenAi extends LLMProvider {
};
}
- let resp: any, hasStreamedText = false, terminal = false;
+ let resp: any, terminal = false;
do {
requestParams.messages = history.map(({timestamp, ...m}) => m);
resp = await this.client.chat.completions.create(requestParams).catch(err => {
@@ -117,14 +119,12 @@ export class OpenAi extends LLMProvider {
});
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}); }
+ 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) {
for(const deltaTC of chunk.choices[0].delta.tool_calls) {
@@ -168,7 +168,7 @@ export class OpenAi extends LLMProvider {
if(chunk.done) { terminal = true; return; }
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()};
} catch (err: any) {
return {role: 'tool', tool_call_id: toolCall.id, content: JSONSanitize({error: err?.message || err?.toString() || 'Unknown'}), timestamp: Date.now()};
@@ -180,13 +180,20 @@ export class OpenAi extends LLMProvider {
} while (!terminal && !controller.signal.aborted && resp.choices?.[0]?.message?.tool_calls?.length);
if(!terminal) {
- const textContent = resp.choices[0].message.content?.trim() || '';
- history.push({role: 'assistant', content: textContent, timestamp: Date.now()});
+ const textContent = resp.choices[0].message.content || '';
+ history.push({role: 'assistant', content: textContent.trim(), 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) => {
+ if(h.role === 'assistant') return str + (h.content || '');
+ if(h.role === 'tool') return str + `${h.name}\n\n`;
+ return str;
+ }, '').trim();
+
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
}), {abort: () => controller.abort()});
}
diff --git a/src/tools.ts b/src/tools.ts
index 79f6f9c..cb3e248 100644
--- a/src/tools.ts
+++ b/src/tools.ts
@@ -41,7 +41,7 @@ export type AiTool = {
/** Tool arguments */
args?: AiToolArg,
/** Callback function */
- fn: (args: any, stream: LLMRequest['stream'], ai: Ai) => any | Promise,
+ fn: (args: any, stream: LLMRequest['stream'], ai: Ai, toolId?: string) => any | Promise,
};
export function convertSchema(schema: any): any {