The pieces more than one part of the core needs at the same time, behind locks held only for the instant they are used. The conversation and the background memory work both write the parents' log and read the child's memory; neither may wait on the other for the length of a model call.
6use std::sync::{Arc, Mutex, MutexGuard, PoisonError};
8use ::log::{debug, error, info, trace, warn}; 9 10use crate::chat::{ChatState, ChatTurn, CompressPlan}; 11use crate::ports::{ChatError, ChatStore, Log, LogError, Memory, MemoryError}; 12use crate::types::{Entry, Fact, NewFact, Speaker}; 13 14fn lock<T: ?Sized>(m: &Mutex<T>) -> MutexGuard<'_, T> { 15 // A panic elsewhere must not stop the parents' log or the child's memory. 16 m.lock().unwrap_or_else(|e| { 17 warn!("a lock was poisoned by a panic elsewhere; carrying on with its contents"); 18 PoisonError::into_inner(e) 19 }) 20} 21 22#[derive(Clone)] 23pub struct SharedLog(Arc<Mutex<Box<dyn Log + Send>>>); 24 25impl SharedLog { 26 pub fn new(log: Box<dyn Log + Send>) -> Self { 27 Self(Arc::new(Mutex::new(log))) 28 } 29 30 pub fn append(&self, entry: &Entry) -> Result<(), LogError> { 31 trace!("shared log append at {} ms", entry.at_ms); 32 lock(&self.0).append(entry) 33 } 34} 35 36#[derive(Clone)] 37pub struct SharedMemory(Arc<Mutex<Box<dyn Memory>>>); 38 39impl SharedMemory { 40 pub fn new(memory: Box<dyn Memory>) -> Self { 41 Self(Arc::new(Mutex::new(memory))) 42 }
Every fact, including the ones she put away: the parents' view and the memory's own work
(duplicates, embeddings). Whatever is built into a prompt or a search uses usable.
The facts Whiskers may use: all but those she put away.
Gives a fact its picture if it has none yet.
She puts a fact away.
A parent restores a fact she put away.
75 pub fn add(&self, new: NewFact, at_ms: u64) -> Result<Fact, MemoryError> { 76 debug!("shared memory: add a fact of {} chars", new.text.len()); 77 lock(&self.0).add(new, at_ms) 78 } 79 80 pub fn set_embedding(&self, id: u64, embedding: Vec<f32>) -> Result<(), MemoryError> { 81 lock(&self.0).set_embedding(id, embedding) 82 } 83 84 pub fn forget(&self, id: u64) -> Result<Fact, MemoryError> { 85 debug!("shared memory: forget {id}"); 86 lock(&self.0).forget(id) 87 } 88 89 pub fn snapshot(&self) -> crate::types::MemorySnapshot { 90 lock(&self.0).snapshot() 91 } 92 93 pub fn merge(&self, remote: crate::types::MemorySnapshot) -> Result<usize, MemoryError> { 94 lock(&self.0).merge(remote) 95 } 96} 97 98struct ChatInner { 99 state: ChatState, 100 store: Box<dyn ChatStore>, 101} 102 103#[derive(Clone)] 104pub struct SharedChat(Arc<Mutex<ChatInner>>); 105 106impl SharedChat {
A chat file that cannot be read is an error: starting empty would overwrite it.
108 pub fn open(store: Box<dyn ChatStore>) -> Result<Self, ChatError> { 109 let state = store.load().inspect_err(|e| error!("the chat cannot be opened: {}", e.0))?; 110 debug!("shared chat opened: {} turns, version {}", state.turns.len(), state.version); 111 Ok(Self(Arc::new(Mutex::new(ChatInner { state, store })))) 112 }
Takes in another device's copy of the one session, if it is the newer. The session is one thing wherever she talks to Whiskers; what the older copy lacks is still in the log.
120 pub fn adopt(&self, remote: ChatState) -> Result<bool, ChatError> { 121 let mut inner = lock(&self.0); 122 if !remote.is_newer_than(&inner.state) { 123 debug!("chat adopt: the remote copy (v{}) is not newer than ours (v{})", remote.version, inner.state.version); 124 return Ok(false); 125 } 126 info!("chat adopt: taking the remote copy v{} over our v{}", remote.version, inner.state.version); 127 inner.store.save(&remote).inspect_err(|e| error!("chat adopt: saving the remote copy failed: {}", e.0))?; 128 inner.state = remote; 129 Ok(true) 130 }
132 fn mutate(&self, f: impl FnOnce(&mut ChatState)) -> Result<(), ChatError> { 133 let mut inner = lock(&self.0); 134 let before = inner.state.clone(); 135 f(&mut inner.state); 136 inner.state.version += 1; 137 let state = inner.state.clone(); 138 if let Err(e) = inner.store.save(&state) { 139 // What is on disk and what is in memory must not disagree. 140 error!("chat change rolled back, the save failed: {}", e.0); 141 inner.state = before; 142 return Err(e); 143 } 144 Ok(()) 145 }
Records a finished exchange and that something just happened.
148 pub fn push_exchange(&self, heard: &str, said: &str, now_ms: u64) -> Result<(), ChatError> { 149 debug!("chat: push an exchange ({} + {} chars)", heard.len(), said.len()); 150 self.mutate(|s| { 151 s.turns.push(ChatTurn { speaker: Speaker::Child, text: heard.to_owned() }); 152 s.turns.push(ChatTurn { speaker: Speaker::Whiskers, text: said.to_owned() }); 153 s.last_active_ms = now_ms; 154 }) 155 }
If the chat is longer than after turns, the oldest turns to fold into the summary,
leaving keep behind.
164 pub fn compress_plan(&self, after: usize, keep: usize) -> Option<CompressPlan> { 165 let s = lock(&self.0).state.clone(); 166 if s.turns.len() <= after || s.turns.len() <= keep { 167 trace!("chat: {} turns, no compression due", s.turns.len()); 168 return None; 169 } 170 let n = s.turns.len() - keep; 171 debug!("chat: compression due, folding {n} of {} turns", s.turns.len()); 172 Some(CompressPlan { summary: s.summary, oldest: s.turns[..n].to_vec() }) 173 }
Replaces the oldest n turns with summary. Turns only ever arrive at the end, so
the first n are still the ones that were summarised.
177 pub fn apply_compression(&self, n: usize, summary: String) -> Result<(), ChatError> { 178 info!("chat: applying a summary of {} chars over the oldest {n} turns", summary.len()); 179 self.mutate(|s| { 180 let n = n.min(s.turns.len()); 181 s.turns.drain(..n); 182 s.summary = summary; 183 }) 184 } 185}