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};

What one finished exchange offers the memory: her words, what the cat said back, and 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}
26pub struct ReflectorConfig {

The most things Whiskers will remember; past it, new facts are refused and logged.

28    pub memory_limit: usize,

The chat is folded into its summary when it is longer than this many turns...

30    pub compress_after: usize,

...leaving this many recent turns word for word.

32    pub keep_turns: usize,
33}

The slow work that follows a conversation turn: deciding what to remember, filing it where it can be found by meaning, and folding the old chat into its summary. It runs off to the side of the conversation and takes no lock for longer than an instant, so 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,

Picks the picture for each fact. A device with no service to ask has NoIcons.

50    pub icons: Arc<dyn IconChooser>,

Facts (by identity) for which icons said that nothing fits, so they are not asked about again at every sync (Jev answers thirty questions a minute). Forgotten when the process ends.

53    pub iconified: Mutex<HashSet<String>>,
54}

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;

Two facts this close in meaning are the same fact.

60const DUPLICATE_COSINE: f32 = 0.92;

What one pass of the reflector found.

63#[derive(Debug, Default)]
64pub struct Reflected {
65    pub learned: Vec<Fact>,

The model could not be asked what to remember: the exchange is worth offering again. (A refusal by the guard is an answer, not a failure, and is not retried.)

68    pub retry: bool,
69}
71impl Reflector {

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    }

Like 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    }

Gives a picture to a few facts that have none. A failure leaves the fact as it was (its card shows a generic 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    }
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 ------------------------------------------------------

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    }
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    }

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    }

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    }
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}

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}