Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 89619e211e | |||
| afc6653364 | |||
| 68e72445a2 | |||
| 1aa6cdf329 | |||
| d022a5ef4d |
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@ztimson/ai-utils",
|
"name": "@ztimson/ai-utils",
|
||||||
"version": "1.2.12",
|
"version": "1.3.3",
|
||||||
"description": "AI Utility library",
|
"description": "AI Utility library",
|
||||||
"author": "Zak Timson",
|
"author": "Zak Timson",
|
||||||
"license": "MIT",
|
"license": "MIT",
|
||||||
|
|||||||
@@ -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()});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,8 +3,6 @@ export * from './antrhopic';
|
|||||||
export * from './audio';
|
export * from './audio';
|
||||||
export * from './llm';
|
export * from './llm';
|
||||||
export * from './memory';
|
export * from './memory';
|
||||||
export * from './memory-cache';
|
|
||||||
export * from './memory-graph';
|
|
||||||
export * from './open-ai';
|
export * from './open-ai';
|
||||||
export * from './provider';
|
export * from './provider';
|
||||||
export * from './tools';
|
export * from './tools';
|
||||||
|
|||||||
193
src/llm.ts
193
src/llm.ts
@@ -1,17 +1,30 @@
|
|||||||
|
import {snakeCase} from '@ztimson/utils';
|
||||||
import {AbortablePromise, Ai} from './ai.ts';
|
import {AbortablePromise, Ai} from './ai.ts';
|
||||||
import {Anthropic} from './antrhopic.ts';
|
import {Anthropic} from './antrhopic.ts';
|
||||||
import {MemoryCache} from './memory-cache.ts';
|
|
||||||
import {OpenAi} from './open-ai.ts';
|
import {OpenAi} from './open-ai.ts';
|
||||||
import {LLMProvider} from './provider.ts';
|
import {LLMProvider} from './provider.ts';
|
||||||
import {AiTool, AiToolArg} from './tools.ts';
|
import {AiTool, AiToolArg} from './tools.ts';
|
||||||
import {fileURLToPath} from 'url';
|
import {fileURLToPath} from 'url';
|
||||||
import {dirname, join} from 'path';
|
import {dirname, join} from 'path';
|
||||||
import {spawn} from 'node:child_process';
|
import {spawn} from 'node:child_process';
|
||||||
import {Memory, MemoryManager} from './memory.ts';
|
import {Memory, MemoryCache, MemoryManager, MemoryOptions} from './memory.ts';
|
||||||
|
|
||||||
export type AnthropicConfig = {proto: 'anthropic', token: string};
|
export type AnthropicConfig = {proto: 'anthropic', token: string};
|
||||||
export type OpenAiConfig = {proto: 'openai', host?: string, token: string};
|
export type OpenAiConfig = {proto: 'openai', host?: string, token: string};
|
||||||
|
|
||||||
|
export type Agent = {
|
||||||
|
name: string;
|
||||||
|
description?: string;
|
||||||
|
model?: string | null;
|
||||||
|
temperature?: number;
|
||||||
|
system: string;
|
||||||
|
delegate?: boolean;
|
||||||
|
skills?: Skill[] | null;
|
||||||
|
tools?: AiTool[] | null;
|
||||||
|
mcp?: McpServer[] | null;
|
||||||
|
agents?: string[] | null;
|
||||||
|
}
|
||||||
|
|
||||||
export type LLMMessage = {
|
export type LLMMessage = {
|
||||||
/** Message originator */
|
/** Message originator */
|
||||||
role: 'assistant' | 'system' | 'user';
|
role: 'assistant' | 'system' | 'user';
|
||||||
@@ -56,13 +69,17 @@ export type LLMRequest = {
|
|||||||
/** Compress old messages in the chat to free up context */
|
/** Compress old messages in the chat to free up context */
|
||||||
compress?: {max: number; min: number};
|
compress?: {max: number; min: number};
|
||||||
/** User's memory documents - RAG injected automatically each turn */
|
/** User's memory documents - RAG injected automatically each turn */
|
||||||
memory?: Memory[] | MemoryCache;
|
memory?: Memory[] | MemoryCache | MemoryOptions;
|
||||||
/** Model to use for memory operations */
|
/** Model to use for memory operations */
|
||||||
memoryModel?: string;
|
memoryModel?: string;
|
||||||
/** Skill documents the AI can browse and read on demand */
|
/** Skill documents the AI can browse and read on demand */
|
||||||
skills?: Skill[];
|
skills?: Skill[];
|
||||||
/** MCP servers to connect and expose as tools */
|
/** MCP servers to connect and expose as tools */
|
||||||
mcp?: McpServer[];
|
mcp?: McpServer[];
|
||||||
|
/** Subagents exposed as delegatable/wrapped tools */
|
||||||
|
agents?: Agent[];
|
||||||
|
/** @internal recursion guard for nested agent delegation */
|
||||||
|
_agentDepth?: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
export type McpServer = {
|
export type McpServer = {
|
||||||
@@ -83,6 +100,7 @@ export type Skill = {
|
|||||||
content: string;
|
content: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const MAX_AGENT_DEPTH = 5;
|
||||||
|
|
||||||
class LLM {
|
class LLM {
|
||||||
private memoryManager!: MemoryManager;
|
private memoryManager!: MemoryManager;
|
||||||
@@ -100,6 +118,59 @@ class LLM {
|
|||||||
this.memoryManager = new MemoryManager(this);
|
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, any>, aborts: (() => void)[], depth = 0): 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},
|
||||||
|
instructions: {type: 'string', description: 'Detailed instructions for subagent to complete', required: true},
|
||||||
|
},
|
||||||
|
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.
|
||||||
|
|
||||||
|
${a.system}`,
|
||||||
|
model: a.model || undefined,
|
||||||
|
temperature: a.temperature,
|
||||||
|
stream: a.delegate ? stream : undefined,
|
||||||
|
history: subHistory,
|
||||||
|
mcp: a.mcp || undefined,
|
||||||
|
skills: a.skills || undefined,
|
||||||
|
tools: a.tools || undefined,
|
||||||
|
agents: nested,
|
||||||
|
_agentDepth: depth + 1,
|
||||||
|
} as any);
|
||||||
|
aborts.push(request.abort);
|
||||||
|
const resp = await request;
|
||||||
|
|
||||||
|
if(a.delegate) {
|
||||||
|
pending.set(<string>id, {resp, subHistory});
|
||||||
|
return '';
|
||||||
|
}
|
||||||
|
return resp;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
private async setupMcp(servers: McpServer[] = []): Promise<{prompt: string, tools: AiTool[]}> {
|
private async setupMcp(servers: McpServer[] = []): Promise<{prompt: string, tools: AiTool[]}> {
|
||||||
if(!servers?.length) return {prompt: '', tools: []};
|
if(!servers?.length) return {prompt: '', tools: []};
|
||||||
const allTools: AiTool[] = [];
|
const allTools: AiTool[] = [];
|
||||||
@@ -170,9 +241,11 @@ class LLM {
|
|||||||
if(!this.models[m]) throw new Error(`Model does not exist: ${m}`);
|
if(!this.models[m]) throw new Error(`Model does not exist: ${m}`);
|
||||||
let request: AbortablePromise<string> | null = null;
|
let request: AbortablePromise<string> | null = null;
|
||||||
let aborted = false;
|
let aborted = false;
|
||||||
|
const nestedAborts: (() => void)[] = [];
|
||||||
const abort = () => {
|
const abort = () => {
|
||||||
aborted = true;
|
aborted = true;
|
||||||
request?.abort?.();
|
request?.abort?.();
|
||||||
|
nestedAborts.forEach(a => a());
|
||||||
};
|
};
|
||||||
|
|
||||||
const promise = (async () => {
|
const promise = (async () => {
|
||||||
@@ -196,22 +269,46 @@ class LLM {
|
|||||||
tools.push(...s.tools);
|
tools.push(...s.tools);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Agents
|
||||||
|
const agents = options.agents || this.ai.options?.llm?.agents;
|
||||||
|
const pendingDelegates = new Map<string, any>();
|
||||||
|
if(agents?.length) tools.push(...this.setupAgent(agents, agents, pendingDelegates, nestedAborts, options._agentDepth || 0));
|
||||||
|
|
||||||
// Memory
|
// Memory
|
||||||
if (options.memory) {
|
const mem = MemoryManager.normalize(options.memory);
|
||||||
const mems = options.memory instanceof MemoryCache ? options.memory.memories : options.memory;
|
if(mem) {
|
||||||
|
const mems = mem.memory instanceof MemoryCache ? mem.memory.memories : mem.memory;
|
||||||
if(mems.length) {
|
if(mems.length) {
|
||||||
const relevant = await this.memoryManager.recollect(message, options.memory, 5);
|
if(mem.inject) {
|
||||||
prompts.unshift(`You have access to the following memory files:
|
const pool = 15; // candidates considered, cheap since only refs are listed
|
||||||
${mems.map(m => `- ${m.name}: ${m.description}`).join('\n')}
|
const budget = mem.maxTokens ?? 2000; // actual content injected
|
||||||
${relevant.length ? `
|
const relevant = await this.memoryManager.recollect(message, mem.memory, pool);
|
||||||
Relevant memories have been preloaded:
|
|
||||||
${relevant.map(r => `
|
let used = 0;
|
||||||
**${r.name}**
|
const preloaded: typeof relevant = [];
|
||||||
${r.description}
|
const listed: typeof relevant = [];
|
||||||
${r.content}
|
for(const r of relevant) {
|
||||||
`).join('\n---\n')}
|
const t = this.estimateTokens(r.content);
|
||||||
` : ''}`.trim());
|
if(used + t <= budget || preloaded.length === 0) {
|
||||||
tools.push(this.memoryManager.tools.read(options.memory));
|
preloaded.push(r);
|
||||||
|
used += t;
|
||||||
|
} 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());
|
||||||
|
}
|
||||||
|
if(mem.tool) tools.push(this.memoryManager.tools.read(mem.memory));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -219,16 +316,32 @@ class LLM {
|
|||||||
|
|
||||||
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')});
|
||||||
const resp = await request;
|
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: 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;
|
||||||
|
i += insert.length;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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;
|
||||||
|
|
||||||
// Trim memory injections from history
|
// Trim memory injections from history
|
||||||
if(options.memory) {
|
if(mem?.tool) history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== 'memory_recall'));
|
||||||
history.splice(0, history.length, ...history.filter(h => h.role !== 'tool' || h.name !== '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) {
|
||||||
if(options.memory) await this.memoryManager.memorize(history, options.memory, {model: options.memoryModel || this.defaultModel, ...options});
|
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);
|
const compressed = await this.compressHistory(history, options.compress.max, options.compress.min, options);
|
||||||
if(options.history) options.history.splice(0, options.history.length, ...compressed);
|
if(options.history) options.history.splice(0, options.history.length, ...compressed);
|
||||||
}
|
}
|
||||||
@@ -397,15 +510,33 @@ class LLM {
|
|||||||
* @param {string} searchTerms Multiple search terms to check against target
|
* @param {string} searchTerms Multiple search terms to check against target
|
||||||
* @returns {{avg: number, max: number, similarities: number[]}} Similarity values 0-1: 0 = unique, 1 = identical
|
* @returns {{avg: number, max: number, similarities: number[]}} Similarity values 0-1: 0 = unique, 1 = identical
|
||||||
*/
|
*/
|
||||||
fuzzyMatch(target: string, ...searchTerms: string[]) {
|
fuzzyMatch(target, ...searchTerms) {
|
||||||
if(searchTerms.length < 2) throw new Error('Requires at least 2 strings to compare');
|
if (searchTerms.length < 2) throw new Error('Requires at least 2 strings to compare');
|
||||||
const vector = (text: string, dimensions: number = 10): number[] => {
|
const levenshtein = (a, b) => {
|
||||||
return text.toLowerCase().split('').map((char, index) =>
|
const m = a.length, n = b.length;
|
||||||
(char.charCodeAt(0) * (index + 1)) % dimensions / dimensions).slice(0, dimensions);
|
if (!m) return n;
|
||||||
}
|
if (!n) return m;
|
||||||
const v = vector(target);
|
const dp = Array.from({length: m + 1}, (_, i) => [i, ...Array(n).fill(0)]);
|
||||||
const similarities = searchTerms.map(t => vector(t)).map(refVector => this.cosineSimilarity(v, refVector));
|
for (let j = 0; j <= n; j++) dp[0][j] = j;
|
||||||
return {avg: similarities.reduce((acc, s) => acc + s, 0) / similarities.length, max: Math.max(...similarities), similarities};
|
for (let i = 1; i <= m; i++) {
|
||||||
|
for (let j = 1; j <= n; j++) {
|
||||||
|
dp[i][j] = a[i - 1] === b[j - 1]
|
||||||
|
? dp[i - 1][j - 1]
|
||||||
|
: 1 + Math.min(dp[i - 1][j - 1], dp[i - 1][j], dp[i][j - 1]);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return dp[m][n];
|
||||||
|
};
|
||||||
|
const similarity = (a, b) => {
|
||||||
|
a = a.toLowerCase(); b = b.toLowerCase();
|
||||||
|
return 1 - levenshtein(a, b) / Math.max(a.length, b.length, 1);
|
||||||
|
};
|
||||||
|
const similarities = searchTerms.map(t => similarity(target, t));
|
||||||
|
return {
|
||||||
|
avg: similarities.reduce((acc, s) => acc + s, 0) / similarities.length,
|
||||||
|
max: Math.max(...similarities),
|
||||||
|
similarities
|
||||||
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -1,59 +0,0 @@
|
|||||||
import {KDPoint, KDTree} from './kd-tree.ts';
|
|
||||||
import {Memory, MemoryRef} from './memory.ts';
|
|
||||||
|
|
||||||
export class MemoryCache {
|
|
||||||
private tree: KDTree<MemoryRef>;
|
|
||||||
public memories: Memory[];
|
|
||||||
|
|
||||||
get length() { return this.memories.length; }
|
|
||||||
|
|
||||||
constructor(memories: Memory[]) {
|
|
||||||
this.memories = memories;
|
|
||||||
this.tree = this.buildTree();
|
|
||||||
}
|
|
||||||
|
|
||||||
private buildTree(): KDTree<MemoryRef> {
|
|
||||||
const embedded = this.memories.filter(m => m.embedding?.length);
|
|
||||||
if (!embedded.length) return new KDTree<MemoryRef>(0);
|
|
||||||
|
|
||||||
const dims = embedded[0].embedding.length;
|
|
||||||
const points: KDPoint<MemoryRef>[] = embedded.map(m => ({
|
|
||||||
vector: m.embedding,
|
|
||||||
payload: {name: m.name, description: m.description},
|
|
||||||
}));
|
|
||||||
|
|
||||||
return new KDTree<MemoryRef>(dims, 'cosine', points);
|
|
||||||
}
|
|
||||||
|
|
||||||
search(query: number[], limit: number): MemoryRef[] {
|
|
||||||
const results = this.tree.knn(query, limit);
|
|
||||||
return results.map(r => r.point.payload);
|
|
||||||
}
|
|
||||||
|
|
||||||
add(memory: Memory): void {
|
|
||||||
this.memories.push(memory);
|
|
||||||
this.rebuild();
|
|
||||||
}
|
|
||||||
|
|
||||||
update(memory: Memory): void {
|
|
||||||
const idx = this.memories.findIndex(m => m.name === memory.name);
|
|
||||||
if (idx !== -1) {
|
|
||||||
this.memories[idx] = memory;
|
|
||||||
} else {
|
|
||||||
this.memories.push(memory);
|
|
||||||
}
|
|
||||||
this.rebuild();
|
|
||||||
}
|
|
||||||
|
|
||||||
remove(name: string): void {
|
|
||||||
const idx = this.memories.findIndex(m => m.name === name);
|
|
||||||
if (idx !== -1) {
|
|
||||||
this.memories.splice(idx, 1);
|
|
||||||
this.rebuild();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
rebuild(): void {
|
|
||||||
this.tree = this.buildTree();
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,69 +0,0 @@
|
|||||||
import {MemoryCache} from './memory-cache.ts';
|
|
||||||
import {extractMetadata, Memory, MemoryNode} from './memory.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();
|
|
||||||
}
|
|
||||||
389
src/memory.ts
389
src/memory.ts
@@ -1,6 +1,143 @@
|
|||||||
import {LLMRequest, LLMMessage} from './llm.ts';
|
import {LLMRequest, LLMMessage} from './llm.ts';
|
||||||
import {MemoryCache} from './memory-cache.ts';
|
|
||||||
import {AiTool} from './tools.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<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 {
|
||||||
|
private tree: KDTree<MemoryRef>;
|
||||||
|
public memories: Memory[];
|
||||||
|
|
||||||
|
get length() { return this.memories.length; }
|
||||||
|
|
||||||
|
constructor(memories: Memory[]) {
|
||||||
|
this.memories = memories;
|
||||||
|
this.tree = this.buildTree();
|
||||||
|
}
|
||||||
|
|
||||||
|
private buildTree(): KDTree<MemoryRef> {
|
||||||
|
const embedded = this.memories.filter(m => m.embedding?.length);
|
||||||
|
if (!embedded.length) return new KDTree<MemoryRef>(0);
|
||||||
|
|
||||||
|
const dims = embedded[0].embedding.length;
|
||||||
|
const points: KDPoint<MemoryRef>[] = embedded.map(m => ({
|
||||||
|
vector: m.embedding,
|
||||||
|
payload: {name: m.name, description: m.description},
|
||||||
|
}));
|
||||||
|
|
||||||
|
return new KDTree<MemoryRef>(dims, 'cosine', points);
|
||||||
|
}
|
||||||
|
|
||||||
|
search(query: number[], limit: number): MemoryRef[] {
|
||||||
|
const results = this.tree.knn(query, limit);
|
||||||
|
return results.map(r => r.point.payload);
|
||||||
|
}
|
||||||
|
|
||||||
|
add(memory: Memory): void {
|
||||||
|
this.memories.push(memory);
|
||||||
|
this.rebuild();
|
||||||
|
}
|
||||||
|
|
||||||
|
update(memory: Memory): void {
|
||||||
|
const idx = this.memories.findIndex(m => m.name === memory.name);
|
||||||
|
if (idx !== -1) {
|
||||||
|
this.memories[idx] = memory;
|
||||||
|
} else {
|
||||||
|
this.memories.push(memory);
|
||||||
|
}
|
||||||
|
this.rebuild();
|
||||||
|
}
|
||||||
|
|
||||||
|
remove(name: string): void {
|
||||||
|
const idx = this.memories.findIndex(m => m.name === name);
|
||||||
|
if (idx !== -1) {
|
||||||
|
this.memories.splice(idx, 1);
|
||||||
|
this.rebuild();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
rebuild(): void {
|
||||||
|
this.tree = this.buildTree();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export type MemoryOptions = {
|
||||||
|
/** Memory object */
|
||||||
|
memory: Memory[] | MemoryCache;
|
||||||
|
/** Inject N memories into the system prompt */
|
||||||
|
inject?: boolean;
|
||||||
|
/** expose recall tool to LLM */
|
||||||
|
tool?: boolean;
|
||||||
|
/** Update memory on compression */
|
||||||
|
update?: boolean;
|
||||||
|
/** Max context size of memories to inject to each call (removed immediately after use) */
|
||||||
|
maxTokens?: number;
|
||||||
|
}
|
||||||
|
|
||||||
export type Memory = {
|
export type Memory = {
|
||||||
name: string;
|
name: string;
|
||||||
@@ -83,9 +220,9 @@ function getWeekSunday(monday: string): string {
|
|||||||
return d.toISOString().slice(0, 10);
|
return d.toISOString().slice(0, 10);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
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,
|
||||||
@@ -109,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;
|
||||||
},
|
},
|
||||||
}),
|
}),
|
||||||
@@ -128,6 +266,12 @@ export class MemoryManager {
|
|||||||
|
|
||||||
constructor(private llm: any) {}
|
constructor(private llm: any) {}
|
||||||
|
|
||||||
|
static normalize(m?: Memory[] | MemoryCache | MemoryOptions) {
|
||||||
|
if(!m) return null;
|
||||||
|
const raw = m instanceof MemoryCache || Array.isArray(m);
|
||||||
|
return raw ? {memory: <Memory[] | MemoryCache>m, inject: true, tool: true, update: true} : {inject: true, tool: true, update: true, ...m};
|
||||||
|
}
|
||||||
|
|
||||||
private async createTempMemory(conversation: string): Promise<Memory> {
|
private async createTempMemory(conversation: string): Promise<Memory> {
|
||||||
const timestamp = Date.now();
|
const timestamp = Date.now();
|
||||||
const content = `---
|
const content = `---
|
||||||
@@ -151,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);
|
||||||
@@ -175,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[]> {
|
||||||
@@ -229,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 {
|
||||||
|
|||||||
@@ -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()});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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