Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1f1a4662d4 | ||
|
|
ee4147e24e | ||
|
|
d29c0ca389 | ||
|
|
4203cb34ef | ||
|
|
d42c240362 | ||
|
|
c1a16096ae | ||
|
|
ff0ee0b60e | ||
|
|
0a6f1e4d62 | ||
|
|
08a351e028 | ||
|
|
85c01d3ef1 |
Generated
+6
-6
@@ -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
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@ztimson/ai-utils",
|
||||
"version": "1.6.1",
|
||||
"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",
|
||||
|
||||
@@ -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 = {
|
||||
|
||||
@@ -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));
|
||||
|
||||
+48
-7
@@ -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;
|
||||
|
||||
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,10 +324,12 @@ export class KDTree<T = unknown> {
|
||||
): void {
|
||||
if (node === null) return;
|
||||
|
||||
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;
|
||||
const diff = query[axis] - node.point.vector[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);
|
||||
}
|
||||
|
||||
+49
-24
@@ -243,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;
|
||||
@@ -267,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 {
|
||||
@@ -399,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;
|
||||
@@ -413,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;
|
||||
@@ -442,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;
|
||||
@@ -482,17 +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') {
|
||||
lastMsg.files = files;
|
||||
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;
|
||||
}
|
||||
@@ -500,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) {
|
||||
|
||||
+305
-114
@@ -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: {
|
||||
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 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'},
|
||||
},
|
||||
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 () => {
|
||||
let current = node, merged = false;
|
||||
try {
|
||||
do {
|
||||
entry.dirty = false;
|
||||
await this.docAgent(node, store.list, options, entry);
|
||||
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(() => {
|
||||
} finally {
|
||||
store.commit(merged ? undefined : [node]);
|
||||
this.queues.delete(key);
|
||||
store.commit();
|
||||
});
|
||||
}
|
||||
})();
|
||||
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
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user