Compare commits

..
11 Commits
Author SHA1 Message Date
ztimson 1f1a4662d4 Entity based notes
Publish Library / Build NPM Project (push) Successful in 41s
Publish Library / Tag Version (push) Successful in 14s
2026-09-18 22:37:05 -04:00
ztimson ee4147e24e Fixed opanai early termination from tool calls
Publish Library / Build NPM Project (push) Successful in 46s
Publish Library / Tag Version (push) Successful in 15s
2026-09-18 16:07:58 -04:00
ztimson d29c0ca389 Fixed opanai early termination from tool calls
Publish Library / Build NPM Project (push) Successful in 55s
Publish Library / Tag Version (push) Successful in 9s
2026-09-18 02:15:02 -04:00
ztimson 4203cb34ef Better fact organization
Publish Library / Build NPM Project (push) Successful in 40s
Publish Library / Tag Version (push) Successful in 10s
2026-09-14 12:22:49 -04:00
ztimson d42c240362 Memorization optimziations
Publish Library / Build NPM Project (push) Successful in 1m2s
Publish Library / Tag Version (push) Successful in 10s
2026-08-31 12:38:40 -04:00
ztimson c1a16096ae Keep message progress on abort
Publish Library / Build NPM Project (push) Successful in 43s
Publish Library / Tag Version (push) Successful in 14s
2026-08-29 21:18:28 -04:00
ztimson ff0ee0b60e Patched memory merging
Publish Library / Build NPM Project (push) Successful in 45s
Publish Library / Tag Version (push) Successful in 10s
2026-08-28 16:48:46 -04:00
ztimson 0a6f1e4d62 Refined memory management prompts
Publish Library / Build NPM Project (push) Successful in 1m18s
Publish Library / Tag Version (push) Successful in 20s
2026-08-25 10:03:36 -04:00
ztimson 08a351e028 Better memory management
Publish Library / Build NPM Project (push) Successful in 59s
Publish Library / Tag Version (push) Successful in 11s
2026-08-24 14:42:10 -04:00
ztimson 85c01d3ef1 Added official file support
Publish Library / Build NPM Project (push) Successful in 30s
Publish Library / Tag Version (push) Successful in 10s
2026-08-17 15:50:48 -04:00
ztimson 5826573d5c Added official file support
Publish Library / Build NPM Project (push) Successful in 50s
Publish Library / Tag Version (push) Successful in 13s
2026-08-17 15:16:32 -04:00
8 changed files with 489 additions and 166 deletions
+6 -6
View File
@@ -1,19 +1,19 @@
{
"name": "@ztimson/ai-utils",
"version": "1.5.0",
"version": "1.6.6",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "@ztimson/ai-utils",
"version": "1.5.0",
"version": "1.6.6",
"license": "MIT",
"dependencies": {
"@anthropic-ai/sdk": "^0.102.0",
"@huggingface/transformers": "^4.2.0",
"@tensorflow/tfjs": "^4.22.0",
"@ztimson/node-utils": "^1.0.7",
"@ztimson/utils": "^0.29.4",
"@ztimson/utils": "^0.30.8",
"cheerio": "^1.2.0",
"openai": "^6.42.0",
"pdf-parse": "^2.4.5",
@@ -1525,9 +1525,9 @@
"license": "MIT"
},
"node_modules/@ztimson/utils": {
"version": "0.29.7",
"resolved": "https://registry.npmjs.org/@ztimson/utils/-/utils-0.29.7.tgz",
"integrity": "sha512-cjQ9+RjC5X7gKNA/hJHDf7OtyYCa+5E0PDc76lIaATwNAxXCSx2IO9r2wHiHtZGV5bldjnFmw7aOV8Jmq7SgKQ==",
"version": "0.30.8",
"resolved": "https://registry.npmjs.org/@ztimson/utils/-/utils-0.30.8.tgz",
"integrity": "sha512-+vBjcinqckqMHkP95xWiQeQz2E7Q1oS0b+Odjp+F9rvQ4z0US4JdodjqhEaqh+RAO/yP77x4Eu0a04yBB4HNdw==",
"license": "MIT",
"dependencies": {
"var-persist": "^1.0.1"
+2 -2
View File
@@ -1,6 +1,6 @@
{
"name": "@ztimson/ai-utils",
"version": "1.6.0",
"version": "1.6.11",
"description": "AI Utility library",
"author": "Zak Timson",
"license": "MIT",
@@ -29,7 +29,7 @@
"@huggingface/transformers": "^4.2.0",
"@tensorflow/tfjs": "^4.22.0",
"@ztimson/node-utils": "^1.0.7",
"@ztimson/utils": "^0.29.4",
"@ztimson/utils": "^0.30.8",
"cheerio": "^1.2.0",
"openai": "^6.42.0",
"pdf-parse": "^2.4.5",
+1 -1
View File
@@ -4,7 +4,7 @@ import { Audio } from './audio.ts';
import {Vision} from './vision.ts';
export type AbortablePromise<T> = Promise<T> & {
abort: () => any
abort: (keep?: boolean) => any
};
export type AiOptions = {
+50
View File
@@ -13,6 +13,56 @@ 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]));
const ensureNode = (name: string): MemoryNode => {
let n = byName.get(name);
if (!n) {
n = {name, missing: !nameSet.has(name), links: [], backlinks: []};
byName.set(name, n);
}
return n;
};
for (const m of changed) {
const node = ensureNode(m.name);
node.missing = false; // real memory, promotes any pre-existing ghost entry
const oldLinks = m.links ?? [];
const newLinks = extractLinks(m.content).filter(l => l !== m.name);
for (const target of oldLinks.filter(l => !newLinks.includes(l))) {
const t = byName.get(target);
if (!t) continue;
t.backlinks = t.backlinks.filter(n => n !== m.name);
if (t.missing && !t.backlinks.length) byName.delete(target); // fully dereferenced ghost
}
for (const target of newLinks.filter(l => !oldLinks.includes(l))) {
const t = ensureNode(target);
if (!t.backlinks.includes(m.name)) t.backlinks.push(m.name);
}
m.links = newLinks;
node.links = newLinks;
}
for (const m of mems) {
const n = byName.get(m.name);
if (n) m.backlinks = n.backlinks;
}
return [...byName.values()];
}
export function rebuildGraph(memories: Memory[] | MemoryCache): MemoryNode[] {
const mems = memories instanceof MemoryCache ? memories.memories : memories;
const nameSet = new Set(mems.map(m => m.name));
+53 -12
View File
@@ -15,6 +15,7 @@ interface KDNode<T> {
axis: number;
left: KDNode<T> | null;
right: KDNode<T> | null;
deleted?: boolean;
}
// ─── Distance helpers ─────────────────────────────────────────────────────────
@@ -95,6 +96,7 @@ class BoundedMaxHeap<T> {
*
* Supports:
* - Insertion of labeled points
* - Lazy (tombstone) removal, physically purged on rebalance()
* - k-nearest-neighbor (KNN) search
* - Radius search (all points within a given distance)
* - Euclidean and cosine distance metrics
@@ -103,6 +105,7 @@ class BoundedMaxHeap<T> {
export class KDTree<T = unknown> {
private root: KDNode<T> | null = null;
private _size = 0;
private _tombstones = 0;
private readonly distanceFn: (a: number[], b: number[]) => number;
readonly dims: number;
@@ -129,9 +132,15 @@ export class KDTree<T = unknown> {
}
}
/** Total number of points stored in the tree. */
/** Total number of live points stored in the tree (excludes tombstoned). */
get size(): number { return this._size; }
/** Fraction of physical nodes that are tombstoned (pending removal on next rebalance). */
get tombstoneRatio(): number {
const total = this._size + this._tombstones;
return total ? this._tombstones / total : 0;
}
// ── Insertion ──────────────────────────────────────────────────────────────
/**
@@ -144,10 +153,36 @@ export class KDTree<T = unknown> {
this._size++;
}
// ── Removal ────────────────────────────────────────────────────────────────
/**
* Lazily remove all live points whose payload matches `predicate`.
* O(n) traversal, but avoids a full tree rebuild. Call `rebalance()`
* periodically (e.g. once tombstoneRatio crosses ~0.25) to reclaim space
* and restore optimal query depth.
* @returns number of points removed
*/
remove(predicate: (payload: T) => boolean): number {
let removed = 0;
const visit = (node: KDNode<T> | null): void => {
if (!node) return;
if (!node.deleted && predicate(node.point.payload)) {
node.deleted = true;
removed++;
}
visit(node.left);
visit(node.right);
};
visit(this.root);
this._size -= removed;
this._tombstones += removed;
return removed;
}
// ── KNN search ─────────────────────────────────────────────────────────────
/**
* Find the k nearest neighbors to `query`.
* Find the k nearest live neighbors to `query`.
* Returns results sorted by distance ascending.
*/
knn(query: number[], k: number): KNNResult<T>[] {
@@ -171,7 +206,7 @@ export class KDTree<T = unknown> {
// ── Radius search ──────────────────────────────────────────────────────────
/**
* Return all points whose distance to `query` is ≤ `radius`,
* Return all live points whose distance to `query` is ≤ `radius`,
* sorted by distance ascending.
*/
radiusSearch(query: number[], radius: number): KNNResult<T>[] {
@@ -186,7 +221,7 @@ export class KDTree<T = unknown> {
// ── Conversion ─────────────────────────────────────────────────────────────
/** Collect all points in the tree (order not guaranteed). */
/** Collect all live points in the tree (order not guaranteed). */
toArray(): KDPoint<T>[] {
const out: KDPoint<T>[] = [];
this.collect(this.root, out);
@@ -194,12 +229,14 @@ export class KDTree<T = unknown> {
}
/**
* Rebuild the tree from its current points as a balanced tree.
* Useful after many individual insertions to restore O(log n) query time.
* Rebuild the tree from its current live points as a balanced tree.
* Physically purges tombstones and restores O(log n) query time.
*/
rebalance(): void {
const points = this.toArray();
this.root = points.length ? this.buildBalanced(points, 0) : null;
this._size = points.length;
this._tombstones = 0;
}
// ── Private: build ─────────────────────────────────────────────────────────
@@ -251,8 +288,10 @@ export class KDTree<T = unknown> {
): void {
if (node === null) return;
const dist = this.distanceFn(query, node.point.vector);
heap.push({ point: node.point, distance: dist });
if (!node.deleted) {
const dist = this.distanceFn(query, node.point.vector);
heap.push({ point: node.point, distance: dist });
}
const axis = node.axis;
const diff = query[axis] - node.point.vector[axis];
@@ -285,9 +324,11 @@ export class KDTree<T = unknown> {
): void {
if (node === null) return;
const dist = this.distanceFn(query, node.point.vector);
if (dist <= radius) {
results.push({ point: node.point, distance: dist });
if (!node.deleted) {
const dist = this.distanceFn(query, node.point.vector);
if (dist <= radius) {
results.push({ point: node.point, distance: dist });
}
}
const axis = node.axis;
@@ -310,7 +351,7 @@ export class KDTree<T = unknown> {
private collect(node: KDNode<T> | null, out: KDPoint<T>[]): void {
if (node === null) return;
out.push(node.point);
if (!node.deleted) out.push(node.point);
this.collect(node.left, out);
this.collect(node.right, out);
}
+51 -23
View File
@@ -50,6 +50,8 @@ export type LLMMessage = {
role: 'assistant' | 'system' | 'user';
/** Message content */
content: string | any;
/** Files attached to request */
files?: LLMFile[];
/** Timestamp */
timestamp?: number;
/** Response duration in ms */
@@ -241,10 +243,8 @@ class LLM {
} else if(isText) {
text = (await this.loadBuffer(file, true)).toString('utf-8');
} else {
text = `Unsupported file type: ${ext || mime}`;
text = typeof file.content === 'string' ? file.content : `[Binary file, unable to extract: ${name}]`;
}
// Cache result, skip re-extraction on future turns of the same conversation
file.content = text;
file.extracted = true;
delete file.path;
@@ -265,7 +265,7 @@ class LLM {
};
}
private setupAgent(agents: Agent[] = [], allAgents: Agent[], history: LLMMessage[], aborts: (() => void)[], depth = 0, delegateState: {resp: string | null}): AiTool[] {
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)}`;
return {
@@ -397,11 +397,13 @@ ${a.system}`,
if(!this.models[m]) throw new Error(`Model does not exist: ${m}`);
let request: AbortablePromise<string> | null = null;
let aborted = false;
const nestedAborts: (() => void)[] = [];
const abort = () => {
let keepOnAbort = true;
const nestedAborts: ((keep?: boolean) => void)[] = [];
const abort = (keep = true) => {
aborted = true;
request?.abort?.();
nestedAborts.forEach(a => a());
keepOnAbort = keep;
request?.abort?.(keep);
nestedAborts.forEach(a => a(keep));
};
let promise: any;
@@ -411,7 +413,24 @@ ${a.system}`,
let tools: AiTool[] = options.tools || this.ai.options.llm?.tools || [];
const prompts: string[] = [];
let history = options.history || [];
if(message) history.push({role: 'user', content: message, timestamp: Date.now()});
const historyStart = history.length;
const files = options.files || [];
if(message || files.length) history.push({role: 'user', content: message || '', timestamp: Date.now()});
// Accumulate streamed text so it can be committed to history if aborted mid-generation
let partialText = '';
const onStream = options.stream;
const stream = (chunk: {text?: string, tool?: string, done?: true}) => {
if(chunk.text) partialText += chunk.text;
return onStream?.(chunk);
};
/** Commit (keep) or discard this turn's progress on abort, then throw */
const abortNow = (): never => {
if(keepOnAbort) { if(partialText) history.push({role: 'assistant', content: partialText, timestamp: Date.now()}); }
else history.splice(historyStart, history.length - historyStart);
throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
};
// MCP
const mcp = options.mcp || this.ai.options?.llm?.mcp;
@@ -440,8 +459,8 @@ ${a.system}`,
const mems = mem.memory instanceof MemoryCache ? mem.memory.memories : mem.memory;
if(mems.length) {
if(mem.inject) {
const pool = 15; // candidates considered, cheap since only refs are listed
const budget = mem.maxTokens ?? 2000; // actual content injected
const pool = 15;
const budget = mem.maxTokens ?? 2000;
const relevant = await this.memoryManager.recollect(message, mem.memory, pool);
let used = 0;
@@ -480,16 +499,18 @@ Linked: ${makeUnique([...r.links, ...r.backlinks]).join(', ')}
}
}
if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
if(aborted) abortNow();
// Files
const files = options.files || [];
const lastMsg = history[history.length - 1];
const originalContent = lastMsg?.content;
if(files.length && lastMsg?.role === 'user') {
const {text, images} = await this.resolveFiles(files);
const merged = text ? `${originalContent}\n\n${text}` : originalContent;
lastMsg.content = images.length
if(files.length && lastMsg?.role === 'user') lastMsg.files = files;
const restores: {msg: LLMMessage, content: any}[] = [];
for(const msg of history) {
if(msg.role !== 'user' || !msg.files?.length) continue;
const {text, images} = await this.resolveFiles(msg.files);
if(!text && !images.length) continue;
restores.push({msg, content: msg.content});
const merged = text ? [msg.content, text].filter(Boolean).join('\n\n') : msg.content;
msg.content = images.length
? [...images.map(i => ({type: 'image', mime: i.mime, data: i.data})), {type: 'text', text: merged}]
: merged;
}
@@ -497,13 +518,20 @@ Linked: ${makeUnique([...r.links, ...r.backlinks]).join(', ')}
const toolTimings = new Map<string, {duration: number, tps: number}>();
tools = this.wrapToolTiming(tools, toolTimings);
if(aborted) throw Object.assign(new Error('Aborted'), {name: 'AbortError'});
if(aborted) abortNow();
prompts.unshift(options.system || this.ai.options.llm?.system || '');
request = this.models[m].ask('', {...options, tools, system: prompts.filter(Boolean).join('\n\n')});
let resp = await request;
request = this.models[m].ask('', {...options, tools, stream, system: prompts.filter(Boolean).join('\n\n')});
let resp: string;
try {
resp = await request;
} catch(err: any) {
if(aborted) return abortNow();
throw err;
}
if(files.length && lastMsg?.role === 'user') lastMsg.content = originalContent;
// Strip the file injection shim
restores.forEach(({msg, content}) => msg.content = content);
// Capture meta (duration / tps)
for(const h of history) {
+310 -119
View File
@@ -1,23 +1,24 @@
import {MemoryNode, rebuildGraph} from './helpers.ts';
import {MemoryNode, patchGraph, rebuildGraph} from './helpers.ts';
import {LLMRequest, LLMMessage} from './llm.ts';
import {AiTool} from './tools.ts';
import {KDPoint, KDTree} from './kd-tree.ts';
import {KDTree} from './kd-tree.ts';
import {escapeRegex} from '@ztimson/utils';
const FACTS_HEADING = '## Facts';
const GENERIC_TEMPLATE = `# {{Title}}
## Summary
## Details
## Related`;
const MERGE_THRESHOLD = 0.12;
const PENDING_HEADING = '## Pending';
const TREE_TOMBSTONE_LIMIT = 0.25;
const ALIAS_MATCH_THRESHOLD = 0.55;
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[];
}
@@ -25,6 +26,8 @@ export type Memory = {
type MemoryRef = {
name: string;
description: string;
/** Cosine distance from the query, present when returned from a search */
distance?: number;
}
type FactBucket = {
@@ -32,6 +35,11 @@ type FactBucket = {
facts: string[];
}
type FactAgentResult = {
buckets: FactBucket[];
journal: string;
}
function dedupeFacts(facts: string[]): string[] {
const seen = new Map<string, string>();
for (const f of facts) {
@@ -55,10 +63,20 @@ function cosineDistance(a: number[], b: number[]): number {
function cosineSearch(query: number[], memories: Memory[], limit: number): MemoryRef[] {
return memories
.filter(m => m.embedding?.length)
.map(m => ({ref: {name: m.name, description: m.description}, distance: cosineDistance(query, m.embedding)}))
.map(m => ({name: m.name, description: m.description, distance: cosineDistance(query, m.embedding)}))
.sort((a, b) => a.distance - b.distance)
.slice(0, limit)
.map(s => s.ref);
.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 {
@@ -67,6 +85,8 @@ export function stripHeader(content: string): string {
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[] = [];
@@ -74,37 +94,48 @@ export class MemoryCache {
constructor(memories: Memory[]) {
this.memories = memories;
this.tree = new KDTree<MemoryRef>(0);
this.rebuild();
}
private buildTree(): KDTree<MemoryRef> {
const embedded = this.memories.filter(m => m.embedding?.length);
if (!embedded.length) return new KDTree<MemoryRef>(0);
/** 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));
const dims = embedded[0].embedding.length;
const points: KDPoint<MemoryRef>[] = embedded.map(m => ({
vector: m.embedding,
payload: {name: m.name, description: m.description},
}));
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);
}
}
return new KDTree<MemoryRef>(dims, 'cosine', points);
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);
return this.tree.knn(query, limit).map(r => ({...r.point.payload, distance: r.distance}));
}
add(memory: Memory): void {
this.memories.push(memory);
this.rebuild();
this.rebuild([memory]);
}
update(memory: Memory): void {
const idx = this.memories.findIndex(m => m.name === memory.name);
if (idx !== -1) this.memories[idx] = memory;
const existing = this.memories.find(m => m.name === memory.name);
if (existing) Object.assign(existing, memory);
else this.memories.push(memory);
this.rebuild();
this.rebuild([existing ?? memory]);
}
remove(name: string): void {
@@ -115,9 +146,11 @@ export class MemoryCache {
}
}
rebuild(): void {
this.nodes = rebuildGraph(this.memories);
this.tree = this.buildTree();
rebuild(changed?: Memory[]): void {
this.nodes = (changed?.length && this.nodes.length)
? patchGraph(this.memories, this.nodes, changed)
: rebuildGraph(this.memories);
this.syncTree();
}
}
@@ -134,9 +167,9 @@ class MemoryAccessor {
return this.list.find(m => m.name === name);
}
commit(): MemoryNode[] {
commit(changed?: Memory[]): MemoryNode[] {
if (this.cache) {
this.cache.rebuild();
this.cache.rebuild(changed);
return this.cache.nodes;
}
return rebuildGraph(this.list);
@@ -147,6 +180,7 @@ class MemoryAccessor {
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);
}
@@ -162,10 +196,7 @@ class MemoryAccessor {
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(async node => {
const [e] = await llm.embedding(node.content);
if (e) node.embedding = e.embedding;
}));
await Promise.all(missing.map(node => embedMemoryFields(node, llm)));
this.commit();
return missing.length;
}
@@ -185,13 +216,13 @@ export type MemoryOptions = {
}
export class MemoryManager {
private recentlyTouched = new Map<string, number>();
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 => ({
@@ -251,14 +282,13 @@ ${m.content}
return new MemoryAccessor(memories);
}
private appendFacts(node: Memory, facts: string[]): void {
private stage(node: Memory, block: string): void {
this.ensureDoc(node);
const body = stripHeader(node.content);
const bullets = facts.map(f => `- ${f}`).join('\n');
const idx = body.indexOf(FACTS_HEADING);
const idx = body.indexOf(PENDING_HEADING);
const newBody = idx === -1
? `${body.trimEnd()}\n\n${FACTS_HEADING}\n${bullets}\n`
: `${body.slice(0, idx + FACTS_HEADING.length)}\n${bullets}${body.slice(idx + FACTS_HEADING.length)}`;
? `${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);
}
@@ -268,42 +298,92 @@ ${m.content}
node.content = this.touchHeader(node, `# ${title}\n`);
}
private async factAgent(conversation: string, store: MemoryAccessor, options: LLMRequest, weekKey: string): Promise<FactBucket[]> {
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;
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 to build obsidian knowledge vaults.
Analyze this conversation and extract facts worth remembering long-term.
system: `Extract durable memory from this conversation
Rules:
- Always extract facts that the user explicitly told you to remember
- ONLY extract current facts the USER explicitly stated about themselves, their work, projects or decisions that were MADE during this conversation
- DO NOT extract greetings, pleasantries, or generic exchanges
- DO NOT extract deltas or changes in facts; ONLY the end fact
- DO NOT extract anything the AI/assistant itself said
- If nothing worth remembering was said, return an empty buckets array
1. Journal recap
- Brief "Captain's Log" of what happened, including useful context, decisions, or events
- Leave empty for trivial exchanges
When extracting facts, you MUST also decide the exact destination path:
- Reuse node names (including ghost) as much as possible IF the facts belongs there
- All information primarily about the user should go under "People/User"
- When required, create a new path following collection/subject format (e.g., People/Sarah, Projects/Oxide) — you are not limited to any fixed list of collections, use whatever fits
- For journal entries, use "Journal"
2. Fact buckets
- Extract only durable facts explicitly stated by the USER
- Record the final/end state, not intermediate changes
- Do not extract assistant claims, guesses, greetings, or temporary conversation details
For each fact, identify its HOME ENTITY:
- The HOME ENTITY name should always be a [abstract|pro]noun
- The grammatical subject/owner of the fact is the strongest clue
- Prefer an existing entity over creating a new one
- A document represents a persistent entity, not a topic, feature, bug, event, decision, setting, or conversation fragment
- Put project facts under the project they belong to, person facts under the person, etc
- New child entities are appropriate only when they are themselves distinct persistent entities
Example Paths:
- Projects/[Name]
- People/[Name]
- History/[Name]
- Science/[Name]
- [Subject]/[Name]
- Class/[Name]/[Child]
Use [[WikiLinks]] to express relationships between entities. NEVER create documents just to hold relationships
Keep journal material in the journal; don't turn journal events into entities unless they represent something persistent
Available nodes:
- Journal
${this.listNodes(store.list).filter(n => !n.name.includes('Journal')).map(n => `- ${n.name}: ${n.description}`).join('\n') || 'None yet.'}
${this.listNodes(store.list).map(n => `- ${n.name}: ${n.description}`).join('\n') || 'None yet.'}
${ghosts.length ? `${ghosts.map(g => `- ${g}: (Ghost)`).join('\n')}` : ''}`,
schema: {
buckets: {type: 'array', description: 'Groups of facts to remember, each assigned to a different node. Return an empty array if there is nothing worth storing in an obsidian vault', items: {
type: 'object', items: {
subject: {type: 'string', description: 'Exact existing node name OR new path (e.g. "People/Sarah", "Projects/Oxide"), or "Journal"', required: true},
facts: {
type: 'array',
description: 'Facts to store at this destination',
items: {type: 'string', description: 'A single fact'},
},
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 persistent entity path', required: true},
facts: {type: 'array', description: 'Facts to store here', items: {type: 'string'}},
},
},
},
@@ -311,15 +391,17 @@ ${ghosts.length ? `${ghosts.map(g => `- ${g}: (Ghost)`).join('\n')}` : ''}`,
});
const buckets = new Map<string, string[]>();
for(const bucket of response.buckets ?? []) {
const subject = bucket.subject.trim().toLowerCase() === 'journal'
? `Journal/${weekKey}` : bucket.subject.trim();
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.entries().toArray().map(([subject, facts]) => ({subject, facts}));
return {
buckets: buckets.entries().toArray().map(([subject, facts]) => ({subject, facts})),
journal: (response.journal ?? '').trim(),
};
}
private getWeekMonday(date: Date = new Date()): string {
@@ -334,6 +416,36 @@ ${ghosts.length ? `${ghosts.map(g => `- ${g}: (Ghost)`).join('\n')}` : ''}`,
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);
@@ -347,18 +459,25 @@ ${ghosts.length ? `${ghosts.map(g => `- ${g}: (Ghost)`).join('\n')}` : ''}`,
this.queues.set(key, entry);
const store = this.access(memories);
entry.task = (async () => {
do {
entry.dirty = false;
await this.docAgent(node, store.list, options, entry);
} while (entry.dirty);
})().finally(() => {
this.queues.delete(key);
store.commit();
});
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 {
@@ -367,27 +486,25 @@ ${ghosts.length ? `${ghosts.map(g => `- ${g}: (Ghost)`).join('\n')}` : ''}`,
model: options.model,
temperature: 0.3,
schema: {
description: {type: 'string', description: 'One-line description of what this document covers, no formatting or emojis', required: true},
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 document in an Obsidian-style vault.
system: `You maintain one persistent knowledge-base document
If the document has a "${FACTS_HEADING}" section, integrate every bullet under it into the appropriate part of the document, then remove the "${FACTS_HEADING}" section entirely. If there is no such section, just tidy the document per the rules below.
Rewrite the ENTIRE document, folding "## Pending" into the existing content. Remove the Pending section when finished
Structure: follow this generic shape loosely, adapting section names/order to what the content actually needs (e.g. journal-style docs may want a timeline instead of "Details"):
\`\`\`markdown
${GENERIC_TEMPLATE}
\`\`\`
Document design:
- The document represents one entity. Keep information about that entity together
- Let the structure fit the entity; there is NO fixed template
- Preserve useful existing headings and organization. Don't redesign the document without reason
- Add headings only when they meaningfully organize recurring information; don't create headings for one-off facts
- Keep the document concise and information-dense without removing useful technical specifics
- Current truth wins when facts conflict. Preserve older context only when it adds useful meaning
- Use [[WikiLinks]] for specific related entities; don't create redundant content for linked entities
- Avoid generic filler sections such as Notes, Miscellaneous, Recent, Updates, or Conversation
- No frontmatter, preamble, filler, or AI commentary
Formatting rules:
- Use Obsidian-style markdown: # headings, **bold** for emphasis, bullet & numbered lists for grouped 1D data, tables for 2D data
- Link related concepts with [[WikiLink]] notation using full paths like [[People/Sarah]] or [[Projects/Website]]
- Create links for specific entities (person, place, project, program) and abstract concepts, but skip generics (car, red, dog)
- Keep the document concise, factual, and human-readable
- Resolve contradictions: newer facts always win — delete the outdated statement entirely, never keep both
- Do not add frontmatter blocks, filler, preamble, or AI commentary
Other nodes in the vault (link to these instead of duplicating their content):
Available nodes to link to:
${this.listNodes(memories).filter(n => n.name !== node.name).map(n => n.name).join(', ') || 'none'}
Current document:
@@ -406,10 +523,45 @@ ${currentBody}
}
if (!update?.content) return;
node.description = node.name !== 'People/User' ? update.description : 'All information about the current user';
node.description = node.name !== 'People/User' ? this.sanitizeDescription(update.description) : 'All information about the current user';
node.content = this.touchHeader(node, update.content);
const [e] = await this.llm.embedding(node.content);
if (e) node.embedding = e.embedding;
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: 'Canonical path for the merged entity', 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: `Determine whether these two documents represent the SAME persistent entity.
Similarity of subject matter is NOT enough. Do not merge documents merely because they discuss the same project, person, technology, topic, or related work.
Merge only when the evidence indicates they are duplicate identities, aliases, renamed entities, or two documents accidentally created for the same real-world entity. If they are distinct entities, they must remain separate.
If they are the same entity:
- Choose the canonical/most established path.
- Combine their information into one document and remove duplication.
- Preserve useful structure, technical specifics, history, and [[WikiLinks]].
- Prefer newer information when facts conflict.
- Return the canonical entity name and the fully reconciled document.
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} {
@@ -419,21 +571,30 @@ ${currentBody}
for (const line of match[1].split('\n')) {
const i = line.indexOf(':');
if (i === -1) continue;
fm.set(line.slice(0, i).trim(), line.slice(i + 1).trim());
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, body);
return this.writeFrontmatter(fm, stripHeader(body));
}
private writeFrontmatter(fm: Map<string, string>, body: string): string {
const lines = [...fm.entries()].map(([k, v]) => `${k}: ${v}`);
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()}`;
}
@@ -452,6 +613,19 @@ ${currentBody}
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 [];
@@ -461,8 +635,11 @@ ${currentBody}
const [e] = await this.llm.embedding(query);
if (!e) return [];
const vectorResults = store.search(e.embedding, limit);
const found = new Set<string>(vectorResults.map(r => r.name));
// 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];
@@ -482,9 +659,9 @@ ${currentBody}
}
}
const vectorOrder = vectorResults.map(r => r.name);
const graphExpansions = [...found].filter(n => !vectorOrder.includes(n));
return [...vectorOrder, ...graphExpansions].map(n => store.find(n)!).filter(Boolean);
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[]> {
@@ -498,26 +675,40 @@ ${currentBody}
history.push(pending);
const store = this.access(memories);
const buckets = await this.factAgent(conversation, store, options, this.getWeekMonday());
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) {
let node = store.find(subject);
const resolved = this.resolveSubject(subject, store);
let node = store.find(resolved);
if (!node) {
node = {name: subject, description: '', content: '', embedding: [], links: [], backlinks: []};
node = {name: resolved, description: '', content: '', embedding: [], links: [], backlinks: []};
store.list.push(node);
}
this.appendFacts(node, facts);
const [e] = await this.llm.embedding(node.content);
if (e) node.embedding = e.embedding;
this.touch(node.name);
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();
store.commit(touched);
(pending as any).content = `Saved to ${touched.map(n => `[[${n.name}]]`).join(', ')}`;
await Promise.all(touched.map(node => this.reconcile(node, memories, options).catch(() => {})));
Promise.all(touched.map(node => this.reconcile(node, memories, options).catch(() => {})));
} else {
(pending as any).content = 'Nothing worth remembering.';
}
@@ -526,9 +717,9 @@ ${currentBody}
return touched;
}
async reconcileVault(memories: Memory[] | MemoryCache, options: LLMRequest, scope: 'touched' | 'all' = 'touched'): Promise<void> {
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(FACTS_HEADING));
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();
}
+16 -3
View File
@@ -98,18 +98,21 @@ export class OpenAi extends LLMProvider {
throw err;
});
let usage: any, msg: any = {content: '', tool_calls: []};
let usage: any, finishReason: string | undefined, msg: any = {content: '', tool_calls: []};
if(options.stream) {
for await (const chunk of resp) {
if(controller.signal.aborted) break;
if(chunk.usage) usage = chunk.usage;
if(chunk.choices[0]?.finish_reason) finishReason = chunk.choices[0].finish_reason;
if(chunk.choices[0]?.delta?.content) {
msg.content += chunk.choices[0].delta.content;
options.stream({text: chunk.choices[0].delta.content});
}
if(chunk.choices[0]?.delta?.tool_calls) {
for(const deltaTC of chunk.choices[0].delta.tool_calls) {
const existing = msg.tool_calls.find((tc: any) => tc.index === deltaTC.index);
const existing = deltaTC.index != null
? msg.tool_calls.find((tc: any) => tc.index === deltaTC.index)
: (deltaTC.id ? msg.tool_calls.find((tc: any) => tc.id === deltaTC.id) : undefined);
if(existing) {
if(deltaTC.id) existing.id = deltaTC.id;
if(deltaTC.function?.name) existing.function.name = deltaTC.function.name;
@@ -126,11 +129,21 @@ export class OpenAi extends LLMProvider {
}
} else {
usage = resp.usage;
finishReason = resp.choices[0].finish_reason;
msg = resp.choices[0].message;
}
const duration = Date.now() - callStart;
const tps = usage?.completion_tokens && duration > 0 ? usage.completion_tokens / (duration / 1000) : 0;
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] Stream ended prematurely - connection likely dropped');
}
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});
@@ -147,7 +160,7 @@ export class OpenAi extends LLMProvider {
if(!tool) { entry.error = 'Tool not found'; return; }
try {
const toolStream = options.stream && ((chunk: any) => {
if(chunk.done) { terminal = true; return; }
if(chunk.done) return;
options.stream!(chunk);
});
const result = await tool.fn(entry.args, toolStream, this.ai, tc.id);