1//! The pieces more than one part of the core needs at the same time, behind locks held 2//! only for the instant they are used. The conversation and the background memory work 3//! both write the parents' log and read the child's memory; neither may wait on the 4//! other for the length of a model call. 5 6use std::sync::{Arc, Mutex, MutexGuard, PoisonError}; 7 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 } 43 44 /// Every fact, including the ones she put away: the parents' view and the memory's own work 45 /// (duplicates, embeddings). Whatever is built into a prompt or a search uses [`usable`](Self::usable). 46 pub fn facts(&self) -> Vec<Fact> { 47 lock(&self.0).facts() 48 } 49 50 /// The facts Whiskers may use: all but those she put away. 51 pub fn usable(&self) -> Vec<Fact> { 52 let mut facts = lock(&self.0).facts(); 53 facts.retain(Fact::whiskers_may_use); 54 facts 55 } 56 57 /// Gives a fact its picture if it has none yet. 58 pub fn give_icon(&self, id: u64, icon: whiskers_icons::IconId) -> Result<Fact, MemoryError> { 59 debug!("shared memory: give fact {id} an icon"); 60 lock(&self.0).give_icon(id, icon) 61 } 62 63 /// She puts a fact away. 64 pub fn hide(&self, id: u64, at_ms: u64) -> Result<Fact, MemoryError> { 65 debug!("shared memory: hide {id}"); 66 lock(&self.0).hide(id, at_ms) 67 } 68 69 /// A parent restores a fact she put away. 70 pub fn restore(&self, id: u64, at_ms: u64) -> Result<Fact, MemoryError> { 71 debug!("shared memory: restore {id}"); 72 lock(&self.0).restore(id, at_ms) 73 } 74 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 { 107 /// 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 } 113 114 pub fn snapshot(&self) -> ChatState { 115 lock(&self.0).state.clone() 116 } 117 118 /// Takes in another device's copy of the one session, if it is the newer. The session is one 119 /// 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 } 131 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 } 146 147 /// 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 } 156 157 pub fn touch(&self, now_ms: u64) -> Result<(), ChatError> { 158 trace!("chat: touched at {now_ms}"); 159 self.mutate(|s| s.last_active_ms = now_ms) 160 } 161 162 /// If the chat is longer than `after` turns, the oldest turns to fold into the summary, 163 /// 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 } 174 175 /// Replaces the oldest `n` turns with `summary`. Turns only ever arrive at the end, so 176 /// 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}