Compare commits

..

3 Commits
1.3.0 ... 1.3.3

Author SHA1 Message Date
89619e211e Fixed message history and response
All checks were successful
Publish Library / Build NPM Project (push) Successful in 59s
Publish Library / Tag Version (push) Successful in 22s
2026-08-03 19:30:39 -04:00
afc6653364 fixed openai system calls in history breaking anthropic calls
All checks were successful
Publish Library / Build NPM Project (push) Successful in 55s
Publish Library / Tag Version (push) Successful in 21s
2026-08-02 22:35:17 -04:00
68e72445a2 Keep recent memories in context
All checks were successful
Publish Library / Build NPM Project (push) Successful in 53s
Publish Library / Tag Version (push) Successful in 17s
2026-08-01 21:42:05 -04:00
6 changed files with 183 additions and 153 deletions

View File

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

View File

@@ -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,7 +86,7 @@ export class Anthropic extends LLMProvider {
}; };
} }
let resp: any, isFirstMessage = true, terminal = false; let resp: any, terminal = false;
do { do {
requestParams.messages = history.map(({timestamp, ...m}) => m); requestParams.messages = history.map(({timestamp, ...m}) => m);
resp = await this.client.messages.create(requestParams).catch(err => { resp = await this.client.messages.create(requestParams).catch(err => {
@@ -93,8 +96,6 @@ export class Anthropic extends LLMProvider {
// Streaming mode // Streaming mode
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;
@@ -135,7 +136,7 @@ export class Anthropic extends LLMProvider {
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 +149,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()});
} }
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 || '');
if(h.role === 'tool') return str + `<tool>${h.name}</tool>\n\n`;
return str;
}, '').trim();
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent); res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
}), {abort: () => controller.abort()}); }), {abort: () => controller.abort()});
} }

View File

@@ -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;
} }
@@ -125,7 +124,7 @@ class LLM {
* Delegate results are queued in `pending` and spliced into history by `ask()` after * 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. * 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[], pending: Map<string, any>, 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,10 +134,10 @@ 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))
@@ -163,8 +162,7 @@ ${a.system}`,
const resp = await request; const resp = await request;
if(a.delegate) { if(a.delegate) {
if(!pending.has(toolName)) pending.set(toolName, []); pending.set(<string>id, {resp, subHistory});
pending.get(toolName)!.push({resp, subHistory});
return ''; return '';
} }
return resp; return resp;
@@ -273,7 +271,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
@@ -324,11 +322,10 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam
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} = 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()}]; 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;
@@ -340,7 +337,7 @@ Also relevant but not preloaded (use \`memory_recall\`): ${listed.map(r => r.nam
if(typeof resp === 'string' && !resp.trim() && lastDelegateResp !== null) resp = lastDelegateResp; if(typeof resp === 'string' && !resp.trim() && lastDelegateResp !== null) resp = lastDelegateResp;
// Trim memory injections from history // Trim memory injections from history
if(mem?.tool) history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'recall')); if(mem?.tool) history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'memory_recall'));
// Auto-memorize before compressing // Auto-memorize before compressing
if(options.compress && this.estimateTokens(history) >= options.compress.max) { if(options.compress && this.estimateTokens(history) >= options.compress.max) {

View File

@@ -221,6 +221,8 @@ function getWeekSunday(monday: string): string {
} }
export class MemoryManager { export class MemoryManager {
private recentlyTouched = new Map<string, number>();
private pendingMemorizations = new Map<string, { private pendingMemorizations = new Map<string, {
memories: Memory[] | MemoryCache, memories: Memory[] | MemoryCache,
tempMemoryName: string, tempMemoryName: string,
@@ -244,6 +246,7 @@ export class MemoryManager {
const mems = memories instanceof MemoryCache ? memories.memories : memories; const mems = memories instanceof MemoryCache ? memories.memories : memories;
const mem = mems.find(m => m.name === args.name); const mem = mems.find(m => m.name === args.name);
if (!mem) return 'Document not found'; if (!mem) return 'Document not found';
this.touch(mem.name);
return mem.content; return mem.content;
}, },
}), }),
@@ -292,6 +295,96 @@ ${conversation}`;
}; };
} }
private applyHeader(content: string, header: string): string {
return `${header}\n\n${this.stripHeader(content)}`;
}
private async backgroundMemorization(conversation: string, memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string): Promise<void> {
const mem = memories instanceof MemoryCache ? memories.memories : memories;
const monday = getWeekMonday();
const sunday = getWeekSunday(monday);
const buckets = await this.factAgent(conversation, mem, options, monday);
if(!buckets.length) return;
const jobs = [...buckets].map(({subject, facts}) => {
let node = mem.find(m => m.name === subject);
if(!node) {
node = {name: subject, description: '', content: '', embedding: [],};
mem.push(node);
}
const week = subject.startsWith('Journal/') ? {monday, sunday} : undefined;
return this.enqueue(node, facts, mem, options, tempName, week);
});
await Promise.all(jobs);
}
private buildHeader(node: Memory, week?: {monday: string, sunday: string}, links: string[] = [], backlinks: string[] = []): string {
const tags = node.name.split('/')[0]?.toLowerCase();
const lines = [
'---',
`name: ${node.name}`,
`description: ${node.description || ''}`,
tags ? `tags: [${tags}]` : '',
links.length ? `links: [${links.map(l => `"${l}"`).join(', ')}]` : 'links: []',
backlinks.length ? `backlinks: [${backlinks.map(l => `"${l}"`).join(', ')}]` : 'backlinks: []',
week ? `week: ${week.monday} ${week.sunday}` : '',
`modified: ${new Date().toISOString()}`,
'---',
].filter(Boolean);
return lines.join('\n');
}
private cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] {
const scored = memories
.filter(m => m.embedding?.length)
.map(m => ({
ref: {name: m.name, description: m.description},
distance: cosineDistance(query, m.embedding),
}))
.sort((a, b) => a.distance - b.distance)
.slice(0, limit);
return scored.map(s => s.ref);
}
/**
* Coalescing queue: if a doc is already compiling, abort the in-flight run, merge its
* facts with the new ones and restart. Never blocks a pending update, never drops facts.
*/
private enqueue(node: Memory, facts: string[], memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string, week?: {monday: string, sunday: string}): Promise<void> {
const key = node.name;
const existing = this.queues.get(key);
if (existing) {
existing.pending.push(...facts);
existing.request?.abort?.();
return existing.task;
}
const entry: {pending: string[], request: {abort?: () => void} | null, task: Promise<void>} = {pending: [...facts], request: null, task: Promise.resolve()};
this.queues.set(key, entry);
const m = memories instanceof MemoryCache ? memories.memories : memories;
entry.task = (async () => {
while (entry.pending.length) {
const batch = dedupeFacts(entry.pending.splice(0, entry.pending.length));
const written = await this.docAgent(node, batch, m, options, tempName, week, entry);
if (!written) entry.pending.unshift(...batch);
}
})().finally(() => {
this.queues.delete(key);
if(!this.queues.size && memories instanceof MemoryCache) memories.rebuild();
});
return entry.task;
}
private listNodes(memories: Memory[]): MemoryRef[] {
return memories.map(m => ({name: m.name, description: m.description}));
}
decay() {
for(const [name, ttl] of this.recentlyTouched) {
if(ttl <= 1) this.recentlyTouched.delete(name);
else this.recentlyTouched.set(name, ttl - 1);
}
}
forget(name: string, memories: Memory[] | MemoryCache): boolean { forget(name: string, memories: Memory[] | MemoryCache): boolean {
const mem = memories instanceof MemoryCache ? memories.memories : memories; const mem = memories instanceof MemoryCache ? memories.memories : memories;
const idx = mem.findIndex(m => m.name === name); const idx = mem.findIndex(m => m.name === name);
@@ -316,20 +409,43 @@ ${conversation}`;
return true; return true;
} }
private cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] { getTouched(): string[] {
const scored = memories return [...this.recentlyTouched.keys()];
.filter(m => m.embedding?.length)
.map(m => ({
ref: {name: m.name, description: m.description},
distance: cosineDistance(query, m.embedding),
}))
.sort((a, b) => a.distance - b.distance)
.slice(0, limit);
return scored.map(s => s.ref);
} }
private listNodes(memories: Memory[]): MemoryRef[] { async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise<Memory[]> {
return memories.map(m => ({name: m.name, description: m.description})); const conversation = history
.filter(h => h.role === 'user' || h.role === 'assistant')
.map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim();
if(!conversation) return [];
const trackingId = `${Date.now()}_${Math.random()}`;
const tempMemory = await this.createTempMemory(conversation);
const mem = memories instanceof MemoryCache ? memories.memories : memories;
mem.push(tempMemory);
if (memories instanceof MemoryCache) memories.rebuild();
this.pendingMemorizations.set(trackingId, {
memories,
tempMemoryName: tempMemory.name,
timestamp: Date.now(),
});
try {
await this.backgroundMemorization(conversation, memories, options, tempMemory.name);
const finalMem = memories instanceof MemoryCache ? memories.memories : memories;
return finalMem.filter(m => !m.name.startsWith('_temp_'));
} finally {
const pending = this.pendingMemorizations.get(trackingId);
if (pending) {
const cleanMem = pending.memories instanceof MemoryCache
? pending.memories.memories
: pending.memories;
const idx = cleanMem.findIndex(m => m.name === pending.tempMemoryName);
if (idx !== -1) cleanMem.splice(idx, 1);
if (pending.memories instanceof MemoryCache) pending.memories.rebuild();
}
this.pendingMemorizations.delete(trackingId);
}
} }
async recollect(query: string, memories: Memory[] | MemoryCache, limit = 5, graphDepth = 1): Promise<Memory[]> { async recollect(query: string, memories: Memory[] | MemoryCache, limit = 5, graphDepth = 1): Promise<Memory[]> {
@@ -370,106 +486,8 @@ ${conversation}`;
return ordered.map(n => mem.find(m => m.name === n)!).filter(Boolean); return ordered.map(n => mem.find(m => m.name === n)!).filter(Boolean);
} }
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise<Memory[]> { touch(name: string, ttl = 2) {
const conversation = history this.recentlyTouched.set(name, ttl);
.filter(h => h.role === 'user' || h.role === 'assistant')
.map(h => `[${h.role}]: ${h.content}`).join('\n\n').trim();
if(!conversation) return [];
const trackingId = `${Date.now()}_${Math.random()}`;
const tempMemory = await this.createTempMemory(conversation);
const mem = memories instanceof MemoryCache ? memories.memories : memories;
mem.push(tempMemory);
if (memories instanceof MemoryCache) memories.rebuild();
this.pendingMemorizations.set(trackingId, {
memories,
tempMemoryName: tempMemory.name,
timestamp: Date.now(),
});
try {
await this._memorizeBackground(conversation, memories, options, tempMemory.name);
const finalMem = memories instanceof MemoryCache ? memories.memories : memories;
return finalMem.filter(m => !m.name.startsWith('_temp_'));
} finally {
const pending = this.pendingMemorizations.get(trackingId);
if (pending) {
const cleanMem = pending.memories instanceof MemoryCache
? pending.memories.memories
: pending.memories;
const idx = cleanMem.findIndex(m => m.name === pending.tempMemoryName);
if (idx !== -1) cleanMem.splice(idx, 1);
if (pending.memories instanceof MemoryCache) pending.memories.rebuild();
}
this.pendingMemorizations.delete(trackingId);
}
}
private async _memorizeBackground(conversation: string, memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string): Promise<void> {
const mem = memories instanceof MemoryCache ? memories.memories : memories;
const monday = getWeekMonday();
const sunday = getWeekSunday(monday);
const buckets = await this.factAgent(conversation, mem, options, monday);
if(!buckets.length) return;
const jobs = [...buckets].map(({subject, facts}) => {
let node = mem.find(m => m.name === subject);
if(!node) {
node = {name: subject, description: '', content: '', embedding: [],};
mem.push(node);
}
const week = subject.startsWith('Journal/') ? {monday, sunday} : undefined;
return this.enqueue(node, facts, mem, options, tempName, week);
});
await Promise.all(jobs);
}
/**
* Coalescing queue: if a doc is already compiling, abort the in-flight run, merge its
* facts with the new ones and restart. Never blocks a pending update, never drops facts.
*/
private enqueue(node: Memory, facts: string[], memories: Memory[] | MemoryCache, options: LLMRequest, tempName: string, week?: {monday: string, sunday: string}): Promise<void> {
const key = node.name;
const existing = this.queues.get(key);
if (existing) {
existing.pending.push(...facts);
existing.request?.abort?.();
return existing.task;
}
const entry: {pending: string[], request: {abort?: () => void} | null, task: Promise<void>} = {pending: [...facts], request: null, task: Promise.resolve()};
this.queues.set(key, entry);
const m = memories instanceof MemoryCache ? memories.memories : memories;
entry.task = (async () => {
while (entry.pending.length) {
const batch = dedupeFacts(entry.pending.splice(0, entry.pending.length));
const written = await this.docAgent(node, batch, m, options, tempName, week, entry);
if (!written) entry.pending.unshift(...batch);
}
})().finally(() => {
this.queues.delete(key);
if(!this.queues.size && memories instanceof MemoryCache) memories.rebuild();
});
return entry.task;
}
private buildHeader(node: Memory, week?: {monday: string, sunday: string}, links: string[] = [], backlinks: string[] = []): string {
const tags = node.name.split('/')[0]?.toLowerCase();
const lines = [
'---',
`name: ${node.name}`,
`description: ${node.description || ''}`,
tags ? `tags: [${tags}]` : '',
links.length ? `links: [${links.map(l => `"${l}"`).join(', ')}]` : 'links: []',
backlinks.length ? `backlinks: [${backlinks.map(l => `"${l}"`).join(', ')}]` : 'backlinks: []',
week ? `week: ${week.monday} ${week.sunday}` : '',
`modified: ${new Date().toISOString()}`,
'---',
].filter(Boolean);
return lines.join('\n');
}
private applyHeader(content: string, header: string): string {
return `${header}\n\n${this.stripHeader(content)}`;
} }
private updateFrontmatter(content: string, updates: {links?: string[], backlinks?: string[]}): string { private updateFrontmatter(content: string, updates: {links?: string[], backlinks?: string[]}): string {

View File

@@ -20,15 +20,17 @@ 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});
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
})); })));
history.splice(i, 1, ...tools); history.splice(i, 1, ...items);
i += tools.length - 1; 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 +71,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,7 +110,7 @@ export class OpenAi extends LLMProvider {
}; };
} }
let resp: any, isFirstMessage = true, terminal = false; let resp: any, terminal = false;
do { do {
requestParams.messages = history.map(({timestamp, ...m}) => m); requestParams.messages = history.map(({timestamp, ...m}) => m);
resp = await this.client.chat.completions.create(requestParams).catch(err => { resp = await this.client.chat.completions.create(requestParams).catch(err => {
@@ -116,8 +119,6 @@ export class OpenAi extends LLMProvider {
}); });
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;
@@ -125,7 +126,6 @@ export class OpenAi extends LLMProvider {
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);
@@ -168,7 +168,7 @@ export class OpenAi extends LLMProvider {
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 +180,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()});
} }
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 || '');
if(h.role === 'tool') return str + `<tool>${h.name}</tool>\n\n`;
return str;
}, '').trim();
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent); res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
}), {abort: () => controller.abort()}); }), {abort: () => controller.abort()});
} }

View File

@@ -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 {