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