1use std::sync::Arc; 2 3use ::log::{debug, error, info, trace, warn}; 4 5use crate::chat::CompressPlan; 6use crate::persona::{compress_prompt, remember_prompt}; 7use crate::profile::SharedProfile; 8use std::collections::HashSet; 9use std::sync::Mutex; 10 11use crate::ports::{Clock, Embedder, IconChooser, Guard, Model}; 12use crate::recall::cosine; 13use crate::shared::{SharedChat, SharedLog, SharedMemory}; 14use crate::types::{Direction, Entry, Event, Fact, Kind, NewFact, PictureId, Speaker, Turn, Verdict}; 15 16/// What one finished exchange offers the memory: her words, what the cat said back, and 17/// what a look at any picture found. 18#[derive(Clone, Debug, PartialEq, Eq)] 19pub struct ReflectInput { 20 pub heard: String, 21 pub said: String, 22 pub description: Option<String>, 23 pub pictures: Vec<PictureId>, 24} 25 26pub struct ReflectorConfig { 27 /// The most things Whiskers will remember; past it, new facts are refused and logged. 28 pub memory_limit: usize, 29 /// The chat is folded into its summary when it is longer than this many turns... 30 pub compress_after: usize, 31 /// ...leaving this many recent turns word for word. 32 pub keep_turns: usize, 33} 34 35/// The slow work that follows a conversation turn: deciding what to remember, filing it 36/// where it can be found by meaning, and folding the old chat into its summary. It runs 37/// off to the side of the conversation and takes no lock for longer than an instant, so 38/// the next thing she says is never waiting on it. 39pub struct Reflector { 40 pub model: Arc<dyn Model>, 41 pub guard: Arc<dyn Guard>, 42 pub embedder: Arc<dyn Embedder>, 43 pub log: SharedLog, 44 pub memory: SharedMemory, 45 pub chat: SharedChat, 46 pub clock: Arc<dyn Clock>, 47 pub profile: SharedProfile, 48 pub config: ReflectorConfig, 49 /// Picks the picture for each fact. A device with no service to ask has [`NoIcons`](crate::NoIcons). 50 pub icons: Arc<dyn IconChooser>, 51 /// Facts (by identity) for which `icons` said that nothing fits, so they are not asked about again at every sync 52 /// (Jev answers thirty questions a minute). Forgotten when the process ends. 53 pub iconified: Mutex<HashSet<String>>, 54} 55 56/// The most facts given a picture in one pass; the rest wait for the next. Each is an embedding and a Jev question. 57const ICONS_PER_PASS: usize = 3; 58 59/// Two facts this close in meaning are the same fact. 60const DUPLICATE_COSINE: f32 = 0.92; 61 62/// What one pass of the reflector found. 63#[derive(Debug, Default)] 64pub struct Reflected { 65 pub learned: Vec<Fact>, 66 /// The model could not be asked what to remember: the exchange is worth offering again. 67 /// (A refusal by the guard is an answer, not a failure, and is not retried.) 68 pub retry: bool, 69} 70 71impl Reflector { 72 /// Returns the facts learned from `input`, if any. Everything else it does is housekeeping. 73 pub fn run(&self, input: Option<&ReflectInput>) -> Vec<Fact> { 74 self.reflect(input).learned 75 } 76 77 /// Like [`run`](Self::run), saying whether `input` should be tried again later. 78 pub fn reflect(&self, input: Option<&ReflectInput>) -> Reflected { 79 debug!("reflect begins (input present = {})", input.is_some()); 80 self.backfill_embeddings(); 81 let (learned, retry) = match input.map(|i| self.remember(i)) { 82 Some(Ok(learned)) => (learned, false), 83 Some(Err(())) => (Vec::new(), true), 84 None => (Vec::new(), false), 85 }; 86 self.compress(); 87 self.give_icons(); 88 info!("reflect done: {} fact(s) learned, retry = {retry}", learned.len()); 89 Reflected { learned, retry } 90 } 91 92 /// Gives a picture to a few facts that have none. A failure leaves the fact as it was (its card shows a generic 93 /// star), is logged by its count and never reaches the child. 94 pub fn give_icons(&self) { 95 let wanting: Vec<Fact> = { 96 let tried = self.iconified.lock().unwrap_or_else(std::sync::PoisonError::into_inner); 97 self.memory.facts().into_iter().filter(|f| f.icon.is_none() && !tried.contains(&f.gid)).take(ICONS_PER_PASS).collect() 98 }; 99 for fact in wanting { 100 match self.icons.choose(&fact.text) { 101 Ok(Some(icon)) => match self.memory.give_icon(fact.id, icon) { 102 Ok(_) => debug!("icons: fact {} has a picture", fact.id), 103 Err(e) => warn!("icons: fact {} kept no picture: {}", fact.id, e.0), 104 }, 105 Ok(None) => { 106 // An answer: nothing fits. Not asked again until the process restarts. 107 debug!("icons: none fits fact {}", fact.id); 108 self.iconified.lock().unwrap_or_else(std::sync::PoisonError::into_inner).insert(fact.gid.clone()); 109 } 110 Err(e) => { 111 // "Could not decide": the service is probably down, so the rest wait for the next pass. 112 warn!("icons: fact {} has no picture for now: {}", fact.id, e.0); 113 return; 114 } 115 } 116 } 117 } 118 119 fn record(&self, event: Event) -> Result<(), ()> { 120 self.log.append(&Entry { at_ms: self.clock.now_ms(), event }).map_err(|e| error!("reflector: the parents' record could not be written: {}", e.0)) 121 } 122 123 fn failed(&self, error: String) { 124 warn!("reflector failure: {error}"); 125 let _ = self.record(Event::MemoryFailed { error }); 126 } 127 128 // ---- remembering ------------------------------------------------------ 129 130 /// `Err` only when the model could not be asked. 131 fn remember(&self, input: &ReflectInput) -> Result<Vec<Fact>, ()> { 132 let known: Vec<String> = self.memory.facts().into_iter().map(|f| f.text).collect(); 133 debug!("remember: {} known facts, heard {} chars, said {} chars", known.len(), input.heard.len(), input.said.len()); 134 let audience = self.profile.audience(); 135 let mut ask = format!("Already known: {known:?}\n{} said: {}\nWhiskers said: {}", audience.speaker_label(), input.heard, input.said); 136 if let Some(d) = &input.description { 137 ask.push_str(&format!("\nThe child showed a picture. A plain look at it found: {d}")); 138 } 139 let reply = match self.model.complete(&remember_prompt(&audience), &[Turn::said(Speaker::Child, ask)]) { 140 Ok(r) => r, 141 Err(e) => { 142 self.failed(format!("deciding what to remember: {}", e.reason)); 143 return Err(()); 144 } 145 }; 146 // Dull checks first, then ONE guard check for the whole batch: Jev admits thirty 147 // questions a minute across the machine and a guard check is two of them, so a check 148 // per fact would spend the minute on housekeeping. 149 let mut batch: Vec<NewFact> = Vec::new(); 150 let parsed = parse_facts(&reply); 151 debug!("remember: the model proposed {} fact(s)", parsed.len()); 152 for mut candidate in parsed { 153 candidate.pictures = input.pictures.clone(); 154 if known.iter().any(|k| k.eq_ignore_ascii_case(&candidate.text)) { 155 trace!("remember: a proposed fact is already known"); 156 continue; 157 } 158 if candidate.text.chars().count() > 160 { 159 self.not_remembered(&candidate.text, "too long"); 160 continue; 161 } 162 if self.memory.facts().len() + batch.len() >= self.config.memory_limit { 163 self.not_remembered(&candidate.text, "memory is full"); 164 continue; 165 } 166 batch.push(candidate); 167 } 168 if batch.is_empty() { 169 debug!("remember: nothing new to keep"); 170 return Ok(Vec::new()); 171 } 172 // A fact is fed back into the cat's prompt, so it passes the same guard as anything 173 // else the cat could say. 174 let joined = batch.iter().map(|c| c.text.as_str()).collect::<Vec<_>>().join(" "); 175 match self.guard.check(Direction::ToChild, audience.age(), &joined) { 176 Ok(Verdict::Allow) => debug!("remember: guard allowed {} fact(s)", batch.len()), 177 Ok(Verdict::Refuse { reason, .. }) => { 178 warn!("remember: guard refused {} fact(s)", batch.len()); 179 for c in &batch { 180 self.not_remembered(&c.text, &format!("guard: {reason}")); 181 } 182 return Ok(Vec::new()); 183 } 184 Err(e) => { 185 error!("remember: guard unavailable, {} fact(s) dropped: {}", batch.len(), e.0); 186 for c in &batch { 187 self.not_remembered(&c.text, &format!("guard unavailable: {}", e.0)); 188 } 189 return Ok(Vec::new()); 190 } 191 } 192 Ok(batch.into_iter().filter_map(|c| self.file(c)).collect()) 193 } 194 195 fn not_remembered(&self, fact: &str, why: &str) { 196 info!("not remembered ({} chars): {}", fact.len(), why.split(':').next().unwrap_or("")); 197 let _ = self.record(Event::NotRemembered { fact: fact.to_owned(), why: why.to_owned() }); 198 } 199 200 /// Files one approved fact, with its embedding, unless it means the same as something known. 201 fn file(&self, candidate: NewFact) -> Option<Fact> { 202 let text = candidate.text.clone(); 203 let embedding = match self.embedder.embed(&[text.clone()]) { 204 Ok(mut v) => v.pop(), 205 Err(e) => { 206 warn!("file: the embedder failed ({}); the fact is kept without one for now", e.0); 207 None 208 } 209 }; 210 if let Some(e) = &embedding { 211 if self.memory.facts().iter().any(|f| cosine(e, &f.embedding) >= DUPLICATE_COSINE) { 212 info!("file: a fact of {} chars means the same as one already known; dropped", text.len()); 213 return None; 214 } 215 } 216 let fact = match self.memory.add(candidate, self.clock.now_ms()) { 217 Ok(f) => f, 218 Err(e) => { 219 self.not_remembered(&text, &format!("could not save: {}", e.0)); 220 return None; 221 } 222 }; 223 if self.record(Event::Remembered { fact: fact.text.clone() }).is_err() { 224 // Not in the parents' log means not remembered. 225 error!("file: fact {} not in the parents' log, so it is forgotten again", fact.id); 226 let _ = self.memory.forget(fact.id); 227 return None; 228 } 229 if let Some(e) = embedding { 230 if let Err(err) = self.memory.set_embedding(fact.id, e) { 231 self.failed(format!("saving an embedding: {}", err.0)); 232 return Some(fact); 233 } 234 return self.memory.facts().into_iter().find(|f| f.id == fact.id); 235 } 236 Some(fact) 237 } 238 239 /// Facts filed while the embedder was down have no embedding yet; give them one. 240 fn backfill_embeddings(&self) { 241 let missing: Vec<Fact> = self.memory.facts().into_iter().filter(|f| f.embedding.is_empty()).collect(); 242 if missing.is_empty() { 243 trace!("backfill: every fact has an embedding"); 244 return; 245 } 246 debug!("backfill: {} fact(s) lack an embedding", missing.len()); 247 let texts: Vec<String> = missing.iter().map(|f| f.text.clone()).collect(); 248 match self.embedder.embed(&texts) { 249 Ok(vectors) if vectors.len() == missing.len() => { 250 for (f, v) in missing.iter().zip(vectors) { 251 if let Err(e) = self.memory.set_embedding(f.id, v) { 252 self.failed(format!("saving an embedding: {}", e.0)); 253 } 254 } 255 } 256 Ok(_) => self.failed("the embedder returned the wrong number of vectors".into()), 257 Err(e) => self.failed(format!("embedding: {}", e.0)), 258 } 259 } 260 261 // ---- compressing the chat --------------------------------------------- 262 263 fn compress(&self) { 264 let Some(CompressPlan { summary, oldest }) = self.chat.compress_plan(self.config.compress_after, self.config.keep_turns) 265 else { 266 return; 267 }; 268 let audience = self.profile.audience(); 269 debug!("compress: folding {} turns into a summary of {} chars", oldest.len(), summary.len()); 270 let mut ask = format!("Summary so far: {}\n\nNewer messages:\n", if summary.is_empty() { "(none yet)" } else { &summary }); 271 for t in &oldest { 272 let who = if t.speaker == Speaker::Child { audience.speaker_label() } else { "Whiskers" }; 273 ask.push_str(&format!("{who}: {}\n", t.text)); 274 } 275 match self.model.complete(&compress_prompt(&audience), &[Turn::said(Speaker::Child, ask)]) { 276 Ok(new_summary) if !new_summary.trim().is_empty() => { 277 // The summary goes into every later prompt, so it passes the guard like anything 278 // else the cat could say; if it does not, the chat is left as it was. 279 match self.guard.check(Direction::ToChild, audience.age(), new_summary.trim()) { 280 Ok(Verdict::Allow) => {} 281 Ok(Verdict::Refuse { reason, .. }) => return self.failed(format!("the chat summary was refused: {reason}")), 282 Err(e) => return self.failed(format!("the chat summary could not be checked: {}", e.0)), 283 } 284 match self.chat.apply_compression(oldest.len(), new_summary.trim().to_owned()) { 285 Ok(()) => { 286 info!("compress: {} turns folded into the summary", oldest.len()); 287 let _ = self.record(Event::Compressed { turns: oldest.len() }); 288 } 289 Err(e) => self.failed(format!("saving the chat: {}", e.0)), 290 } 291 } 292 Ok(_) => self.failed("the summary came back empty".into()), 293 Err(e) => self.failed(format!("summarising the chat: {}", e.reason)), 294 } 295 } 296} 297 298/// Facts out of the model's reply: the first JSON array in it, of objects or plain strings. 299pub(crate) fn parse_facts(reply: &str) -> Vec<NewFact> { 300 let (Some(start), Some(end)) = (reply.find('['), reply.rfind(']')) else { 301 debug!("parse_facts: no JSON array in a reply of {} chars", reply.len()); 302 return Vec::new(); 303 }; 304 if end < start { 305 debug!("parse_facts: brackets out of order in a reply of {} chars", reply.len()); 306 return Vec::new(); 307 } 308 let Ok(items) = serde_json::from_str::<Vec<serde_json::Value>>(&reply[start..=end]) else { 309 warn!("parse_facts: the bracketed part of a {} char reply is not a JSON array", reply.len()); 310 return Vec::new(); 311 }; 312 items 313 .into_iter() 314 .filter_map(|v| match v { 315 serde_json::Value::String(s) => { 316 let s = s.trim().to_owned(); 317 (!s.is_empty()).then(|| NewFact { text: s, ..NewFact::default() }) 318 } 319 serde_json::Value::Object(o) => { 320 let text = o.get("text")?.as_str()?.trim().to_owned(); 321 if text.is_empty() { 322 return None; 323 } 324 let opt = |k: &str| o.get(k).and_then(|x| x.as_str()).map(str::trim).filter(|s| !s.is_empty()).map(str::to_owned); 325 Some(NewFact { 326 text, 327 kind: o.get("kind").and_then(|k| k.as_str()).map_or(Kind::Other, Kind::from_word), 328 who: o 329 .get("who") 330 .and_then(|w| w.as_array()) 331 .map(|a| a.iter().filter_map(|x| x.as_str().map(|s| s.trim().to_owned())).filter(|s| !s.is_empty()).collect()) 332 .unwrap_or_default(), 333 place: opt("place"), 334 when: opt("when"), 335 pictures: Vec::new(), 336 }) 337 } 338 _ => None, 339 }) 340 .take(3) 341 .collect() 342}