Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2921b208da | ||
|
|
b5aec246ac | ||
|
|
2d6debad86 | ||
|
|
6bed8f20b5 | ||
|
|
dc45a99b04 | ||
|
|
263a65c192 | ||
|
|
1e8c7c6662 | ||
|
|
1f1a4662d4 | ||
|
|
ee4147e24e | ||
|
|
d29c0ca389 |
+1
-1
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@ztimson/ai-utils",
|
||||
"version": "1.6.8",
|
||||
"version": "1.7.2",
|
||||
"description": "AI Utility library",
|
||||
"author": "Zak Timson",
|
||||
"license": "MIT",
|
||||
|
||||
+5
-2
@@ -1,11 +1,14 @@
|
||||
export * from './ai';
|
||||
export * from './antrhopic';
|
||||
export * from './audio';
|
||||
export * from './helpers';
|
||||
export * from './llm';
|
||||
export * from './memory';
|
||||
export * from './memory/graph';
|
||||
export * from './memory/kd-tree';
|
||||
export * from './memory/memory';
|
||||
export * from './memory/memory-state';
|
||||
export * from './open-ai';
|
||||
export * from './provider';
|
||||
export * from './token-pool'
|
||||
export * from './tools';
|
||||
export * from './vision';
|
||||
export * from './utils';
|
||||
|
||||
+23
-33
@@ -1,17 +1,19 @@
|
||||
import {clean, makeUnique, snakeCase} from '@ztimson/utils';
|
||||
import {AbortablePromise, Ai} from './ai.ts';
|
||||
import {Anthropic} from './antrhopic.ts';
|
||||
import {MemoryCache} from './memory/memory-state.ts';
|
||||
import {Memory, MemoryManager, MemoryOptions} from './memory/memory.ts';
|
||||
import {OpenAi} from './open-ai.ts';
|
||||
import {LLMProvider} from './provider.ts';
|
||||
import {AiTool, AiToolArg} from './tools.ts';
|
||||
import {fileURLToPath} from 'url';
|
||||
import {spawn} from 'node:child_process';
|
||||
import {Memory, MemoryCache, MemoryManager, MemoryOptions, stripHeader} from './memory.ts';
|
||||
import {mkdtempSync} from 'node:fs';
|
||||
import fs from 'node:fs/promises';
|
||||
import {tmpdir} from 'node:os';
|
||||
import {dirname, join, basename, extname} from 'path';
|
||||
import { PDFParse } from 'pdf-parse';
|
||||
import {stripHeader} from './utils.ts';
|
||||
|
||||
const MAX_AGENT_DEPTH = 5;
|
||||
const PDF_OCR_PAGE_THRESHOLD = 12; // above this many pages, OCR scanned pages instead of feeding images to the model
|
||||
@@ -19,6 +21,13 @@ const PDF_OCR_PAGE_THRESHOLD = 12; // above this many pages, OCR scanned pages i
|
||||
export type AnthropicConfig = {proto: 'anthropic', token: string | string[]};
|
||||
export type OpenAiConfig = {proto: 'openai', host?: string, token: string | string[]};
|
||||
|
||||
export type AgentRef = {
|
||||
name: string;
|
||||
description?: string;
|
||||
delegate?: boolean;
|
||||
fn: () => Agent | null | Promise<Agent | null>;
|
||||
}
|
||||
|
||||
export type Agent = {
|
||||
name: string;
|
||||
description?: string;
|
||||
@@ -29,7 +38,7 @@ export type Agent = {
|
||||
skills?: Skill[] | null;
|
||||
tools?: AiTool[] | null;
|
||||
mcp?: McpServer[] | null;
|
||||
agents?: string[] | null;
|
||||
agents?: AgentRef[] | null;
|
||||
}
|
||||
|
||||
export type LLMFile = {
|
||||
@@ -106,8 +115,8 @@ export type LLMRequest = {
|
||||
skills?: Skill[];
|
||||
/** MCP servers to connect and expose as tools */
|
||||
mcp?: McpServer[];
|
||||
/** Subagents exposed as delegatable/wrapped tools */
|
||||
agents?: Agent[];
|
||||
/** Subagents exposed as delegatable/wrapped tools, resolved lazily via their `fn` */
|
||||
agents?: AgentRef[];
|
||||
/** Attach files to request */
|
||||
files?: LLMFile[];
|
||||
/** @internal recursion guard for nested agent delegation */
|
||||
@@ -265,22 +274,21 @@ class LLM {
|
||||
};
|
||||
}
|
||||
|
||||
private setupAgent(agents: Agent[] = [], allAgents: Agent[], history: LLMMessage[], aborts: ((keep?: boolean) => void)[], depth = 0, delegateState: {resp: string | null}): AiTool[] {
|
||||
return agents.map(a => {
|
||||
const toolName = `${a.delegate ? '' : 'sub'}agent_${snakeCase(a.name)}`;
|
||||
private setupAgent(stubs: AgentRef[] = [], history: LLMMessage[], aborts: ((keep?: boolean) => void)[], depth = 0, delegateState: {resp: string | null}): AiTool[] {
|
||||
return stubs.map(stub => {
|
||||
const toolName = `${stub.delegate ? '' : 'sub'}agent_${snakeCase(stub.name)}`;
|
||||
return {
|
||||
name: toolName,
|
||||
description: `${a.delegate ? 'Delegate to ' : ''}Subagent: ${a.description || a.name}`,
|
||||
description: `${stub.delegate ? 'Delegate to ' : ''}Subagent: ${stub.description || stub.name}`,
|
||||
args: clean<any>({
|
||||
context: !a.delegate ? {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true} : undefined,
|
||||
context: !stub.delegate ? {type: 'string', description: 'Summary of related messages, samples, files, etc...', required: true} : undefined,
|
||||
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 nested = (a.agents || [])
|
||||
.map(name => allAgents.find(x => x.name === name))
|
||||
.filter((x): x is Agent => !!x && x.name !== a.name);
|
||||
const a = await stub.fn();
|
||||
if(!a) return `Agent "${stub.name}" could not be resolved`;
|
||||
|
||||
const q = a.delegate ? '' : `${args.instructions}${args.context ? `\n\n<context>${args.context}</context>` : ''}`;
|
||||
|
||||
@@ -297,7 +305,7 @@ ${a.system}`,
|
||||
mcp: a.mcp || undefined,
|
||||
skills: a.skills || undefined,
|
||||
tools: a.tools || undefined,
|
||||
agents: nested,
|
||||
agents: a.agents || [],
|
||||
_agentDepth: depth + 1,
|
||||
} as any);
|
||||
aborts.push(request.abort);
|
||||
@@ -451,7 +459,7 @@ ${a.system}`,
|
||||
// Agents
|
||||
const agents = options.agents || this.ai.options?.llm?.agents;
|
||||
const delegateState: {resp: string | null} = {resp: null};
|
||||
if(agents?.length) tools.push(...this.setupAgent(agents, agents, history, nestedAborts, options._agentDepth || 0, delegateState));
|
||||
if(agents?.length) tools.push(...this.setupAgent(agents, history, nestedAborts, options._agentDepth || 0, delegateState));
|
||||
|
||||
// Memory
|
||||
const mem = MemoryManager.normalize(options.memory);
|
||||
@@ -495,7 +503,7 @@ Description: ${r.description}
|
||||
Linked: ${makeUnique([...r.links, ...r.backlinks]).join(', ')}
|
||||
<!-- Truncated -->`).join('\n\n') : ''}`.trim())
|
||||
}
|
||||
if(mem.tool) tools.push(this.memoryManager.tools.read(mem.memory));
|
||||
if(mem.tool) tools.push(...this.memoryManager.tools.read(mem.memory));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -588,24 +596,6 @@ Linked: ${makeUnique([...r.links, ...r.backlinks]).join(', ')}
|
||||
return h;
|
||||
}
|
||||
|
||||
/**
|
||||
* Compare the difference between embeddings (calculates the angle between two vectors)
|
||||
* @param {number[]} v1 First embedding / vector comparison
|
||||
* @param {number[]} v2 Second embedding / vector for comparison
|
||||
* @returns {number} Similarity values 0-1: 0 = unique, 1 = identical
|
||||
*/
|
||||
cosineSimilarity(v1: number[], v2: number[]): number {
|
||||
if (v1.length !== v2.length) throw new Error('Vectors must be same length');
|
||||
let dotProduct = 0, normA = 0, normB = 0;
|
||||
for (let i = 0; i < v1.length; i++) {
|
||||
dotProduct += v1[i] * v2[i];
|
||||
normA += v1[i] * v1[i];
|
||||
normB += v2[i] * v2[i];
|
||||
}
|
||||
const denominator = Math.sqrt(normA) * Math.sqrt(normB);
|
||||
return denominator === 0 ? 0 : dotProduct / denominator;
|
||||
}
|
||||
|
||||
/**
|
||||
* Chunk text into parts for AI digestion
|
||||
* @param {object | string} target Item that will be chunked (objects get converted)
|
||||
|
||||
-732
@@ -1,732 +0,0 @@
|
||||
import {MemoryNode, patchGraph, rebuildGraph} from './helpers.ts';
|
||||
import {LLMRequest, LLMMessage} from './llm.ts';
|
||||
import {AiTool} from './tools.ts';
|
||||
import {KDTree} from './kd-tree.ts';
|
||||
import {escapeRegex} from '@ztimson/utils';
|
||||
|
||||
const MERGE_THRESHOLD = 0.12;
|
||||
const PENDING_HEADING = '## Pending';
|
||||
const TREE_TOMBSTONE_LIMIT = 0.25;
|
||||
const ALIAS_MATCH_THRESHOLD = 0.55;
|
||||
const GENERIC_TEMPLATE = `# {{Title}}
|
||||
|
||||
## Summary
|
||||
|
||||
## Details
|
||||
|
||||
## Related`;
|
||||
|
||||
export type Memory = {
|
||||
name: string;
|
||||
description: string;
|
||||
content: string;
|
||||
/** Description embedding — indexed in the KD tree, used for merge/ANN candidate lookup */
|
||||
embedding: number[];
|
||||
/** Title-only embedding, weighted heaviest during recall ranking */
|
||||
titleEmbedding?: number[];
|
||||
/** Chunked body embeddings, best-chunk match used during recall ranking */
|
||||
bodyEmbeddings?: number[][];
|
||||
links: string[];
|
||||
backlinks: string[];
|
||||
}
|
||||
|
||||
type MemoryRef = {
|
||||
name: string;
|
||||
description: string;
|
||||
/** Cosine distance from the query, present when returned from a search */
|
||||
distance?: number;
|
||||
}
|
||||
|
||||
type FactBucket = {
|
||||
subject: string;
|
||||
facts: string[];
|
||||
}
|
||||
|
||||
type FactAgentResult = {
|
||||
buckets: FactBucket[];
|
||||
journal: string;
|
||||
}
|
||||
|
||||
function dedupeFacts(facts: string[]): string[] {
|
||||
const seen = new Map<string, string>();
|
||||
for (const f of facts) {
|
||||
const clean = f.trim();
|
||||
if (clean) seen.set(clean.toLowerCase(), clean);
|
||||
}
|
||||
return [...seen.values()];
|
||||
}
|
||||
|
||||
function cosineDistance(a: number[], b: number[]): number {
|
||||
let dot = 0, normA = 0, normB = 0;
|
||||
for (let i = 0; i < a.length; i++) {
|
||||
dot += a[i] * b[i];
|
||||
normA += a[i] * a[i];
|
||||
normB += b[i] * b[i];
|
||||
}
|
||||
const denom = Math.sqrt(normA) * Math.sqrt(normB);
|
||||
return denom === 0 ? 1 : 1 - dot / denom;
|
||||
}
|
||||
|
||||
function cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] {
|
||||
return memories
|
||||
.filter(m => m.embedding?.length)
|
||||
.map(m => ({name: m.name, description: m.description, distance: cosineDistance(query, m.embedding)}))
|
||||
.sort((a, b) => a.distance - b.distance)
|
||||
.slice(0, limit);
|
||||
}
|
||||
|
||||
/** Re-embed a node's title / description / body fields. Description embedding stays the KD-tree index key. */
|
||||
async function embedMemoryFields(node: Memory, llm: any): Promise<void> {
|
||||
const body = stripHeader(node.content);
|
||||
const [titleE] = await llm.embedding(node.name.split('/').pop() || node.name);
|
||||
const [descE] = await llm.embedding(node.description || '');
|
||||
const bodyChunks = body ? await llm.embedding(body) : [];
|
||||
if (titleE) node.titleEmbedding = titleE.embedding;
|
||||
if (descE) node.embedding = descE.embedding;
|
||||
node.bodyEmbeddings = bodyChunks.map((c: any) => c.embedding).filter(Boolean);
|
||||
}
|
||||
|
||||
export function stripHeader(content: string): string {
|
||||
return content.replace(/^---[\s\S]*?\n---\n?/, '').trimStart();
|
||||
}
|
||||
|
||||
export class MemoryCache {
|
||||
private tree!: KDTree<MemoryRef>;
|
||||
/** Tracks which memories are currently indexed in the tree, keyed by name -> embedding reference */
|
||||
private indexed = new Map<string, number[]>();
|
||||
public memories: Memory[];
|
||||
public nodes: MemoryNode[] = [];
|
||||
|
||||
get length() { return this.memories.length; }
|
||||
|
||||
constructor(memories: Memory[]) {
|
||||
this.memories = memories;
|
||||
this.tree = new KDTree<MemoryRef>(0);
|
||||
this.rebuild();
|
||||
}
|
||||
|
||||
/** Incrementally sync the KD tree against `this.memories` instead of rebuilding from scratch */
|
||||
private syncTree(): void {
|
||||
const current = new Set(this.memories.map(m => m.name));
|
||||
|
||||
for (const [name, emb] of [...this.indexed]) {
|
||||
const mem = this.memories.find(m => m.name === name);
|
||||
if (!mem || !current.has(name) || mem.embedding !== emb) {
|
||||
this.tree.remove(p => p.name === name);
|
||||
this.indexed.delete(name);
|
||||
}
|
||||
}
|
||||
|
||||
for (const mem of this.memories) {
|
||||
if (!mem.embedding?.length || this.indexed.has(mem.name)) continue;
|
||||
if (this.tree.dims === 0) this.tree = new KDTree<MemoryRef>(mem.embedding.length, 'cosine');
|
||||
if (mem.embedding.length !== this.tree.dims) continue; // guard against embedding model/dim drift
|
||||
this.tree.insert({vector: mem.embedding, payload: {name: mem.name, description: mem.description}});
|
||||
this.indexed.set(mem.name, mem.embedding);
|
||||
}
|
||||
|
||||
if (this.tree.tombstoneRatio > TREE_TOMBSTONE_LIMIT) this.tree.rebalance();
|
||||
}
|
||||
|
||||
search(query: number[], limit: number): MemoryRef[] {
|
||||
if (!this.tree || this.tree.dims === 0) return [];
|
||||
return this.tree.knn(query, limit).map(r => ({...r.point.payload, distance: r.distance}));
|
||||
}
|
||||
|
||||
add(memory: Memory): void {
|
||||
this.memories.push(memory);
|
||||
this.rebuild([memory]);
|
||||
}
|
||||
|
||||
update(memory: Memory): void {
|
||||
const existing = this.memories.find(m => m.name === memory.name);
|
||||
if (existing) Object.assign(existing, memory);
|
||||
else this.memories.push(memory);
|
||||
this.rebuild([existing ?? memory]);
|
||||
}
|
||||
|
||||
remove(name: string): void {
|
||||
const idx = this.memories.findIndex(m => m.name === name);
|
||||
if (idx !== -1) {
|
||||
this.memories.splice(idx, 1);
|
||||
this.rebuild();
|
||||
}
|
||||
}
|
||||
|
||||
rebuild(changed?: Memory[]): void {
|
||||
this.nodes = (changed?.length && this.nodes.length)
|
||||
? patchGraph(this.memories, this.nodes, changed)
|
||||
: rebuildGraph(this.memories);
|
||||
this.syncTree();
|
||||
}
|
||||
}
|
||||
|
||||
class MemoryAccessor {
|
||||
readonly list: Memory[];
|
||||
private readonly cache: MemoryCache | null;
|
||||
|
||||
constructor(memories: Memory[] | MemoryCache) {
|
||||
this.cache = memories instanceof MemoryCache ? memories : null;
|
||||
this.list = this.cache ? this.cache.memories : <Memory[]>memories;
|
||||
}
|
||||
|
||||
find(name: string): Memory | undefined {
|
||||
return this.list.find(m => m.name === name);
|
||||
}
|
||||
|
||||
commit(changed?: Memory[]): MemoryNode[] {
|
||||
if (this.cache) {
|
||||
this.cache.rebuild(changed);
|
||||
return this.cache.nodes;
|
||||
}
|
||||
return rebuildGraph(this.list);
|
||||
}
|
||||
|
||||
ghosts(): string[] {
|
||||
const nodes = this.cache ? this.cache.nodes : rebuildGraph(this.list);
|
||||
return nodes.filter(n => n.missing).map(n => n.name);
|
||||
}
|
||||
|
||||
/** Cache path uses the KD tree's knn(); raw-array path (no cache available) falls back to a linear cosine scan */
|
||||
search(vector: number[], limit: number): MemoryRef[] {
|
||||
return this.cache ? this.cache.search(vector, limit) : cosineSearch(vector, this.list, limit);
|
||||
}
|
||||
|
||||
forget(name: string): boolean {
|
||||
const idx = this.list.findIndex(m => m.name === name);
|
||||
if (idx === -1) return false;
|
||||
this.list.splice(idx, 1);
|
||||
this.commit();
|
||||
return true;
|
||||
}
|
||||
|
||||
async backfillEmbeddings(llm: any): Promise<number> {
|
||||
const missing = this.list.filter(m => !m.embedding?.length);
|
||||
if (!missing.length) return 0;
|
||||
await Promise.all(missing.map(node => embedMemoryFields(node, llm)));
|
||||
this.commit();
|
||||
return missing.length;
|
||||
}
|
||||
}
|
||||
|
||||
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 class MemoryManager {
|
||||
private mergeLock: Promise<any> = Promise.resolve();
|
||||
private queues = new Map<string, {
|
||||
dirty: boolean,
|
||||
request: {abort?: () => void} | null,
|
||||
task: Promise<void>,
|
||||
}>();
|
||||
private recentlyTouched = new Map<string, number>();
|
||||
|
||||
tools = {
|
||||
forget: (memories: Memory[] | MemoryCache): AiTool => ({
|
||||
name: 'memory_forget',
|
||||
description: 'Permanently delete a memory document and clean up all references to it',
|
||||
args: {
|
||||
name: {type: 'string', description: 'Exact memory name to forget', required: true}
|
||||
},
|
||||
fn: (args: any) => {
|
||||
const result = this.forget(args.name, memories);
|
||||
return result ? `Forgotten: ${args.name}` : `Not found: ${args.name}`;
|
||||
},
|
||||
}),
|
||||
|
||||
read: (memories: Memory[] | MemoryCache): AiTool => ({
|
||||
name: 'memory_recall',
|
||||
description: 'Read the full content of a memory document',
|
||||
args: {
|
||||
name: {type: 'string', description: 'Exact memory name', required: true},
|
||||
},
|
||||
fn: (args: any) => {
|
||||
const mem = this.access(memories).find(args.name);
|
||||
if (!mem) return 'Document not found';
|
||||
this.touch(mem.name);
|
||||
return mem.content;
|
||||
},
|
||||
}),
|
||||
|
||||
search: (memories: Memory[] | MemoryCache): AiTool => ({
|
||||
name: 'memory_search',
|
||||
description: 'Use embeddings to find the MOST relevant memories, even if NOT relevant',
|
||||
args: {
|
||||
query: {type: 'string', description: 'What to look for in the memories', required: true},
|
||||
limit: {type: 'number', description: 'Number of memories to return', default: 1},
|
||||
},
|
||||
fn: async ({query, limit}) => {
|
||||
const mem = await this.recollect(query, memories, limit)
|
||||
return mem.map(m => `Memory: ${m.name}
|
||||
Description: ${m.description}
|
||||
Links: ${[...m.links, ...m.backlinks].join(', ')}
|
||||
\`\`\`
|
||||
${m.content}
|
||||
\`\`\``).join('\n\n');
|
||||
},
|
||||
}),
|
||||
};
|
||||
|
||||
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 access(memories: Memory[] | MemoryCache): MemoryAccessor {
|
||||
return new MemoryAccessor(memories);
|
||||
}
|
||||
|
||||
private stage(node: Memory, block: string): void {
|
||||
this.ensureDoc(node);
|
||||
const body = stripHeader(node.content);
|
||||
const idx = body.indexOf(PENDING_HEADING);
|
||||
const newBody = idx === -1
|
||||
? `${body.trimEnd()}\n\n${PENDING_HEADING}\n${block}\n`
|
||||
: `${body.slice(0, idx + PENDING_HEADING.length)}\n${block}${body.slice(idx + PENDING_HEADING.length)}`;
|
||||
node.content = this.touchHeader(node, newBody);
|
||||
}
|
||||
|
||||
private ensureDoc(node: Memory): void {
|
||||
if (node.content) return;
|
||||
const title = node.name.split('/').pop() ?? node.name;
|
||||
node.content = this.touchHeader(node, `# ${title}\n`);
|
||||
}
|
||||
|
||||
private sanitizeDescription(text: string): string {
|
||||
return (text ?? '').replace(/\s+/g, ' ').trim().slice(0, 240);
|
||||
}
|
||||
|
||||
private relink(memories: Memory[], from: string, to: string): void {
|
||||
const pattern = new RegExp(`\\[\\[${escapeRegex(from)}\\]\\]`, 'g');
|
||||
for (const m of memories) if (pattern.test(m.content)) m.content = m.content.replace(pattern, `[[${to}]]`);
|
||||
}
|
||||
|
||||
private normalizeLeaf(name: string): string {
|
||||
return name.trim().toLowerCase().replace(/\s+/g, ' ');
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve a fact-agent proposed subject to an existing node when it's an alias/rename of one.
|
||||
* Exact match is checked first (cheap, and covers the common case since node names are
|
||||
* already normalized at creation time). Only falls through to fuzzy alias matching against
|
||||
* same-root candidates when there's no existing hit — i.e. only on likely-new-doc creation.
|
||||
*/
|
||||
private resolveSubject(subject: string, store: MemoryAccessor): string {
|
||||
const trimmed = subject.trim();
|
||||
const exact = store.find(trimmed);
|
||||
if (exact) return exact.name;
|
||||
|
||||
const normalized = this.normalizeLeaf(trimmed);
|
||||
const caseInsensitive = store.list.find(m => this.normalizeLeaf(m.name) === normalized);
|
||||
if (caseInsensitive) return caseInsensitive.name;
|
||||
|
||||
const root = trimmed.split('/')[0];
|
||||
const leaf = trimmed.split('/').slice(1).join('/') || trimmed;
|
||||
const candidates = store.list.filter(m => m.name.split('/')[0] === root && m.name !== trimmed);
|
||||
if (!candidates.length) return trimmed;
|
||||
|
||||
// fuzzyMatch requires >=2 terms; pad with an empty string when there's only one candidate
|
||||
const leaves = candidates.map(m => m.name.split('/').slice(1).join('/') || m.name);
|
||||
const probe = leaves.length > 1 ? leaves : [...leaves, ''];
|
||||
const {max, similarities} = this.llm.fuzzyMatch(leaf, ...probe);
|
||||
if (max >= ALIAS_MATCH_THRESHOLD) return candidates[similarities.indexOf(max)].name;
|
||||
|
||||
return trimmed;
|
||||
}
|
||||
|
||||
private async factAgent(conversation: string, store: MemoryAccessor, options: LLMRequest): Promise<FactAgentResult> {
|
||||
const ghosts = store.ghosts();
|
||||
|
||||
const response = await this.llm.ask(conversation, {
|
||||
model: options.model,
|
||||
temperature: 0.2,
|
||||
system: `You are a fact extractor for Obsidian-style knowledge vaults. Analyze the conversation and produce:
|
||||
|
||||
1. Journal recap (single paragraph)
|
||||
- "Captains Log" style record keeping
|
||||
- What was discussed/worked on, decisions, user's events/state/mood, general context
|
||||
- Leave empty only for trivial/empty exchanges/small talk
|
||||
|
||||
2. Fact buckets
|
||||
- ONLY facts the USER explicitly stated about themselves, their work, projects, or decisions made during this conversation
|
||||
- NEVER extract greetings, pleasantries, or anything the assistant itself said
|
||||
- Extract the final/end state, not deltas
|
||||
|
||||
Path assignment (entity) rules:
|
||||
- Use the owning entity of the fact (even if implied): "New bug on project 51 -> Projects/51"
|
||||
- When multiple facts relate to the same entity, pick a primary owner and wikilink related entities
|
||||
- Reuse existing entities when the owner already has a node
|
||||
- Always group under consistent entity roots (always plural):
|
||||
- Projects/[Name] for all initiatives
|
||||
- People/[Name] for all individuals
|
||||
- History/[Name] for all historical figures/events
|
||||
- Science/[Name] for all scientific concepts
|
||||
- Child entities nest under their parent entity:
|
||||
- Projects/51/Memory System, Projects/51/Bug-XYZ, not Bugs/51
|
||||
- Science/AI/Model-X, not Model-X/AI
|
||||
|
||||
Wikilink rules:
|
||||
- Use [[WikiLinks]] to connect related entities (e.g., [[Projects/51]], [[People/Robert]])
|
||||
- Only link specific, existing or implied entity paths — skip generic terms
|
||||
- Don't over-link: each link should add clarity or context, not noise
|
||||
|
||||
Available nodes:
|
||||
${this.listNodes(store.list).map(n => `- ${n.name}: ${n.description}`).join('\n') || 'None yet.'}
|
||||
${ghosts.length ? `${ghosts.map(g => `- ${g}: (Ghost)`).join('\n')}` : ''}`,
|
||||
schema: {
|
||||
journal: {type: 'string', description: 'Short day-to-day recap; empty if nothing happened.', required: false},
|
||||
buckets: {type: 'array', description: 'Groups of facts to remember; empty array if nothing worth storing.', items: {
|
||||
type: 'object', items: {
|
||||
subject: {type: 'string', description: 'Exact node name or new path (e.g. "People/Sarah", "Projects/Oxide")', required: true},
|
||||
facts: {type: 'array', description: 'Facts to store here', items: {type: 'string'}},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
const buckets = new Map<string, string[]>();
|
||||
for (const bucket of response.buckets ?? []) {
|
||||
const subject = bucket.subject.trim();
|
||||
const facts = buckets.get(subject) ?? [];
|
||||
facts.push(...dedupeFacts(bucket.facts));
|
||||
buckets.set(subject, facts);
|
||||
}
|
||||
|
||||
return {
|
||||
buckets: buckets.entries().toArray().map(([subject, facts]) => ({subject, facts})),
|
||||
journal: (response.journal ?? '').trim(),
|
||||
};
|
||||
}
|
||||
|
||||
private getWeekMonday(date: Date = new Date()): string {
|
||||
const d = new Date(Date.UTC(date.getFullYear(), date.getMonth(), date.getDate()));
|
||||
const day = d.getUTCDay();
|
||||
const diff = day === 0 ? -6 : 1 - day;
|
||||
d.setUTCDate(d.getUTCDate() + diff);
|
||||
return d.toISOString().slice(0, 10);
|
||||
}
|
||||
|
||||
private listNodes(memories: Memory[]): MemoryRef[] {
|
||||
return memories.map(m => ({name: m.name, description: m.description}));
|
||||
}
|
||||
|
||||
/** Finds the closest merge candidate via the KD tree's knn() instead of a manual O(n) cosine scan */
|
||||
private async checkMerge(node: Memory, memories: Memory[] | MemoryCache, options: LLMRequest, threshold = MERGE_THRESHOLD): Promise<Memory | null> {
|
||||
if (!node.embedding?.length || node.name.startsWith('Journal/')) return null;
|
||||
const store = this.access(memories);
|
||||
|
||||
const candidate = store.search(node.embedding, 5)
|
||||
.find(r => r.name !== node.name && !r.name.startsWith('Journal/') && r.distance !== undefined && r.distance <= threshold);
|
||||
if (!candidate) return null;
|
||||
const closest = store.find(candidate.name);
|
||||
if (!closest) return null;
|
||||
|
||||
const result = await this.mergeAgent(node, closest, options);
|
||||
const merged: Memory = {name: result.name, description: this.sanitizeDescription(result.description), content: '', embedding: [], links: [], backlinks: []};
|
||||
merged.content = this.touchHeader(merged, result.content);
|
||||
await embedMemoryFields(merged, this.llm);
|
||||
|
||||
this.relink(store.list, node.name, merged.name);
|
||||
this.relink(store.list, closest.name, merged.name);
|
||||
|
||||
this.queues.get(closest.name)?.request?.abort?.();
|
||||
this.queues.delete(closest.name);
|
||||
|
||||
store.forget(node.name);
|
||||
store.forget(closest.name);
|
||||
store.list.push(merged);
|
||||
store.commit();
|
||||
|
||||
return merged;
|
||||
}
|
||||
|
||||
private reconcile(node: Memory, memories: Memory[] | MemoryCache, options: LLMRequest): Promise<void> {
|
||||
const key = node.name;
|
||||
const existing = this.queues.get(key);
|
||||
if (existing) {
|
||||
existing.dirty = true;
|
||||
existing.request?.abort?.();
|
||||
return existing.task;
|
||||
}
|
||||
|
||||
const entry = {dirty: false, request: null, task: Promise.resolve()};
|
||||
this.queues.set(key, entry);
|
||||
const store = this.access(memories);
|
||||
entry.task = (async () => {
|
||||
let current = node, merged = false;
|
||||
try {
|
||||
do {
|
||||
entry.dirty = false;
|
||||
await this.docAgent(current, store.list, options, entry);
|
||||
this.mergeLock = this.mergeLock.then(() => this.checkMerge(current, memories, options));
|
||||
const result = await this.mergeLock;
|
||||
if (result) { current = result; merged = true; }
|
||||
} while (entry.dirty);
|
||||
} finally {
|
||||
store.commit(merged ? undefined : [node]);
|
||||
this.queues.delete(key);
|
||||
}
|
||||
})();
|
||||
return entry.task;
|
||||
}
|
||||
|
||||
private async docAgent(node: Memory, memories: Memory[], options: LLMRequest, entry: {request: {abort?: () => void} | null}): Promise<void> {
|
||||
if(!memories.includes(node)) return;
|
||||
const currentBody = stripHeader(node.content);
|
||||
let update;
|
||||
try {
|
||||
for (let i = 0; i < 2 && !update?.content; i++) {
|
||||
const request = this.llm.ask(currentBody, {
|
||||
model: options.model,
|
||||
temperature: 0.3,
|
||||
schema: {
|
||||
description: {type: 'string', description: 'One factual sentence describing the document\'s ENTIRE SUBJECT MATTER — for use as a search/merge fingerprint', required: true},
|
||||
content: {type: 'string', description: 'Rewritten document body in markdown, without the frontmatter block', required: true},
|
||||
},
|
||||
system: `You are a knowledge base editor maintaining one Obsidian-style document.
|
||||
|
||||
If it has a "## Pending" section, fold all new material into the appropriate part, resolve overlap, then remove the section entirely. If no section, just tidy per the rules below.
|
||||
|
||||
Use this loose structure, adapting headings to what the content needs:
|
||||
\`\`\`markdown
|
||||
${GENERIC_TEMPLATE}
|
||||
\`\`\`
|
||||
|
||||
Rules:
|
||||
- Contradictions: "## Pending" holds the newest information — bias toward it. Fold it in as the standing fact and drop the outdated statement, unless the old context adds meaningful nuance (e.g. "previously X, now Y"). This document should read as a source of truth, not an audit log
|
||||
- Journals (Journal/...): keep entries as a chronological timeline; clean up grammar within entries but never delete history
|
||||
- Use Obsidian markdown: # headings, **bold**, bullet/numbered lists, tables for 2D data
|
||||
- Link specific entities and concepts with [[WikiLink]] (e.g., [[Projects/KiwixServer]]); skip generics
|
||||
- Keep concise, factual, human-readable
|
||||
- NO frontmatter, filler, preamble, or AI commentary
|
||||
|
||||
Available nodes to link to (don't duplicate their content):
|
||||
${this.listNodes(memories).filter(n => n.name !== node.name).map(n => n.name).join(', ') || 'none'}
|
||||
|
||||
Current document:
|
||||
\`\`\`markdown
|
||||
${currentBody}
|
||||
\`\`\``,
|
||||
});
|
||||
entry.request = request;
|
||||
update = await request;
|
||||
}
|
||||
} catch (err: any) {
|
||||
if (err?.name === 'AbortError') return;
|
||||
throw err;
|
||||
} finally {
|
||||
entry.request = null;
|
||||
}
|
||||
|
||||
if (!update?.content) return;
|
||||
node.description = node.name !== 'People/User' ? this.sanitizeDescription(update.description) : 'All information about the current user';
|
||||
node.content = this.touchHeader(node, update.content);
|
||||
await embedMemoryFields(node, this.llm);
|
||||
}
|
||||
|
||||
private async mergeAgent(a: Memory, b: Memory, options: LLMRequest): Promise<{name: string, description: string, content: string}> {
|
||||
const modifiedOf = (m: Memory) => this.parseFrontmatter(m.content).fm.get('modified') || 'unknown';
|
||||
|
||||
return this.llm.ask('', {
|
||||
model: options.model,
|
||||
temperature: 0.3,
|
||||
schema: {
|
||||
name: {type: 'string', description: 'New path for the merged doc, collection/subject format (e.g. Projects/Oxide) — only reuse an old title if it\'s genuinely the best fit', required: true},
|
||||
description: {type: 'string', description: 'One factual sentence describing the merged document\'s subject matter', required: true},
|
||||
content: {type: 'string', description: 'Fully reconciled body in markdown, without frontmatter', required: true},
|
||||
},
|
||||
system: `You are a knowledge base editor merging two overlapping Obsidian documents into one.
|
||||
|
||||
Structure loosely:
|
||||
\`\`\`markdown
|
||||
${GENERIC_TEMPLATE}
|
||||
\`\`\`
|
||||
|
||||
Combine both documents, resolve duplication. On contradictions, bias toward whichever document was modified more recently; drop the outdated statement unless the old context adds meaningful nuance.
|
||||
|
||||
Document A ("${a.name}", last modified ${modifiedOf(a)}):
|
||||
\`\`\`markdown
|
||||
${stripHeader(a.content)}
|
||||
\`\`\`
|
||||
|
||||
Document B ("${b.name}", last modified ${modifiedOf(b)}):
|
||||
\`\`\`markdown
|
||||
${stripHeader(b.content)}
|
||||
\`\`\``,
|
||||
});
|
||||
}
|
||||
|
||||
private parseFrontmatter(content: string): {fm: Map<string, string>, body: string} {
|
||||
const match = content.match(/^---\n([\s\S]*?)\n---\n?([\s\S]*)$/);
|
||||
if (!match) return {fm: new Map(), body: content};
|
||||
const fm = new Map<string, string>();
|
||||
for (const line of match[1].split('\n')) {
|
||||
const i = line.indexOf(':');
|
||||
if (i === -1) continue;
|
||||
const key = line.slice(0, i).trim();
|
||||
const raw = line.slice(i + 1).trim();
|
||||
let value = raw;
|
||||
try { value = JSON.parse(raw); } catch { /* legacy unquoted value, keep raw */ }
|
||||
fm.set(key, value);
|
||||
}
|
||||
return {fm, body: match[2]};
|
||||
}
|
||||
|
||||
/**
|
||||
* Writes the code-owned frontmatter block. `body` is passed through stripHeader() first so a
|
||||
* model that ignores instructions and hallucinates its own `---` block can never corrupt or
|
||||
* duplicate the real frontmatter — the LLM only ever gets to influence the body.
|
||||
*/
|
||||
private touchHeader(node: Memory, body: string): string {
|
||||
const {fm} = this.parseFrontmatter(node.content);
|
||||
fm.set('name', node.name);
|
||||
fm.set('description', node.description || '');
|
||||
fm.set('modified', new Date().toISOString());
|
||||
return this.writeFrontmatter(fm, stripHeader(body));
|
||||
}
|
||||
|
||||
private writeFrontmatter(fm: Map<string, string>, body: string): string {
|
||||
const lines = [...fm.entries()].map(([k, v]) => `${k}: ${JSON.stringify(String(v).replace(/\s+/g, ' ').trim())}`);
|
||||
return `---\n${lines.join('\n')}\n---\n\n${body.trimStart()}`;
|
||||
}
|
||||
|
||||
decay() {
|
||||
for (const [name, ttl] of this.recentlyTouched) {
|
||||
if (ttl <= 1) this.recentlyTouched.delete(name);
|
||||
else this.recentlyTouched.set(name, ttl - 1);
|
||||
}
|
||||
}
|
||||
|
||||
touch(name: string, ttl = 2) {
|
||||
this.recentlyTouched.set(name, ttl);
|
||||
}
|
||||
|
||||
forget(name: string, memories: Memory[] | MemoryCache): boolean {
|
||||
return this.access(memories).forget(name);
|
||||
}
|
||||
|
||||
/** Ranks a candidate pool by weighted title/description/body similarity against the query embedding */
|
||||
private rankByFields(query: number[], candidates: Memory[], limit: number): Memory[] {
|
||||
const scored = candidates.map(m => {
|
||||
const titleSim = m.titleEmbedding?.length ? 1 - cosineDistance(query, m.titleEmbedding) : 0;
|
||||
const descSim = m.embedding?.length ? 1 - cosineDistance(query, m.embedding) : 0;
|
||||
const bodySim = m.bodyEmbeddings?.length
|
||||
? Math.max(...m.bodyEmbeddings.map(b => 1 - cosineDistance(query, b)))
|
||||
: 0;
|
||||
return {memory: m, score: titleSim * 0.5 + descSim * 0.35 + bodySim * 0.15};
|
||||
});
|
||||
return scored.sort((a, b) => b.score - a.score).slice(0, limit).map(s => s.memory);
|
||||
}
|
||||
|
||||
async recollect(query: string, memories: Memory[] | MemoryCache, limit = 5, graphDepth = 1): Promise<Memory[]> {
|
||||
const store = this.access(memories);
|
||||
if (!store.list.length) return [];
|
||||
|
||||
await store.backfillEmbeddings(this.llm);
|
||||
|
||||
const [e] = await this.llm.embedding(query);
|
||||
if (!e) return [];
|
||||
|
||||
// Description embedding is the cheap ANN index key; pull a wider pool then re-rank by field weight
|
||||
const pool = store.search(e.embedding, Math.max(limit * 3, limit));
|
||||
const poolMemories = pool.map(r => store.find(r.name)).filter((m): m is Memory => !!m);
|
||||
const ranked = this.rankByFields(e.embedding, poolMemories, limit);
|
||||
const found = new Set<string>(ranked.map(m => m.name));
|
||||
|
||||
if (graphDepth > 0) {
|
||||
let frontier = [...found];
|
||||
for (let depth = 0; depth < graphDepth && frontier.length; depth++) {
|
||||
const next: string[] = [];
|
||||
for (const name of frontier) {
|
||||
const node = store.find(name);
|
||||
if (!node) continue;
|
||||
for (const link of node.links) {
|
||||
if (!found.has(link) && store.find(link)) {
|
||||
found.add(link);
|
||||
next.push(link);
|
||||
}
|
||||
}
|
||||
}
|
||||
frontier = next;
|
||||
}
|
||||
}
|
||||
|
||||
const rankedOrder = ranked.map(m => m.name);
|
||||
const graphExpansions = [...found].filter(n => !rankedOrder.includes(n));
|
||||
return [...rankedOrder, ...graphExpansions].map(n => store.find(n)!).filter(Boolean);
|
||||
}
|
||||
|
||||
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest): Promise<Memory[]> {
|
||||
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 uid = `${Date.now()}_${Math.random().toString(36).slice(2)}`;
|
||||
const pending = {role: 'tool', name: 'memory_process', id: uid, content: conversation} as unknown as LLMMessage;
|
||||
history.push(pending);
|
||||
|
||||
const store = this.access(memories);
|
||||
const {buckets, journal} = await this.factAgent(conversation, store, options);
|
||||
const touched: Memory[] = [];
|
||||
|
||||
if (journal) {
|
||||
const journalName = `Journal/${this.getWeekMonday()}`;
|
||||
let jnode = store.find(journalName);
|
||||
if (!jnode) {
|
||||
jnode = {name: journalName, description: '', content: '', embedding: [], links: [], backlinks: []};
|
||||
store.list.push(jnode);
|
||||
}
|
||||
this.stage(jnode, `### ${new Date().toISOString().slice(0, 10)}\n${journal}`);
|
||||
touched.push(jnode);
|
||||
}
|
||||
|
||||
for (const {subject, facts} of buckets) {
|
||||
const resolved = this.resolveSubject(subject, store);
|
||||
let node = store.find(resolved);
|
||||
if (!node) {
|
||||
node = {name: resolved, description: '', content: '', embedding: [], links: [], backlinks: []};
|
||||
store.list.push(node);
|
||||
}
|
||||
this.stage(node, facts.map(f => `- ${f}`).join('\n'));
|
||||
touched.push(node);
|
||||
}
|
||||
|
||||
await Promise.all(touched.map(async node => {
|
||||
await embedMemoryFields(node, this.llm);
|
||||
this.touch(node.name);
|
||||
}));
|
||||
|
||||
if (touched.length) {
|
||||
store.commit(touched);
|
||||
(pending as any).content = `Saved to ${touched.map(n => `[[${n.name}]]`).join(', ')}`;
|
||||
Promise.all(touched.map(node => this.reconcile(node, memories, options).catch(() => {})));
|
||||
} else {
|
||||
(pending as any).content = 'Nothing worth remembering.';
|
||||
}
|
||||
|
||||
(touched as any).uid = uid;
|
||||
return touched;
|
||||
}
|
||||
|
||||
async reconcileAll(memories: Memory[] | MemoryCache, options: LLMRequest, scope: 'touched' | 'all' = 'touched'): Promise<void> {
|
||||
const store = this.access(memories);
|
||||
const targets = scope === 'all' ? store.list : store.list.filter(m => m.content.includes(PENDING_HEADING));
|
||||
await Promise.all(targets.map(node => this.reconcile(node, memories, options)));
|
||||
store.commit();
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,5 @@
|
||||
import {Memory, MemoryCache} from './memory.ts';
|
||||
import {MemoryCache} from './memory-state.ts';
|
||||
import type {Memory} from './memory.ts';
|
||||
|
||||
export type MemoryNode = {
|
||||
name: string;
|
||||
@@ -13,14 +14,6 @@ export function extractLinks(content: string): string[] {
|
||||
return [...new Set([...matches].map(m => m[1].trim()))];
|
||||
}
|
||||
|
||||
/**
|
||||
* Incrementally patch the graph for a set of changed memories, instead of
|
||||
* re-scanning every document. Only the changed memories' own content is
|
||||
* re-parsed for links; affected targets have their backlinks patched.
|
||||
* Does NOT handle node deletion — full rebuildGraph() is still required
|
||||
* when a memory is removed, since that needs a backlink sweep across
|
||||
* everyone who might reference it.
|
||||
*/
|
||||
export function patchGraph(mems: Memory[], nodes: MemoryNode[], changed: Memory[]): MemoryNode[] {
|
||||
const nameSet = new Set(mems.map(m => m.name));
|
||||
const byName = new Map(nodes.map(n => [n.name, n]));
|
||||
@@ -1,3 +1,5 @@
|
||||
import {cosineDistance, euclideanDistance} from '../utils.ts';
|
||||
|
||||
export type DistanceMetric = "euclidean" | "cosine";
|
||||
|
||||
export interface KDPoint<T = unknown> {
|
||||
@@ -18,28 +20,6 @@ interface KDNode<T> {
|
||||
deleted?: boolean;
|
||||
}
|
||||
|
||||
// ─── Distance helpers ─────────────────────────────────────────────────────────
|
||||
|
||||
function euclidean(a: number[], b: number[]): number {
|
||||
let sum = 0;
|
||||
for (let i = 0; i < a.length; i++) {
|
||||
const d = a[i] - b[i];
|
||||
sum += d * d;
|
||||
}
|
||||
return Math.sqrt(sum);
|
||||
}
|
||||
|
||||
function cosine(a: number[], b: number[]): number {
|
||||
let dot = 0, normA = 0, normB = 0;
|
||||
for (let i = 0; i < a.length; i++) {
|
||||
dot += a[i] * b[i];
|
||||
normA += a[i] * a[i];
|
||||
normB += b[i] * b[i];
|
||||
}
|
||||
const denom = Math.sqrt(normA) * Math.sqrt(normB);
|
||||
return denom === 0 ? 1 : 1 - dot / denom; // distance = 1 - similarity
|
||||
}
|
||||
|
||||
/**
|
||||
* Keeps the k closest candidates in memory, evicts the furthest when full
|
||||
*/
|
||||
@@ -123,7 +103,7 @@ export class KDTree<T = unknown> {
|
||||
points?: KDPoint<T>[]
|
||||
) {
|
||||
this.dims = dims;
|
||||
this.distanceFn = metric === "cosine" ? cosine : euclidean;
|
||||
this.distanceFn = metric === "cosine" ? cosineDistance : euclideanDistance;
|
||||
|
||||
if (points && points.length > 0) {
|
||||
this.validateAll(points);
|
||||
@@ -300,11 +280,8 @@ export class KDTree<T = unknown> {
|
||||
: [node.right, node.left];
|
||||
|
||||
this.searchKNN(near, query, k, heap, depth + 1);
|
||||
|
||||
// Only explore the far side if it could contain a closer point.
|
||||
// For cosine distance we can't prune by axis gap alone, so always explore.
|
||||
const shouldExplore =
|
||||
this.distanceFn === cosine
|
||||
this.distanceFn === cosineDistance
|
||||
? true
|
||||
: Math.abs(diff) < heap.worstDistance;
|
||||
|
||||
@@ -340,7 +317,7 @@ export class KDTree<T = unknown> {
|
||||
this.searchRadius(near, query, radius, results, depth + 1);
|
||||
|
||||
const shouldExplore =
|
||||
this.distanceFn === cosine ? true : Math.abs(diff) <= radius;
|
||||
this.distanceFn === cosineDistance ? true : Math.abs(diff) <= radius;
|
||||
|
||||
if (shouldExplore) {
|
||||
this.searchRadius(far, query, radius, results, depth + 1);
|
||||
@@ -0,0 +1,181 @@
|
||||
import {MemoryNode, patchGraph, rebuildGraph} from './graph.ts';
|
||||
import {KDTree} from './kd-tree.ts';
|
||||
import type {Memory, MemoryRef, MemoryStore} from './memory.ts';
|
||||
import {cosineDistance, embedMemoryFields} from '../utils.ts';
|
||||
|
||||
const TREE_TOMBSTONE_LIMIT = 0.25;
|
||||
|
||||
export function memoryStore(memories: MemoryStore): {
|
||||
list: Memory[];
|
||||
cache: MemoryCache | null;
|
||||
find: (name: string) => Memory | undefined;
|
||||
ghosts: () => string[];
|
||||
search: (vector: number[], limit: number) => MemoryRef[];
|
||||
forget: (name: string) => boolean;
|
||||
rebuild: (changed?: Memory[]) => MemoryNode[];
|
||||
backfillEmbeddings: (llm: any) => Promise<number>;
|
||||
} {
|
||||
if(memories instanceof MemoryCache) {
|
||||
return {
|
||||
list: memories.memories,
|
||||
cache: memories,
|
||||
find: name => memories.find(name),
|
||||
ghosts: () => memories.ghosts(),
|
||||
search: (vector, limit) => memories.search(vector, limit),
|
||||
forget: name => memories.remove(name),
|
||||
rebuild: changed => memories.rebuild(changed),
|
||||
backfillEmbeddings: llm => memories.backfillEmbeddings(llm),
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
list: memories,
|
||||
cache: null,
|
||||
find: name => memories.find(m => m.name === name),
|
||||
ghosts: () => rebuildGraph(memories).filter(n => n.missing).map(n => n.name),
|
||||
search: (vector, limit) => memories
|
||||
.filter(m => m.embedding?.length)
|
||||
.map(m => ({
|
||||
name: m.name,
|
||||
description: m.description,
|
||||
distance: cosineDistance(vector, m.embedding),
|
||||
}))
|
||||
.sort((a, b) => a.distance - b.distance)
|
||||
.slice(0, limit),
|
||||
forget: name => {
|
||||
const idx = memories.findIndex(m => m.name === name);
|
||||
if(idx === -1) return false;
|
||||
|
||||
memories.splice(idx, 1);
|
||||
return true;
|
||||
},
|
||||
rebuild: changed => rebuildGraph(memories),
|
||||
backfillEmbeddings: async llm => {
|
||||
const missing = memories.filter(m => !m.embedding?.length);
|
||||
if(!missing.length) return 0;
|
||||
|
||||
await Promise.all(missing.map(async node => {
|
||||
await embedMemoryFields(node, llm);
|
||||
}));
|
||||
|
||||
return missing.length;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
export class MemoryCache {
|
||||
private tree!: KDTree<MemoryRef>;
|
||||
private indexed = new Map<string, number[]>();
|
||||
public memories: Memory[];
|
||||
public nodes: MemoryNode[] = [];
|
||||
|
||||
get length() {
|
||||
return this.memories.length;
|
||||
}
|
||||
|
||||
constructor(memories: Memory[]) {
|
||||
this.memories = memories;
|
||||
this.tree = new KDTree<MemoryRef>(0);
|
||||
this.rebuild();
|
||||
}
|
||||
|
||||
find(name: string): Memory | undefined {
|
||||
return this.memories.find(m => m.name === name);
|
||||
}
|
||||
|
||||
private syncTree(): void {
|
||||
const current = new Set(this.memories.map(m => m.name));
|
||||
|
||||
for(const [name, emb] of [...this.indexed]) {
|
||||
const mem = this.memories.find(m => m.name === name);
|
||||
|
||||
if(!mem || !current.has(name) || mem.embedding !== emb) {
|
||||
this.tree.remove(p => p.name === name);
|
||||
this.indexed.delete(name);
|
||||
}
|
||||
}
|
||||
|
||||
for(const mem of this.memories) {
|
||||
if(!mem.embedding?.length || this.indexed.has(mem.name)) continue;
|
||||
|
||||
if(this.tree.dims === 0) {
|
||||
this.tree = new KDTree<MemoryRef>(mem.embedding.length, 'cosine');
|
||||
}
|
||||
|
||||
if(mem.embedding.length !== this.tree.dims) continue;
|
||||
|
||||
this.tree.insert({
|
||||
vector: mem.embedding,
|
||||
payload: {
|
||||
name: mem.name,
|
||||
description: mem.description,
|
||||
},
|
||||
});
|
||||
|
||||
this.indexed.set(mem.name, mem.embedding);
|
||||
}
|
||||
|
||||
if(this.tree.tombstoneRatio > TREE_TOMBSTONE_LIMIT) {
|
||||
this.tree.rebalance();
|
||||
}
|
||||
}
|
||||
|
||||
search(query: number[], limit: number): MemoryRef[] {
|
||||
if(!this.tree || this.tree.dims === 0) return [];
|
||||
|
||||
return this.tree.knn(query, limit).map(r => ({
|
||||
...r.point.payload,
|
||||
distance: r.distance,
|
||||
}));
|
||||
}
|
||||
|
||||
add(memory: Memory): void {
|
||||
this.memories.push(memory);
|
||||
this.rebuild([memory]);
|
||||
}
|
||||
|
||||
update(memory: Memory): void {
|
||||
const existing = this.find(memory.name);
|
||||
|
||||
if(existing) Object.assign(existing, memory);
|
||||
else this.memories.push(memory);
|
||||
|
||||
this.rebuild([existing ?? memory]);
|
||||
}
|
||||
|
||||
remove(name: string): boolean {
|
||||
const idx = this.memories.findIndex(m => m.name === name);
|
||||
if(idx === -1) return false;
|
||||
|
||||
this.memories.splice(idx, 1);
|
||||
this.rebuild();
|
||||
return true;
|
||||
}
|
||||
|
||||
ghosts(): string[] {
|
||||
return this.nodes.filter(n => n.missing).map(n => n.name);
|
||||
}
|
||||
|
||||
rebuild(changed?: Memory[]): MemoryNode[] {
|
||||
this.nodes = changed?.length && this.nodes.length
|
||||
? patchGraph(this.memories, this.nodes, changed)
|
||||
: rebuildGraph(this.memories);
|
||||
|
||||
this.syncTree();
|
||||
return this.nodes;
|
||||
}
|
||||
|
||||
commit(changed?: Memory[]): MemoryNode[] {
|
||||
return this.rebuild(changed);
|
||||
}
|
||||
|
||||
async backfillEmbeddings(llm: any): Promise<number> {
|
||||
const missing = this.memories.filter(m => !m.embedding?.length);
|
||||
if(!missing.length) return 0;
|
||||
|
||||
await Promise.all(missing.map(node => embedMemoryFields(node, llm)));
|
||||
this.commit(missing);
|
||||
|
||||
return missing.length;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,348 @@
|
||||
import {AiTool} from '../tools.ts';
|
||||
import type {LLMMessage, LLMRequest} from '../llm.ts';
|
||||
import {MemoryCache, memoryStore} from './memory-state.ts';
|
||||
import {cosineDistance, embedMemoryFields, stripHeader, updateMemory} from '../utils.ts';
|
||||
|
||||
const FACT_SIMILARITY_THRESHOLD = 0.62;
|
||||
const DUPLICATE_THRESHOLD = 0.68;
|
||||
const PROTECTED_MEMORIES = ['People/User'];
|
||||
const COLLECTION_WORDS = ['project', 'projects', 'people', 'person', 'managed', 'guides', 'guide', 'research', 'class', 'classes'];
|
||||
|
||||
export type Memory = {
|
||||
name: string;
|
||||
description: string;
|
||||
content: string;
|
||||
embedding: number[];
|
||||
titleEmbedding?: number[];
|
||||
bodyEmbeddings?: number[][];
|
||||
links: string[];
|
||||
backlinks: string[];
|
||||
}
|
||||
|
||||
export type MemoryRef = {
|
||||
name: string;
|
||||
description: string;
|
||||
distance?: number;
|
||||
}
|
||||
|
||||
export type MemoryOptions = {
|
||||
memory: Memory[] | MemoryCache;
|
||||
inject?: boolean;
|
||||
tool?: boolean;
|
||||
update?: boolean;
|
||||
maxTokens?: number;
|
||||
}
|
||||
|
||||
export type MemoryStore = Memory[] | MemoryCache;
|
||||
|
||||
/** Create an empty memory shell. */
|
||||
function emptyNode(name: string, description = ''): Memory {
|
||||
return {name, description, content: `# ${name.split('/').pop()}\n`, embedding: [], links: [], backlinks: []};
|
||||
}
|
||||
|
||||
function renderNode(node: Memory): string {
|
||||
return `### ${node.name}
|
||||
Description: ${node.description}
|
||||
Links: ${[...node.links, ...node.backlinks].join(', ') || 'none'}
|
||||
|
||||
\`\`\`markdown
|
||||
${node.content}
|
||||
\`\`\``;
|
||||
}
|
||||
|
||||
function factSimilarity(a: Memory, b: Memory): number {
|
||||
return !a.bodyEmbeddings?.length || !b.bodyEmbeddings?.length ? 0 : Math.max(...a.bodyEmbeddings.flatMap(av => b.bodyEmbeddings!.map(bv => 1 - cosineDistance(av, bv))));
|
||||
}
|
||||
|
||||
function words(text: string): string[] {
|
||||
return [...new Set(text.toLowerCase().replace(/[[\]()/_-]/g, ' ').replace(/[^a-z0-9\s]/g, '').split(/\s+/).filter(w => w && !COLLECTION_WORDS.includes(w)))];
|
||||
}
|
||||
|
||||
function jaccard(a: string[], b: string[]): number {
|
||||
const bs = new Set(b), hit = a.filter(x => bs.has(x)).length, total = new Set([...a, ...b]).size;
|
||||
return total ? hit / total : 0;
|
||||
}
|
||||
|
||||
function duplicateScore(a: Memory, b: Memory): number {
|
||||
const name = Math.max(
|
||||
jaccard(words(a.name), words(b.name)),
|
||||
jaccard(words(a.name.split('/').pop() || a.name), words(b.name.split('/').pop() || b.name)),
|
||||
);
|
||||
const desc = jaccard(words(a.description), words(b.description));
|
||||
const body = factSimilarity(a, b);
|
||||
const emb = a.embedding?.length && b.embedding?.length && a.embedding.length === b.embedding.length ? 1 - cosineDistance(a.embedding, b.embedding) : 0;
|
||||
return Math.max(body, name * 0.9 + desc * 0.06 + emb * 0.04, emb * 0.55 + name * 0.35 + desc * 0.1);
|
||||
}
|
||||
|
||||
function homeScore(node: Memory): number {
|
||||
return (PROTECTED_MEMORIES.includes(node.name) ? 1e9 : 0)
|
||||
+ (node.name.includes('/') ? 4 : 0)
|
||||
+ (node.description && node.description !== 'Persistent memory document' ? 1 : 0)
|
||||
+ Math.min(stripHeader(node.content).length / 1000, 5);
|
||||
}
|
||||
|
||||
function pickMerge(a: Memory, b: Memory, touched: Set<string>): [drop: Memory, home: Memory] {
|
||||
const as = homeScore(a), bs = homeScore(b);
|
||||
if(touched.has(a.name) && !touched.has(b.name)) return as > bs + 2 ? [b, a] : [a, b];
|
||||
if(touched.has(b.name) && !touched.has(a.name)) return bs > as + 2 ? [a, b] : [b, a];
|
||||
return as <= bs ? [a, b] : [b, a];
|
||||
}
|
||||
|
||||
/** Build memory tools and memory index text. */
|
||||
export function memoryTools(llm: any, memories: MemoryStore): {tools: AiTool[]; list: string} {
|
||||
const store = memoryStore(memories);
|
||||
const names = new Map<string, string>();
|
||||
|
||||
for(const node of store.list)
|
||||
if(!names.has(node.name)) names.set(node.name, `${node.name} - ${node.description}`);
|
||||
|
||||
for(const name of store.ghosts())
|
||||
if(!names.has(name)) names.set(name, `${name} - ghost node`);
|
||||
|
||||
return {
|
||||
list: [...names.values()].join('\n'),
|
||||
tools: [
|
||||
{
|
||||
name: 'memory_search',
|
||||
description: 'Semantically search memories for most relevant',
|
||||
args: {
|
||||
query: {type: 'string', description: 'Search query', required: true},
|
||||
limit: {type: 'number', description: 'Maximum results, default 5', default: 5},
|
||||
},
|
||||
fn: async ({query, limit = 5}) => {
|
||||
if(!query?.trim()) return 'Search query is required.';
|
||||
const [chunk] = await llm.embedding(query, {maxTokens: 8000, overlapTokens: 0});
|
||||
if(!chunk?.embedding) return 'Failed to create embedding from query';
|
||||
const results = store.search(chunk.embedding, limit).map(ref => store.find(ref.name)).filter((node): node is Memory => !!node);
|
||||
return results.length ? results.map(renderNode).join('\n\n---\n\n') : 'No relevant memories found.';
|
||||
},
|
||||
},
|
||||
{
|
||||
name: 'memory_read',
|
||||
description: 'Read an entire memory document by name',
|
||||
args: {name: {type: 'string', description: 'Exact document name', required: true}},
|
||||
fn: async ({name}) => {
|
||||
const node = store.find(name);
|
||||
return node ? renderNode(node) : store.ghosts().includes(name) ? `"${name}" is a ghost node with no document of its own.` : `Not found: "${name}".`;
|
||||
},
|
||||
},
|
||||
{
|
||||
name: 'memory_delete',
|
||||
description: 'Delete a duplicate or merged memory',
|
||||
args: {name: {type: 'string', description: 'Exact document name', required: true}},
|
||||
fn: async ({name}) => {
|
||||
store.forget(name);
|
||||
return `Removed: ${name}`;
|
||||
},
|
||||
},
|
||||
{
|
||||
name: 'memory_write',
|
||||
description: 'Create or replace a memory document.',
|
||||
args: {
|
||||
name: {type: 'string', description: 'Document name following the entity naming convention.', required: true},
|
||||
description: {type: 'string', description: 'One factual sentence describing the entire document subject', required: true},
|
||||
content: {type: 'string', description: 'Complete Markdown document body, including the # title', required: true},
|
||||
},
|
||||
fn: async (args: any) => {
|
||||
const name = String(args.name || '').trim();
|
||||
if(!name) return 'A document name is required.';
|
||||
const description = String(args.description || '').trim();
|
||||
if(!description) return 'A document description is required.';
|
||||
const content = String(args.content || '').trim();
|
||||
if(!content) return 'Document content is required.';
|
||||
|
||||
let node = store.find(name);
|
||||
if(!node) {
|
||||
node = emptyNode(name, description);
|
||||
if(store.cache) store.cache.add(node);
|
||||
else store.list.push(node);
|
||||
}
|
||||
|
||||
node.description = name === 'People/User' ? 'All information about the current user' : description.replace(/\s+/g, ' ').trim();
|
||||
node.content = updateMemory(node, content);
|
||||
await embedMemoryFields(node, llm);
|
||||
store.cache?.commit([node]);
|
||||
return `Updated ${name}`;
|
||||
},
|
||||
},
|
||||
],
|
||||
};
|
||||
}
|
||||
|
||||
export class MemoryManager {
|
||||
private memorized = new WeakMap<LLMMessage[], LLMMessage>();
|
||||
|
||||
constructor(private llm: any) {}
|
||||
|
||||
static normalize(memory?: Memory[] | MemoryCache | MemoryOptions): MemoryOptions | null {
|
||||
if(!memory) return null;
|
||||
if(Array.isArray(memory) || memory instanceof MemoryCache) return {memory, inject: true, tool: false, update: false};
|
||||
if(typeof memory === 'object' && 'memory' in memory) return {inject: true, tool: false, update: false, ...memory};
|
||||
return null;
|
||||
}
|
||||
|
||||
private memorySystem(list: string): string {
|
||||
return `You maintain notes written in markdown used for memories from recent conversations using your tools.
|
||||
Only preserve durable information worth remembering established by the USER.
|
||||
Do not store assistant guesses, speculation, suggestions, commentary, temporary state, or details that are not worth remembering.
|
||||
|
||||
## Rules
|
||||
- ALWAYS READ a target memory before changing it, \`memory_write\` does a full replace, it DOES NOT append!
|
||||
- Memories should contain the final state, not deltas
|
||||
- New conversational context is authoritative when it contracts existing information; reconcile it
|
||||
- Only remove information when stale, contradicted or duplicated; always preserve existing information, formatting and keep related information together
|
||||
- Only merge memories when two or more nodes are clearly about the same thing; only split a memory when it is clearly about two distinct subjects
|
||||
- Use [[WikiLinks]] liberally to record aliases and relationships between entities, even ones without pages yet (ghost nodes)
|
||||
- Use headings, subheadings, lists, tables and other markdown formatting to make documents clean
|
||||
- Maintain a \`## Todo List\` of checkboxes AS THE FIRST SUBHEADING when an entity has tasks
|
||||
- Only create todo items for USER tasks, not AI work
|
||||
- Only store each in one place, no duplicates
|
||||
- Use \`People/User\` for personal tasks or as a fallback
|
||||
|
||||
## Naming
|
||||
- Every fact should be grouped with the owning entity
|
||||
- Always follow the naming convention \`Collection/(Pro)Noun\`
|
||||
- Facts about the user belong under People/User
|
||||
- Reuse existing memories when they are clearly the same entity including aliases and ghost references.
|
||||
- Only create deeper paths when there is a real parent/child entity relationship: \`School/Class/Chapter\`
|
||||
|
||||
Valid Examples:
|
||||
- People/User
|
||||
- People/John Smith
|
||||
- Projects/Momentum
|
||||
- Projects/Momentum/Marketing
|
||||
- Research/Object Recognition
|
||||
- Guides/HAM Radio SOP
|
||||
|
||||
## Workflow
|
||||
|
||||
1. Create groups of durable information and todos based on the owning entity & naming rules above
|
||||
2. For each group:
|
||||
1. Read the existing memory(s)
|
||||
2. Merge the information & todos based on the rules above
|
||||
3. Write the entire patched document
|
||||
|
||||
Available memories:
|
||||
|
||||
${list || 'No memory documents exist yet.'}`;
|
||||
}
|
||||
|
||||
private touchedNames(history: LLMMessage[]): string[] {
|
||||
return [...new Set(history
|
||||
.filter((h: any) => h.role === 'tool' && h.name === 'memory_write' && !h.error)
|
||||
.map((h: any) => String(h.args?.name || h.content?.match(/^Updated (.+)$/)?.[1] || '').trim())
|
||||
.filter(Boolean))];
|
||||
}
|
||||
|
||||
private async backfillEmbeddings(store: ReturnType<typeof memoryStore>): Promise<void> {
|
||||
const missing = store.list.filter(m => !m.embedding?.length || !m.titleEmbedding?.length || !m.bodyEmbeddings?.length);
|
||||
await Promise.all(missing.map(m => embedMemoryFields(m, this.llm)));
|
||||
store.cache?.commit(missing);
|
||||
}
|
||||
|
||||
private closestDuplicate(node: Memory, store: ReturnType<typeof memoryStore>): Memory | null {
|
||||
return store.list
|
||||
.filter(m => m.name !== node.name && !m.name.startsWith('Journal/') && !node.name.startsWith('Journal/'))
|
||||
.map(m => ({node: m, score: duplicateScore(node, m)}))
|
||||
.filter(x => x.score >= DUPLICATE_THRESHOLD || factSimilarity(node, x.node) >= FACT_SIMILARITY_THRESHOLD)
|
||||
.sort((a, b) => b.score - a.score)[0]?.node || null;
|
||||
}
|
||||
|
||||
private async rehomeDeleted(drop: Memory, home: Memory, memories: MemoryStore, options: LLMRequest): Promise<void> {
|
||||
const store = memoryStore(memories);
|
||||
const backup = structuredClone(drop);
|
||||
store.forget(drop.name);
|
||||
|
||||
try {
|
||||
const memory = memoryTools(this.llm, memories);
|
||||
await this.llm.ask(`A duplicate memory document was removed automatically.
|
||||
|
||||
Deleted document:
|
||||
${renderNode(backup)}
|
||||
|
||||
Closest surviving home:
|
||||
${renderNode(home)}
|
||||
|
||||
Reinsert every durable unique fact, useful relationship, alias, and user todo from the deleted document into the best remaining memory document.
|
||||
Usually this should be "${home.name}", but use another existing memory if it is a better home.
|
||||
Read before writing. Write full replacement documents only.
|
||||
Do NOT recreate "${backup.name}" unless the deletion was wrong and it is clearly a distinct persistent entity.`, {
|
||||
model: options.memoryModel || options.model,
|
||||
temperature: 0.2,
|
||||
maxTokens: options.maxTokens,
|
||||
tools: memory.tools,
|
||||
history: [],
|
||||
system: this.memorySystem(memory.list),
|
||||
});
|
||||
} catch(err) {
|
||||
if(!store.find(backup.name)) store.cache ? store.cache.add(backup) : store.list.push(backup);
|
||||
throw err;
|
||||
} finally {
|
||||
store.cache?.commit(store.list);
|
||||
}
|
||||
}
|
||||
|
||||
private async reconcileSimilar(history: LLMMessage[], memories: MemoryStore, options: LLMRequest): Promise<void> {
|
||||
const store = memoryStore(memories);
|
||||
const touched = new Set(this.touchedNames(history));
|
||||
const targets = store.list.filter(m => touched.has(m.name) || [...touched].some(t => duplicateScore(m, store.find(t) || m) >= DUPLICATE_THRESHOLD));
|
||||
const deleted = new Set<string>();
|
||||
if(!targets.length) return;
|
||||
await this.backfillEmbeddings(store);
|
||||
|
||||
for(const node of targets) {
|
||||
if(!store.find(node.name) || deleted.has(node.name) || PROTECTED_MEMORIES.includes(node.name)) continue;
|
||||
const closest = this.closestDuplicate(node, store);
|
||||
if(!closest) continue;
|
||||
const [drop, home] = pickMerge(node, closest, touched);
|
||||
if(deleted.has(drop.name) || PROTECTED_MEMORIES.includes(drop.name)) continue;
|
||||
deleted.add(drop.name);
|
||||
await this.rehomeDeleted(drop, home, memories, options);
|
||||
await this.backfillEmbeddings(store);
|
||||
}
|
||||
}
|
||||
|
||||
async recollect(query: string, memory: MemoryStore, limit = 15): Promise<Memory[]> {
|
||||
const store = memoryStore(memory);
|
||||
if(!store.list.length || !query?.trim()) return [];
|
||||
const [chunk] = await this.llm.embedding(query, {maxTokens: 8000, overlapTokens: 0});
|
||||
return !chunk?.embedding ? [] : store.search(chunk.embedding, limit).map(ref => store.find(ref.name)).filter((m: Memory | undefined): m is Memory => !!m);
|
||||
}
|
||||
|
||||
get tools(): {read: (memory: MemoryStore) => AiTool[]} {
|
||||
return {read: (memory: MemoryStore) => memoryTools(this.llm, memory).tools};
|
||||
}
|
||||
|
||||
async memorize(history: LLMMessage[], memories: Memory[] | MemoryCache, options: LLMRequest = {},): Promise<Memory[]> {
|
||||
const store = memoryStore(memories);
|
||||
const previous = this.memorized.get(history);
|
||||
let start = 0;
|
||||
|
||||
if(previous) {
|
||||
const index = history.indexOf(previous);
|
||||
if(index >= 0) start = index + 1;
|
||||
}
|
||||
|
||||
const turns = history.slice(start).filter((h: any) => h.role === 'user' || h.role === 'assistant');
|
||||
const conversation = turns.map((h: any) => `[${h.role}]: ${h.content}`).join('\n\n').trim();
|
||||
if(!conversation) return store.list;
|
||||
|
||||
const memory = memoryTools(this.llm, memories);
|
||||
const memoryHistory: LLMMessage[] = [];
|
||||
|
||||
await this.llm.ask(conversation, {
|
||||
model: options.memoryModel || options.model,
|
||||
temperature: 0.2,
|
||||
maxTokens: options.maxTokens,
|
||||
tools: memory.tools,
|
||||
history: memoryHistory,
|
||||
system: this.memorySystem(memory.list),
|
||||
});
|
||||
|
||||
await this.reconcileSimilar(memoryHistory, memories, options);
|
||||
|
||||
const lastTurn = turns.at(-1);
|
||||
if(lastTurn) this.memorized.set(history, lastTurn);
|
||||
return store.list;
|
||||
}
|
||||
}
|
||||
+140
-135
@@ -36,21 +36,49 @@ export class OpenAi extends LLMProvider {
|
||||
private toWire(history: LLMMessage[], system?: string): any[] {
|
||||
const wire: any[] = [];
|
||||
if(system) wire.push({role: 'system', content: system});
|
||||
for(const h of history) {
|
||||
if(h.role === 'tool') {
|
||||
wire.push({
|
||||
role: 'assistant',
|
||||
content: null,
|
||||
tool_calls: [{id: h.id, type: 'function', function: {name: h.name, arguments: JSON.stringify(h.args)}],
|
||||
}, {
|
||||
role: 'tool',
|
||||
tool_call_id: h.id,
|
||||
content: h.error || h.content || '',
|
||||
});
|
||||
} else {
|
||||
|
||||
for(let i = 0; i < history.length; i++) {
|
||||
const h = history[i];
|
||||
|
||||
if(h.role !== 'tool') {
|
||||
wire.push({role: h.role, content: this.toWireContent(h.content)});
|
||||
continue;
|
||||
}
|
||||
|
||||
const calls: any[] = [];
|
||||
const results: any[] = [];
|
||||
|
||||
while(i < history.length && history[i].role === 'tool') {
|
||||
const tool: any = history[i];
|
||||
|
||||
calls.push({
|
||||
id: tool.id,
|
||||
type: 'function',
|
||||
function: {
|
||||
name: tool.name,
|
||||
arguments: JSON.stringify(tool.args || {})
|
||||
}
|
||||
});
|
||||
|
||||
results.push({
|
||||
role: 'tool',
|
||||
tool_call_id: tool.id,
|
||||
content: tool.error || tool.content || ''
|
||||
});
|
||||
|
||||
i++;
|
||||
}
|
||||
|
||||
wire.push({
|
||||
role: 'assistant',
|
||||
content: null,
|
||||
tool_calls: calls
|
||||
});
|
||||
|
||||
wire.push(...results);
|
||||
i--;
|
||||
}
|
||||
|
||||
return wire;
|
||||
}
|
||||
|
||||
@@ -60,13 +88,12 @@ export class OpenAi extends LLMProvider {
|
||||
if(!options.history) options.history = [];
|
||||
const history = options.history;
|
||||
if(message) history.push({role: 'user', content: message, timestamp: Date.now()});
|
||||
|
||||
const tools = options.tools || this.ai.options.llm?.tools || [];
|
||||
const requestParams: any = {
|
||||
model: options.model || this.model,
|
||||
stream: !!options.stream,
|
||||
max_completion_tokens: options.maxTokens || this.ai.options.llm?.maxTokens || undefined,
|
||||
temperature: options.temperature || this.ai.options.llm?.temperature || undefined,
|
||||
max_completion_tokens: options.maxTokens ?? this.ai.options.llm?.maxTokens,
|
||||
temperature: options.temperature ?? this.ai.options.llm?.temperature,
|
||||
tools: tools.map(t => ({
|
||||
type: 'function',
|
||||
function: {
|
||||
@@ -74,8 +101,12 @@ export class OpenAi extends LLMProvider {
|
||||
description: t.description,
|
||||
parameters: {
|
||||
type: 'object',
|
||||
properties: t.args ? objectMap(t.args, (key, value) => ({...value, required: undefined})) : {},
|
||||
required: t.args ? Object.entries(t.args).filter(t => t[1].required).map(t => t[0]) : []
|
||||
properties: t.args
|
||||
? objectMap(t.args, (key, value) => ({...value, required: undefined}))
|
||||
: {},
|
||||
required: t.args
|
||||
? Object.entries(t.args).filter(t => t[1].required).map(t => t[0])
|
||||
: []
|
||||
}
|
||||
}
|
||||
}))
|
||||
@@ -83,160 +114,134 @@ export class OpenAi extends LLMProvider {
|
||||
|
||||
if(options.schema) {
|
||||
const schema = convertSchema(options.schema);
|
||||
requestParams.response_format = {type: 'json_schema', json_schema: {name: 'response', strict: true, schema}};
|
||||
requestParams.response_format = {
|
||||
type: 'json_schema',
|
||||
json_schema: {name: 'response', strict: true, schema}
|
||||
};
|
||||
}
|
||||
if(options.stream) requestParams.stream_options = {include_usage: true};
|
||||
|
||||
try {
|
||||
let terminal = false;
|
||||
let iteration = 0;
|
||||
|
||||
do {
|
||||
iteration++;
|
||||
requestParams.messages = this.toWire(history.filter(h => h.role !== 'system'), options.system);
|
||||
|
||||
const callStart = Date.now();
|
||||
const resp: any = await this.tokenPool.run(token => this.getClient(token).chat.completions.create(requestParams)).catch(err => {
|
||||
const resp: any = await this.tokenPool.run(token =>
|
||||
this.getClient(token).chat.completions.create(requestParams)
|
||||
).catch(err => {
|
||||
err.message += `\n\nMessages:\n${JSON.stringify(requestParams.messages, null, 2)}`;
|
||||
throw err;
|
||||
});
|
||||
|
||||
let usage: any, msg: any = {content: '', tool_calls: []};
|
||||
let usage: any;
|
||||
let finishReason: string | undefined;
|
||||
let msg: any = {content: '', tool_calls: []};
|
||||
let streamedChars = 0;
|
||||
|
||||
if(options.stream) {
|
||||
const toolCallState: Record<number, any> = {};
|
||||
let pendingToolCalls = 0;
|
||||
let streamCompleted = false;
|
||||
try {
|
||||
for await (const chunk of resp) {
|
||||
if(controller.signal.aborted) break;
|
||||
if(chunk.usage) usage = chunk.usage;
|
||||
|
||||
for await (const chunk of resp) {
|
||||
if(controller.signal.aborted) break;
|
||||
const choice = chunk.choices?.[0];
|
||||
if(choice?.finish_reason) finishReason = choice.finish_reason;
|
||||
|
||||
if(chunk.usage) usage = chunk.usage;
|
||||
if(choice?.delta?.content) {
|
||||
msg.content += choice.delta.content;
|
||||
streamedChars += choice.delta.content.length;
|
||||
options.stream({text: choice.delta.content});
|
||||
}
|
||||
|
||||
// Handle content deltas
|
||||
if(chunk.choices[0]?.delta?.content) {
|
||||
msg.content += chunk.choices[0].delta.content;
|
||||
options.stream?.({text: chunk.choices[0].delta.content});
|
||||
}
|
||||
if(choice?.delta?.tool_calls) {
|
||||
for(const deltaTC of choice.delta.tool_calls) {
|
||||
const index = deltaTC.index ?? msg.tool_calls.length;
|
||||
let existing = msg.tool_calls.find((tc: any) => tc.index === index);
|
||||
|
||||
// Handle tool_call deltas
|
||||
if(chunk.choices[0]?.delta?.tool_calls) {
|
||||
for(const deltaTC of chunk.choices[0].delta.tool_calls) {
|
||||
const tc = toolCallState[deltaTC.index];
|
||||
if(!existing) {
|
||||
existing = {index, id: '', function: {name: '', arguments: ''}};
|
||||
msg.tool_calls.push(existing);
|
||||
}
|
||||
|
||||
if(!tc && deltaTC.id) {
|
||||
// New tool call delta
|
||||
toolCallState[deltaTC.index] = {
|
||||
id: deltaTC.id,
|
||||
name: deltaTC.function?.name || '',
|
||||
arguments: deltaTC.function?.arguments || '',
|
||||
complete: false
|
||||
};
|
||||
msg.tool_calls.push({
|
||||
index: deltaTC.index,
|
||||
id: deltaTC.id || '',
|
||||
function: {name: deltaTC.function?.name || '', arguments: deltaTC.function?.arguments || ''}
|
||||
});
|
||||
pendingToolCalls++;
|
||||
} else if(tc && deltaTC.function?.name) {
|
||||
// Update existing tool call
|
||||
if(deltaTC.id) tc.id = deltaTC.id;
|
||||
if(deltaTC.function.name) tc.name = deltaTC.function.name;
|
||||
if(deltaTC.function.arguments) tc.arguments += deltaTC.function.arguments;
|
||||
}
|
||||
|
||||
if(tc && deltaTC.function?.arguments && !tc.complete) {
|
||||
tc.complete = true;
|
||||
pendingToolCalls--;
|
||||
if(deltaTC.id) existing.id = deltaTC.id;
|
||||
if(deltaTC.function?.name) existing.function.name = deltaTC.function.name;
|
||||
if(deltaTC.function?.arguments) existing.function.arguments += deltaTC.function.arguments;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Execute completed tools immediately when all deltas are received
|
||||
if(!pendingToolCalls && msg.tool_calls.length > 0) {
|
||||
const completedTools: any[] = [];
|
||||
|
||||
for(const tc of msg.tool_calls) {
|
||||
const tcState = toolCallState[tc.index];
|
||||
if(!tcState) continue;
|
||||
|
||||
const entry: any = {role: 'tool', id: tc.id, name: tcState.name, args: JSONAttemptParse(tcState.arguments, {}), content: undefined, timestamp: Date.now()};
|
||||
history.push(entry);
|
||||
completedTools.push({tc, entry});
|
||||
|
||||
const tool = tools.find(findByProp('name', tcState.name));
|
||||
if(options.stream) options.stream?.({tool: tcState.name});
|
||||
|
||||
if(!tool) {
|
||||
entry.error = 'Tool not found';
|
||||
continue;
|
||||
}
|
||||
|
||||
try {
|
||||
const toolStream = options.stream && ((chunk: any) => {
|
||||
if(chunk.done) { terminal = true; return; }
|
||||
options.stream?.(chunk);
|
||||
});
|
||||
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
|
||||
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
|
||||
} catch(err: any) {
|
||||
entry.error = err?.message || err?.toString() || 'Unknown';
|
||||
}
|
||||
}
|
||||
|
||||
if(options.stream) options.stream?.({done: true});
|
||||
|
||||
for(const {entry} of completedTools) {
|
||||
if(entry.error) {
|
||||
options.stream?.({error: entry.error});
|
||||
} else if(entry.content) {
|
||||
options.stream?.({toolResult: true});
|
||||
}
|
||||
}
|
||||
|
||||
msg.tool_calls = []; // Clear after execution
|
||||
}
|
||||
streamCompleted = true;
|
||||
} catch(err) {
|
||||
if(!controller.signal.aborted) throw err;
|
||||
}
|
||||
|
||||
// Handle finish_reason for proper termination detection
|
||||
if(chunk.choices[0]?.finish_reason) {
|
||||
const finishReason = chunk.choices[0].finish_reason;
|
||||
if(finishReason === 'stop' || finishReason === 'tool_calls' || finishReason === 'length') {
|
||||
terminal = true;
|
||||
}
|
||||
}
|
||||
if(streamCompleted && !finishReason) finishReason = msg.tool_calls.length ? 'tool_calls' : 'stop';
|
||||
} else {
|
||||
usage = resp.usage;
|
||||
finishReason = resp.choices[0].finish_reason;
|
||||
msg = resp.choices[0].message;
|
||||
if(msg.tool_calls) {
|
||||
// Non-streaming: execute tools immediately
|
||||
for(const tc of msg.tool_calls) {
|
||||
const entry: any = {role: 'tool', id: tc.id, name: tc.function.name, args: JSONAttemptParse(tc.function.arguments, {}), content: undefined, timestamp: Date.now()};
|
||||
history.push(entry);
|
||||
|
||||
const tool = tools.find(findByProp('name', tc.function.name));
|
||||
if(!tool) {
|
||||
entry.error = 'Tool not found';
|
||||
continue;
|
||||
}
|
||||
|
||||
try {
|
||||
const result = await tool.fn(entry.args, undefined, this.ai, tc.id);
|
||||
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
|
||||
} catch(err: any) {
|
||||
entry.error = err?.message || err?.toString() || 'Unknown';
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const duration = Date.now() - callStart;
|
||||
const tps = usage?.completion_tokens && duration > 0 ? usage.completion_tokens / (duration / 1000) : 0;
|
||||
|
||||
// Capture assistant messages (before or after tools)
|
||||
if(msg.content?.trim()) {
|
||||
history.push({role: 'assistant', content: msg.content.trim(), timestamp: Date.now(), duration, tps});
|
||||
if(finishReason === 'length' && !controller.signal.aborted) {
|
||||
if(msg.content?.trim()) history.push({role: 'assistant', content: msg.content.trim(), timestamp: Date.now(), duration, tps});
|
||||
throw new Error(`[OpenAI] Response hit token limit before completing`);
|
||||
}
|
||||
|
||||
if(!finishReason && !controller.signal.aborted) {
|
||||
throw new Error('[OpenAI] Completion ended without a usable response');
|
||||
}
|
||||
|
||||
const toolCalls = msg.tool_calls || [];
|
||||
|
||||
if(toolCalls.length && !controller.signal.aborted) {
|
||||
if(msg.content?.trim()) history.push({role: 'assistant', content: msg.content.trim(), timestamp: Date.now(), duration, tps});
|
||||
|
||||
const entries = toolCalls.map((tc: any) => {
|
||||
const entry: any = {
|
||||
role: 'tool',
|
||||
id: tc.id,
|
||||
name: tc.function.name,
|
||||
args: JSONAttemptParse(tc.function.arguments, {}),
|
||||
content: undefined,
|
||||
timestamp: Date.now()
|
||||
};
|
||||
|
||||
history.push(entry);
|
||||
return {tc, entry};
|
||||
});
|
||||
|
||||
await Promise.all(entries.map(async ({tc, entry}: any) => {
|
||||
const tool = tools.find(findByProp('name', tc.function.name));
|
||||
if(options.stream) options.stream({tool: tc.function.name});
|
||||
if(!tool) return entry.error = 'Tool not found';
|
||||
try {
|
||||
const toolStream = options.stream && ((chunk: any) => {
|
||||
if(chunk.done) return;
|
||||
options.stream!(chunk);
|
||||
});
|
||||
|
||||
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);
|
||||
entry.content = typeof result === 'object' ? JSONSanitize(result) : result;
|
||||
} catch(err: any) {
|
||||
entry.error = err?.message || err?.toString() || 'Unknown';
|
||||
}
|
||||
}));
|
||||
} else {
|
||||
terminal = true;
|
||||
const text = (msg.content || '').trim();
|
||||
if(text) history.push({role: 'assistant', content: text, timestamp: Date.now(), duration, tps});
|
||||
}
|
||||
} while(!terminal && !controller.signal.aborted);
|
||||
|
||||
if(options.stream) options.stream?.({done: true});
|
||||
|
||||
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) => h.role === 'assistant' ? str + (h.content || '') : str, '').trim();
|
||||
res(options.schema ? JSONAttemptParse(finalContent, finalContent) : finalContent);
|
||||
@@ -245,4 +250,4 @@ export class OpenAi extends LLMProvider {
|
||||
}
|
||||
}), {abort: () => controller.abort()});
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
import {Memory} from './memory/memory.ts';
|
||||
|
||||
export function cosineDistance(a: number[], b: number[]): number {
|
||||
let dot = 0, normA = 0, normB = 0;
|
||||
for(let i = 0; i < a.length; i++) {
|
||||
dot += a[i] * b[i];
|
||||
normA += a[i] * a[i];
|
||||
normB += b[i] * b[i];
|
||||
}
|
||||
const denom = Math.sqrt(normA) * Math.sqrt(normB);
|
||||
return denom === 0 ? 1 : 1 - dot / denom;
|
||||
}
|
||||
|
||||
export async function embedMemoryFields(node: Memory, llm: any): Promise<void> {
|
||||
const body = stripHeader(node.content);
|
||||
const [titleE] = await llm.embedding(node.name.split('/').pop() || node.name);
|
||||
const [descE] = await llm.embedding(node.description || '');
|
||||
const bodyChunks = body ? await llm.embedding(body) : [];
|
||||
if(titleE) node.titleEmbedding = titleE.embedding;
|
||||
if(descE) node.embedding = descE.embedding;
|
||||
node.bodyEmbeddings = bodyChunks.map((c: any) => c.embedding).filter(Boolean);
|
||||
}
|
||||
|
||||
export function euclideanDistance(a: number[], b: number[]): number {
|
||||
let sum = 0;
|
||||
for(let i = 0; i < a.length; i++) {
|
||||
const d = a[i] - b[i];
|
||||
sum += d * d;
|
||||
}
|
||||
return Math.sqrt(sum);
|
||||
}
|
||||
|
||||
export function getWeekStart(date: Date = new Date()): string {
|
||||
const d = new Date(Date.UTC(date.getFullYear(), date.getMonth(), date.getDate()));
|
||||
const day = d.getUTCDay();
|
||||
const diff = day === 0 ? -6 : 1 - day;
|
||||
d.setUTCDate(d.getUTCDate() + diff);
|
||||
return d.toISOString().slice(0, 10);
|
||||
}
|
||||
|
||||
export function journalDescription(journalName?: string): string {
|
||||
const start = journalName?.split('/').pop() || getWeekStart();
|
||||
const d = new Date(`${start}T00:00:00Z`);
|
||||
d.setUTCDate(d.getUTCDate() + 6);
|
||||
const end = d.toISOString().slice(0, 10);
|
||||
return `Log from ${start} - ${end}`;
|
||||
}
|
||||
|
||||
function parseFrontmatter(content: string): {fm: Map<string, string>, body: string} {
|
||||
const match = content.match(/^---\n([\s\S]*?)\n---\n?([\s\S]*)$/);
|
||||
if(!match) return {fm: new Map(), body: content};
|
||||
const fm = new Map<string, string>();
|
||||
for(const line of match[1].split('\n')) {
|
||||
const i = line.indexOf(':');
|
||||
if(i === -1) continue;
|
||||
const key = line.slice(0, i).trim();
|
||||
const raw = line.slice(i + 1).trim();
|
||||
let value = raw;
|
||||
try { value = JSON.parse(raw); } catch { }
|
||||
fm.set(key, value);
|
||||
}
|
||||
return {fm, body: match[2]};
|
||||
}
|
||||
|
||||
export function writeFrontmatter(fm: Map<string, string>, body: string): string {
|
||||
const lines = [...fm.entries()].map(([k, v]) =>
|
||||
`${k}: ${JSON.stringify(String(v).replace(/\s+/g, ' ').trim())}`);
|
||||
return `---\n${lines.join('\n')}\n---\n\n${body.trimStart()}`;
|
||||
}
|
||||
|
||||
export function stripHeader(content: string): string {
|
||||
return content.replace(/^---[\s\S]*?\n---\n?/, '').trimStart();
|
||||
}
|
||||
|
||||
export function updateMemory(node: Memory, body: string): string {
|
||||
const {fm} = parseFrontmatter(node.content);
|
||||
fm.set('name', node.name);
|
||||
fm.set('description', (node.name.startsWith('Journal/') ? journalDescription(node.name) : node.description)
|
||||
|| 'Persistent memory document');
|
||||
fm.set('modified', new Date().toISOString());
|
||||
return writeFrontmatter(fm, stripHeader(body));
|
||||
}
|
||||
Reference in New Issue
Block a user