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}