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}