Everything a shell needs, behind one object: the cat, the conversation, the parents' view and the voice allowance. The shell reports touches and what the microphone heard, and speaks and draws what comes back. This is the boundary the UniFFI crate exports.
10use log::{debug, error, info, trace, warn}; 11 12use whiskers_core::{HouseholdConfig, TimeKeeper, TimeStatus, 13 ChatError, Clock, Config, Conversation, Digest, DirPictures, Embedder, Fact, Guard, Heard, 14 Image, JsonChat, JsonMemory, JsonlLog, Model, Outcome, Parts, PictureId, Recall, ReflectInput, Reflector, 15 ReflectorConfig, SharedChat, SharedProfile, SharedLog, SharedMemory, read_shared_log, 16}; 17use whiskers_gateway::Gateway; 18use whiskers_guard::{RemoteEmbedder, RemoteGuard, RemoteRanker, RemoteSync}; 19pub use whiskers_pet::{Mood, Phase, PetView, Touch, TouchKind}; 20 21pub struct EngineConfig { 22 pub data_dir: PathBuf, 23 pub gateway_url: String, 24 pub guard_url: String, 25 pub model: String,
The parents' own prompt; None uses the default, written for the child's profile.
What Whiskers says. The text is the only thing a shell may speak.
The parents' view of a stretch of time.
53pub struct ExchangeView { 54 pub at_ms: u64, 55 pub heard: String, 56 pub pictures: Vec<String>, 57 pub model_wrote: Option<String>, 58 pub said: String, 59 pub outcome: Outcome, 60 pub needs_a_grown_up: bool, 61 pub notes: Vec<String>, 62} 63 64pub struct Engine { 65 pet: Mutex<whiskers_pet::Pet>, 66 convo: Mutex<Conversation>,
Exchanges waiting for the memory to consider them (see reflect).
None for an engine built without a service (tests).
76 sync: Option<RemoteSync>,
One sync at a time: two at once would each append the same lines.
The copy of the service's log that sync keeps (every device's entries, one line each).
81 shared_log: PathBuf,
The child, kept level with the household document (settings change it, sync can too).
Names this device in the day's time count and nowhere else. Made once and kept.
87use std::hash::BuildHasher as _; 88fn device_id(dir: &std::path::Path) -> Result<String, String> { 89 let path = dir.join("device.id"); 90 if let Ok(id) = std::fs::read_to_string(&path) { 91 if !id.trim().is_empty() { 92 trace!("device id read from {}", path.display()); 93 return Ok(id.trim().to_owned()); 94 } 95 } 96 let id = format!("{:x}{:x}", SystemClock.now_ms(), std::process::id()) 97 + &std::collections::hash_map::RandomState::new().hash_one(SystemClock.now_ms()).to_string(); 98 std::fs::write(&path, &id).map_err(|e| { 99 error!("cannot write the device id to {}: {e}", path.display()); 100 e.to_string() 101 })?; 102 info!("this device has a new id ({} chars)", id.len()); 103 Ok(id) 104}
The whole lines of a log file; a last line without its newline is a write in progress.
107fn complete_lines(path: &std::path::Path) -> Vec<String> { 108 let bytes = std::fs::read(path).unwrap_or_else(|e| { 109 debug!("{} not readable ({e}); treated as empty", path.display()); 110 Vec::new() 111 }); 112 let text = String::from_utf8_lossy(&bytes); 113 match text.rfind('\n') { 114 Some(i) => text[..=i].lines().map(str::to_owned).collect(), 115 None => Vec::new(), 116 } 117}
How many exchanges wait to be thought about: a long outage must not queue the whole evening.
120const PENDING_LIMIT: usize = 8;
122fn lock<T>(m: &Mutex<T>) -> MutexGuard<'_, T> { 123 // A panic while holding a lock must not take the cat's face down with it. 124 m.lock().unwrap_or_else(PoisonError::into_inner) 125} 126 127impl Engine { 128 pub fn new(cfg: EngineConfig) -> Result<Self, String> { 129 info!("engine starting: model {}, data dir {}", cfg.model, cfg.data_dir.display()); 130 std::fs::create_dir_all(&cfg.data_dir).map_err(|e| { 131 error!("cannot create the data directory {}: {e}", cfg.data_dir.display()); 132 e.to_string() 133 })?; 134 let summarizer = Box::new(Gateway::new(&cfg.gateway_url, &cfg.model, 600)); 135 let embedder: Arc<RemoteEmbedder> = Arc::new(RemoteEmbedder::new(&cfg.guard_url)); 136 let ranker = Arc::new(RemoteRanker::new(&cfg.guard_url)); 137 let model: Arc<dyn Model> = Arc::new(Gateway::new(&cfg.gateway_url, &cfg.model, 400)); 138 let guard: Arc<dyn Guard> = Arc::new(RemoteGuard::new(&cfg.guard_url)); 139 // A state file that cannot be read is kept aside and the app starts clean, rather than 140 // being unable to start at all; the parents' log says it happened. 141 let mut set_aside = Vec::new(); 142 let (chat_store, aside) = JsonChat::open_or_set_aside(&cfg.data_dir.join("chat.json")).map_err(|e: ChatError| e.0)?; 143 set_aside.extend(aside.map(|p| ("the chat", p))); 144 let chat = SharedChat::open(Box::new(chat_store)).map_err(|e: ChatError| e.0)?; 145 let (memory_store, aside) = JsonMemory::open_or_set_aside(&cfg.data_dir.join("memory.json")).map_err(|e| e.0)?; 146 set_aside.extend(aside.map(|p| ("what Whiskers remembers", p))); 147 let log = SharedLog::new(Box::new(JsonlLog::open(&cfg.data_dir.join("log.jsonl")).map_err(|e| e.0)?)); 148 for (what, kept_at) in set_aside { 149 warn!("{what}: the damaged file was set aside at {}", kept_at.display()); 150 let _ = log.append(&whiskers_core::Entry { 151 at_ms: SystemClock.now_ms(), 152 event: whiskers_core::Event::MemoryFailed { 153 error: format!("the file for {what} was damaged; it is kept at {} and Whiskers started that part afresh", kept_at.display()), 154 }, 155 }); 156 } 157 let parts = Parts { 158 model, 159 guard, 160 log, 161 pictures: Box::new(DirPictures::open(&cfg.data_dir.join("pictures")).map_err(|e| e.0)?), 162 memory: SharedMemory::new(Box::new(memory_store)), 163 chat, 164 recall: Arc::new(Recall::new(embedder.clone(), ranker)), 165 clock: Arc::new(SystemClock), 166 profile: SharedProfile::default(), 167 }; 168 let service = cfg.guard_url.clone(); 169 let mut engine = Self::with_parts(cfg, parts, embedder, summarizer)?; 170 engine.sync = Some(RemoteSync::new(service)); 171 info!("engine ready (device {})", engine.device); 172 Ok(engine) 173 }
For tests and shells that bring their own adapters.
176 pub fn with_parts( 177 cfg: EngineConfig, 178 parts: Parts, 179 embedder: Arc<dyn Embedder>, 180 summarizer: Box<dyn Model>, 181 ) -> Result<Self, String> { 182 debug!("engine parts assembled; system prompt is {}", if cfg.system_prompt.is_some() { "the parents'" } else { "the default" }); 183 let cfg_data_dir = cfg.data_dir.clone(); 184 let profile = parts.profile.clone(); 185 let reflector = Reflector { 186 model: parts.model.clone(), 187 guard: parts.guard.clone(), 188 embedder, 189 log: parts.log.clone(), 190 memory: parts.memory.clone(), 191 chat: parts.chat.clone(), 192 clock: parts.clock.clone(), 193 profile: parts.profile.clone(), 194 config: ReflectorConfig { memory_limit: 200, compress_after: 24, keep_turns: 10 }, 195 }; 196 let convo = Conversation::new( 197 Config { 198 system_prompt: cfg.system_prompt, 199 history_turns: 24, 200 greet_after_ms: 30 * 60 * 1000, 201 }, 202 parts, 203 ); 204 let time = TimeKeeper::open(&cfg_data_dir.join("household.json"), &device_id(&cfg_data_dir)?); 205 profile.set(time.config().child); 206 Ok(Self { 207 pet: Mutex::new(whiskers_pet::Pet::new()), 208 convo: Mutex::new(convo), 209 pending: Mutex::new(Vec::new()), 210 reflect_lock: Mutex::new(()), 211 reflector, 212 log_path: cfg.data_dir.join("log.jsonl"), 213 pictures: DirPictures::open(&cfg.data_dir.join("pictures")).map_err(|e| e.0)?, 214 summarizer, 215 time: Mutex::new(time), 216 profile, 217 sync: None, 218 sync_lock: Mutex::new(()), 219 shared_log: cfg_data_dir.join("shared-log.jsonl"), 220 device: device_id(&cfg_data_dir)?, 221 }) 222 }
224 // ---- the cat -------------------------------------------------------- 225 226 pub fn touch(&self, t: Touch) { 227 lock(&self.pet).touch(t); 228 } 229 230 pub fn set_phase(&self, phase: Phase, now_ms: u64) { 231 trace!("engine set_phase {phase:?}"); 232 lock(&self.pet).set_phase(phase, now_ms); 233 } 234 235 pub fn tick(&self, now_ms: u64) -> PetView { 236 lock(&self.pet).tick(now_ms) 237 } 238 239 // ---- the conversation (blocking: call off the main thread) -----------
Whether Whiskers should say hello: the first time ever, or after a long quiet. A phone folded or unfolded, or an app brought back within half an hour, is not a new chat.
247 pub fn greet(&self) -> Spoken { 248 debug!("engine greet"); 249 let r = lock(&self.convo).greet(); 250 Spoken { text: r.say.as_str().to_owned(), outcome: r.outcome } 251 } 252 253 pub fn hear(&self, text: &str, pictures: Vec<Image>) -> Spoken { 254 // Held for the turn only; the slow memory work is `reflect`, which does not take this lock. 255 let started = Instant::now(); 256 debug!("engine hear: {} chars, {} picture(s)", text.len(), pictures.len()); 257 let r = lock(&self.convo).respond(Heard { text: text.to_owned(), pictures }); 258 info!("engine hear done in {} ms: {:?}", started.elapsed().as_millis(), r.outcome); 259 if let Some(input) = r.reflect { 260 let mut pending = lock(&self.pending); 261 pending.push(input); 262 // Bounded: a long outage must not queue the whole evening. 263 if pending.len() > PENDING_LIMIT { 264 warn!("{} exchanges wait for the memory; dropping the oldest", pending.len()); 265 pending.remove(0); 266 } 267 debug!("{} exchange(s) wait for the memory", pending.len()); 268 } 269 Spoken { text: r.say.as_str().to_owned(), outcome: r.outcome } 270 }
The slow work after a turn: deciding what to remember, filing it, folding the old chat into its summary. Blocking, but it never holds the conversation, so she can talk again at once. Returns what was learned. Only one runs at a time; a call that finds another running returns nothing and leaves its work queued for the next.
276 pub fn reflect(&self) -> Vec<Fact> { 277 let _running = match self.reflect_lock.try_lock() { 278 Ok(g) => g, 279 // A panic inside an earlier pass poisons the lock; that must not end remembering for good. 280 Err(std::sync::TryLockError::Poisoned(p)) => p.into_inner(), 281 Err(std::sync::TryLockError::WouldBlock) => { 282 debug!("reflect: another pass is running; leaving the work queued"); 283 return Vec::new(); 284 } 285 }; 286 let started = Instant::now(); 287 let mut learned = Vec::new(); 288 loop { 289 let inputs: Vec<ReflectInput> = std::mem::take(&mut *lock(&self.pending)); 290 debug!("reflect pass over {} exchange(s)", inputs.len()); 291 if inputs.is_empty() { 292 self.reflector.run(None); 293 } 294 let mut retry = Vec::new(); 295 for input in inputs { 296 let r = self.reflector.reflect(Some(&input)); 297 learned.extend(r.learned); 298 if r.retry { 299 retry.push(input); 300 } 301 } 302 // Exchanges that arrived while this pass was busy: a caller that found the lock held 303 // returned at once, so nobody else will pick these up. 304 let newer = lock(&self.pending).len(); 305 if !retry.is_empty() { 306 warn!("reflect: {} exchange(s) will be tried again later (model unreachable)", retry.len()); 307 // The model was unreachable: put these back ahead of anything newer, within the bound. 308 let mut pending = lock(&self.pending); 309 retry.append(&mut pending); 310 *pending = retry; 311 let excess = pending.len().saturating_sub(PENDING_LIMIT); 312 pending.drain(..excess); 313 } 314 // Retried ones wait for the next turn; going round again for them would spin on an outage. 315 if newer == 0 { 316 break; 317 } 318 } 319 info!("reflect finished in {} ms: {} fact(s) learned", started.elapsed().as_millis(), learned.len()); 320 learned 321 }
323 // ---- the grown-ups' time limits -----------------------------------------
Every choice the grown-ups have made, as this device last heard of them.
The grown-ups changed a choice. Call sync soon so the other devices hear.
The pictures that go with remembered facts follow the facts: the ones this device lacks are fetched, and the ones the service lacks are sent. Returns how many arrived here.
339 fn sync_pictures(&self, remote: &RemoteSync, wanted: &[PictureId]) -> Result<u32, String> { 340 let mut arrived = 0; 341 debug!("sync pictures: {} wanted", wanted.len()); 342 for id in wanted.iter().filter(|id| !self.pictures.has(id)) { 343 if let Some(bytes) = remote.get_picture(&id.0)? { 344 self.pictures.store(id, &bytes).map_err(|e| { 345 error!("sync pictures: cannot keep {}: {}", id.0, e.0); 346 e.0 347 })?; 348 arrived += 1; 349 } else { 350 warn!("sync pictures: {} is wanted but the service does not have it", id.0); 351 } 352 } 353 let here: Vec<String> = wanted.iter().filter(|id| self.pictures.has(id)).map(|id| id.0.clone()).collect(); 354 if !here.is_empty() { 355 for id in remote.pictures_missing(&here)? { 356 if let Some(bytes) = self.pictures.read(&PictureId(id.clone())) { 357 remote.put_picture(&id, &bytes)?; 358 } else { 359 warn!("sync pictures: the service lacks {id} but it is not readable here"); 360 } 361 } 362 } 363 if arrived > 0 { 364 info!("sync pictures: {arrived} picture(s) arrived"); 365 } 366 Ok(arrived) 367 }
Brings this device level with the service: what Whiskers remembers, and the grown-ups'
choices and today's time. Returns how many things changed here. Blocking, and safe to
call whenever there is a network; an unreachable service is an Err and changes nothing.
Every picture some device kept for a memory or for a conversation, which every device wants.
373 fn pictures_everyone_wants(&self) -> Vec<PictureId> { 374 let mut v: Vec<PictureId> = self.reflector.memory.snapshot().facts.iter().flat_map(|f| f.pictures.iter().cloned()).collect(); 375 if let Ok((sources, _)) = read_shared_log(&self.log_path, &self.shared_log, &self.device) { 376 for e in sources.iter().flatten() { 377 if let whiskers_core::Event::Heard { pictures, .. } = &e.event { 378 v.extend(pictures.iter().cloned()); 379 } 380 } 381 } 382 v.retain(PictureId::is_safe); 383 v.sort_by(|a, b| a.0.cmp(&b.0)); 384 v.dedup(); 385 v 386 }
This device's new log entries go up to the service and the one log comes down, so the parents see every conversation wherever it happened. The log only ever grows, so this is just "send what the other side lacks". Returns how many lines arrived here.
391 fn sync_journal(&self, remote: &RemoteSync) -> Result<u32, String> { 392 let mine = complete_lines(&self.log_path); 393 let mut have = remote.journal_have(&self.device)?; 394 debug!("sync journal: {} local line(s), the service has {have}", mine.len()); 395 if have > mine.len() { 396 warn!("sync journal: the service has {have} lines from this device but only {} exist here", mine.len()); 397 } 398 while have < mine.len() { 399 let end = (have + 200).min(mine.len()); 400 let now = remote.journal_push(&self.device, have, &mine[have..end])?; 401 if now <= have { 402 error!("sync journal: pushed from line {have} but the service now reports {now}"); 403 return Err("the service did not take the log".into()); 404 } 405 have = now; 406 } 407 let mut arrived = 0; 408 loop { 409 let held = complete_lines(&self.shared_log).len(); 410 let got = remote.journal_pull(held)?; 411 if got.total < held { 412 // Our copy is longer than the log it copies, so it is not a copy of it (it can only 413 // have been damaged): throw it away and take the log again from the start. 414 warn!("sync journal: our copy has {held} lines but the log has {}; rebuilding the copy from the start", got.total); 415 if let Err(e) = std::fs::remove_file(&self.shared_log) { 416 error!("sync journal: cannot remove the damaged copy: {e}"); 417 } 418 continue; 419 } 420 // Only lines that continue exactly where our copy ends are taken. 421 if got.lines.is_empty() || got.from != held { 422 if got.from != held { 423 warn!("sync journal: pulled lines start at {} but our copy ends at {held}; taking none", got.from); 424 } 425 break; 426 } 427 use std::io::Write as _; 428 let mut text = got.lines.join("\n"); 429 text.push('\n'); 430 let mut f = std::fs::OpenOptions::new().create(true).append(true).open(&self.shared_log).map_err(|e| { 431 error!("sync journal: cannot open the shared log copy: {e}"); 432 e.to_string() 433 })?; 434 f.write_all(text.as_bytes()).and_then(|()| f.sync_data()).map_err(|e| { 435 error!("sync journal: cannot append to the shared log copy: {e}"); 436 e.to_string() 437 })?; 438 arrived += got.lines.len() as u32; 439 } 440 info!("sync journal: {arrived} line(s) arrived"); 441 Ok(arrived) 442 }
Brings this device level with the service: what Whiskers remembers, every device's
conversations, the pictures that go with both, and the grown-ups' choices and today's
time. Returns how many things changed here. Blocking. Each part is tried even if another
failed; an unreachable service is an Err and changes nothing.
448 pub fn sync(&self) -> Result<u32, String> { 449 let Some(remote) = &self.sync else { 450 warn!("sync asked for but there is no service configured"); 451 return Err("no service to sync with".into()); 452 }; 453 // If another sync is under way it is already doing this; there is nothing to add. 454 let Ok(_one_at_a_time) = self.sync_lock.try_lock() else { 455 debug!("sync skipped: another is under way"); 456 return Ok(0); 457 }; 458 info!("sync starts"); 459 let started = Instant::now(); 460 let mut changed = 0u32; 461 let mut errors: Vec<String> = Vec::new(); 462 let mut step = |what: &str, r: Result<u32, String>| match r { 463 Ok(n) => { 464 debug!("sync {what}: {n} change(s)"); 465 changed += n; 466 } 467 Err(e) => { 468 warn!("sync {what} failed: {e}"); 469 errors.push(format!("{what}: {e}")); 470 } 471 }; 472 let memory = self.reflector.memory.clone(); 473 step("memory", remote.memory(&memory.snapshot()).and_then(|theirs| memory.merge(theirs).map(|n| n as u32).map_err(|e| e.0))); 474 // The one session: whichever copy has had the most said in it is the one everyone continues. 475 let chat = self.reflector.chat.clone(); 476 step("session", remote.chat(&chat.snapshot()).and_then(|theirs| chat.adopt(theirs).map(u32::from).map_err(|e| e.0))); 477 step("conversations", self.sync_journal(remote)); 478 step("pictures", self.sync_pictures(remote, &self.pictures_everyone_wants())); 479 let mine = lock(&self.time).household(); 480 step("settings", remote.household(&mine).map(|theirs| u32::from(lock(&self.time).merge(&theirs)))); 481 // Settings from another device may have changed the child; the next turn uses it. 482 self.profile.set(lock(&self.time).config().child); 483 if errors.is_empty() { 484 info!("sync finished in {} ms: {changed} change(s)", started.elapsed().as_millis()); 485 Ok(changed) 486 } else { 487 error!("sync finished in {} ms with {} failed part(s), {changed} change(s)", started.elapsed().as_millis(), errors.len()); 488 Err(errors.join("; ")) 489 } 490 }
delta_ms was spent with Whiskers on day (a number that changes at midnight).
The grown-ups give minutes more today.
510 pub fn time_used_minutes(&self, day: u32) -> u32 { 511 lock(&self.time).used_minutes(day) 512 } 513 514 // ---- the parents' view ------------------------------------------------ 515 516 pub fn digest(&self, from_ms: u64, to_ms: u64) -> Result<DigestView, String> { 517 debug!("engine digest [{from_ms}, {to_ms})"); 518 let (sources, unreadable) = read_shared_log(&self.log_path, &self.shared_log, &self.device).map_err(|e| e.0)?; 519 let d = Digest::between_all(&sources, from_ms, to_ms); 520 let v = view(&d, unreadable); 521 debug!("engine digest: {} exchange(s), {unreadable} unreadable line(s)", v.exchanges.len()); 522 Ok(v) 523 } 524 525 pub fn summarize(&self, from_ms: u64, to_ms: u64) -> Result<String, String> { 526 info!("engine summarize [{from_ms}, {to_ms})"); 527 let (sources, _) = read_shared_log(&self.log_path, &self.shared_log, &self.device).map_err(|e| e.0)?; 528 Digest::between_all(&sources, from_ms, to_ms).summarize(self.summarizer.as_ref(), &self.profile.audience()) 529 .inspect(|s| debug!("engine summarize: note of {} chars", s.len())) 530 .map_err(|e| { 531 error!("engine summarize failed: {}", e.reason); 532 e.reason 533 }) 534 }
The child's journal. Reads the memory directly, so the parents' view never waits on a turn.
A parent removes something Whiskers remembers; the log keeps a record that they did.
543 pub fn forget(&self, id: u64) -> bool { 544 match self.reflector.memory.forget(id) { 545 Ok(fact) => { 546 info!("engine forget: fact {id} removed by a parent"); 547 let entry = whiskers_core::Entry { at_ms: SystemClock.now_ms(), event: whiskers_core::Event::Forgot { fact: fact.text } }; 548 if let Err(e) = self.reflector.log.append(&entry) { 549 warn!("engine forget: fact {id} removed but the record could not be written: {}", e.0); 550 } 551 true 552 } 553 Err(e) => { 554 warn!("engine forget: fact {id} not removed: {}", e.0); 555 false 556 } 557 } 558 }
560 pub fn picture_path(&self, id: &str) -> PathBuf { 561 self.pictures.path_of(&PictureId(id.to_owned())) 562 } 563 564} 565 566fn view(d: &Digest, unreadable: usize) -> DigestView { 567 let (answered, stopped) = d.counts(); 568 let urgent: Vec<u64> = d.needs_a_grown_up().iter().map(|e| e.at_ms).collect(); 569 let mut exchanges: Vec<ExchangeView> = d 570 .exchanges 571 .iter() 572 .map(|e| ExchangeView { 573 at_ms: e.at_ms, 574 heard: e.heard.clone(), 575 pictures: e.pictures.iter().map(|p| p.0.clone()).collect(), 576 model_wrote: e.model_wrote.clone(), 577 said: e.said.clone(), 578 outcome: e.outcome, 579 needs_a_grown_up: urgent.contains(&e.at_ms), 580 notes: e.notes.clone(), 581 }) 582 .collect(); 583 // Needs-a-grown-up first, then the rest in the order they happened. 584 exchanges.sort_by_key(|e| (!e.needs_a_grown_up, e.at_ms)); 585 DigestView { 586 exchanges, 587 facts_learned: d.facts_learned.clone(), 588 answered: answered as u32, 589 stopped: stopped as u32, 590 unreadable_lines: unreadable as u32, 591 } 592}