diff --git a/crates/tinymemory-tools/src/lifecycle/mod.rs b/crates/tinymemory-tools/src/lifecycle/mod.rs index 50d3f9fa..ffe30c34 100644 --- a/crates/tinymemory-tools/src/lifecycle/mod.rs +++ b/crates/tinymemory-tools/src/lifecycle/mod.rs @@ -7,8 +7,9 @@ //! //! | When | Call | Engine work | //! | --- | --- | --- | -//! | session start or resume | [`AgentMemory::start_session`] | reads only | +//! | session start, before any user turn | [`AgentMemory::start_session`] | reads only | //! | user turn, before the model | [`AgentMemory::pre_turn`] | logs the turn (accepted, not indexed) while fetching the pack | +//! | first user turn after a compaction | [`AgentMemory::pre_turn_resumed`] | as `pre_turn`, with the thread's earlier turns leading the same pack | //! | after the reply | [`AgentMemory::post_turn`] | logs the reply; may return a belief build (never on an engine that builds on its own) | //! | prompt truncated | [`AgentMemory::recall_for_compaction`] | an answered summary of the thread, plus related memory | //! | any time | [`AgentMemory::recall`] | a pre-turn pack without logging | @@ -330,14 +331,7 @@ impl AgentMemory { let mut sections = Vec::new(); if let Some(thread_id) = &start.thread_id { let thread_id = non_blank(thread_id, "thread id")?; - sections.push(ScopeSection::latest( - THREAD_HEADING, - MetaFilter { - thread_id: Some(thread_id.to_string()), - ..self.layout.conversations_filter(Some(&self.agent_id)) - }, - self.policy.history_limit.max(1), - )); + sections.push(self.thread_section(thread_id, self.policy.history_limit.max(1))); } sections.extend(self.standard_sections()); self.read(start.focus, sections).await @@ -355,6 +349,24 @@ impl AgentMemory { /// [`Error::InvalidRequest`] for a blank thread id or text. Engine /// failures never fail the call. pub async fn pre_turn(&self, turn: PreTurn) -> Result { + self.pre_turn_with(turn, false).await + } + + /// [`AgentMemory::pre_turn`] for a session that resumes this thread after + /// a compaction: the one pack leads with the thread's earlier turns (those + /// before [`PreTurn::in_prompt_from`]), sharing the turn's budget and the + /// cross-section dedupe, instead of a separate + /// [`AgentMemory::start_session`] pack pasted beside it. A policy with + /// `history_limit == 0` gets no thread section, as everywhere else. + /// + /// # Errors + /// + /// As [`AgentMemory::pre_turn`]. + pub async fn pre_turn_resumed(&self, turn: PreTurn) -> Result { + self.pre_turn_with(turn, true).await + } + + async fn pre_turn_with(&self, turn: PreTurn, resumed: bool) -> Result { let thread_id = non_blank(&turn.thread_id, "thread id")?; let text = non_blank(&turn.user_text, "user text")?; let item = self.turn_item( @@ -370,7 +382,12 @@ impl AgentMemory { thread_id: thread_id.to_string(), from_turn: turn.in_prompt_from, }; - let mut request = self.request(Some(text.to_string()), self.standard_sections()); + let mut sections = Vec::new(); + if resumed && self.policy.history_limit > 0 { + sections.push(self.thread_section(thread_id, self.policy.history_limit)); + } + sections.extend(self.standard_sections()); + let mut request = self.request(Some(text.to_string()), sections); request.exclude_ids = vec![id]; request.exclude_thread = Some(window); let (logged, pack) = join( @@ -507,6 +524,18 @@ impl AgentMemory { /// Learnings, each core scope, brain, this agent's history, then the /// team's, each filled by fetch; a zero limit leaves its section out. + /// The thread's own latest turns, up to `limit`. + fn thread_section(&self, thread_id: &str, limit: usize) -> ScopeSection { + ScopeSection::latest( + THREAD_HEADING, + MetaFilter { + thread_id: Some(thread_id.to_string()), + ..self.layout.conversations_filter(Some(&self.agent_id)) + }, + limit, + ) + } + fn standard_sections(&self) -> Vec { let policy = &self.policy; let learnings = ( diff --git a/crates/tinymemory-tools/src/lifecycle/mod_tests.rs b/crates/tinymemory-tools/src/lifecycle/mod_tests.rs index 4acaf855..aefd1029 100644 --- a/crates/tinymemory-tools/src/lifecycle/mod_tests.rs +++ b/crates/tinymemory-tools/src/lifecycle/mod_tests.rs @@ -648,3 +648,144 @@ async fn core_build_consolidates_exactly_the_core_node() { assert_eq!(request.reach, Reach::exact(acme())); assert!(request.kinds.is_empty()); } + +// ── a resumed pre-turn: one pack, one budget (openhuman#7023) ─────────────── + +/// Logs `turns` exchanges on `thread`, each with distinct text. +async fn log_exchanges(memory: &AgentMemory, thread: &str, turns: u32) { + for turn in 0..turns { + memory + .pre_turn(PreTurn::new( + thread, + turn * 2, + format!("leg {turn}: how far is Porto"), + )) + .await + .unwrap(); + memory + .post_turn(PostTurn::new( + thread, + turn * 2 + 1, + format!("Leg {turn} is about 310 km."), + )) + .await + .unwrap(); + } +} + +fn bullets(markdown: &str) -> Vec<&str> { + markdown + .lines() + .filter(|line| line.trim_start().starts_with("- ")) + .collect() +} + +#[tokio::test] +async fn a_resumed_pre_turn_leads_with_the_thread_in_one_pack_and_repeats_nothing() { + let engine = Arc::new(ReferenceEngine::new()); + let memory = memory(&engine, "assistant"); + engine + .store(StoreItem::learning( + "The user prefers metric units", + LearningKind::Preference, + 0.9, + MemoryMeta::default(), + )) + .await + .unwrap(); + log_exchanges(&memory, "t1", 3).await; + + let resumed = memory + .pre_turn_resumed(PreTurn { + in_prompt_from: 4, + ..PreTurn::new("t1", 6, "which units does the user prefer for the distance") + }) + .await + .unwrap() + .pack; + let markdown = &resumed.markdown; + assert!(markdown.contains(THREAD_HEADING), "{markdown}"); + assert!(markdown.contains("metric units"), "{markdown}"); + assert!( + markdown.contains("leg 0"), + "an older turn leads: {markdown}" + ); + assert!( + !markdown.contains("leg 2:"), + "turns in the prompt stay out: {markdown}" + ); + let lines = bullets(markdown); + let unique: std::collections::HashSet<&str> = lines.iter().copied().collect(); + assert_eq!( + lines.len(), + unique.len(), + "a line was injected twice:\n{markdown}" + ); + assert!(resumed.tokens <= RecallPolicy::default().budget_tokens); + + // An ordinary pre-turn has no thread section, as before. + let plain = memory + .pre_turn(PreTurn { + in_prompt_from: 4, + ..PreTurn::new("t1", 8, "which units does the user prefer for the distance") + }) + .await + .unwrap() + .pack; + assert!( + !plain.markdown.contains(THREAD_HEADING), + "{}", + plain.markdown + ); +} + +#[tokio::test] +async fn a_resumed_pre_turn_honours_a_zero_history_limit() { + let engine = Arc::new(ReferenceEngine::new()); + let memory = memory(&engine, "assistant"); + log_exchanges(&memory, "t1", 3).await; + let quiet = memory.clone().with_policy(RecallPolicy { + history_limit: 0, + ..memory.policy().clone() + }); + + let pack = quiet + .pre_turn_resumed(PreTurn { + in_prompt_from: 4, + ..PreTurn::new("t1", 6, "how far is Porto") + }) + .await + .unwrap() + .pack; + assert!(!pack.markdown.contains(THREAD_HEADING), "{}", pack.markdown); +} + +#[tokio::test] +async fn older_turns_survive_a_prompt_window_bigger_than_the_overfetch() { + // 30 exchanges = turns 0..59; the prompt holds turns 8..59 (52 turns), + // far more than the section's overfetch allowance. Before the fix the + // newest candidates were cut to that allowance first, all of them were + // in the window, and the thread section came back empty. + let engine = Arc::new(ReferenceEngine::new()); + let memory = memory(&engine, "assistant"); + log_exchanges(&memory, "t1", 30).await; + + let pack = memory + .pre_turn_resumed(PreTurn { + in_prompt_from: 8, + ..PreTurn::new("t1", 60, "how far is Porto") + }) + .await + .unwrap() + .pack; + let markdown = &pack.markdown; + assert!(markdown.contains(THREAD_HEADING), "{markdown}"); + assert!( + markdown.contains("leg 3"), + "the newest turn before the window: {markdown}" + ); + assert!( + !markdown.contains("leg 4:"), + "turns in the window stay out: {markdown}" + ); +} diff --git a/crates/tinymemory-tools/src/recall/gather.rs b/crates/tinymemory-tools/src/recall/gather.rs index 033c3cc1..6883bf6e 100644 --- a/crates/tinymemory-tools/src/recall/gather.rs +++ b/crates/tinymemory-tools/src/recall/gather.rs @@ -80,6 +80,7 @@ pub(super) async fn section( section: &ScopeSection, beliefs: usize, ) -> Gathered { + let keep = |hit: &Hit| !request.excludes(hit); let want = wanted(request, section); let outcome = match §ion.query { SectionQuery::Answer { @@ -94,7 +95,7 @@ pub(super) async fn section( "[recall] answer failed, fetching instead heading={:?} error={error}", section.heading ); - fetch(engine, §ion.filter, question, want, 0).await + fetch(engine, §ion.filter, question, want, 0, &keep).await } Err(error) => Err(error), }, @@ -104,11 +105,11 @@ pub(super) async fn section( .or(request.query.as_deref()) .filter(|query| !query.trim().is_empty()) { - Some(query) => fetch(engine, §ion.filter, query, want, beliefs).await, - None => with_listed_beliefs(engine, section, want).await, + Some(query) => fetch(engine, §ion.filter, query, want, beliefs, &keep).await, + None => with_listed_beliefs(engine, section, want, &keep).await, } } - SectionQuery::Latest => with_listed_beliefs(engine, section, want).await, + SectionQuery::Latest => with_listed_beliefs(engine, section, want, &keep).await, }; match outcome { Ok((hits, beliefs)) => Gathered::Hits { hits, beliefs }, @@ -235,9 +236,13 @@ async fn with_listed_beliefs( engine: &dyn MemoryEngine, section: &ScopeSection, want: usize, + keep: &dyn Fn(&Hit) -> bool, ) -> tinymemory_api::Result<(Vec, Vec)> { if !reads_learnings(section) { - return Ok((latest(engine, §ion.filter, want).await?, Vec::new())); + return Ok(( + latest(engine, §ion.filter, want, keep).await?, + Vec::new(), + )); } let reach = section .filter @@ -245,7 +250,7 @@ async fn with_listed_beliefs( .clone() .unwrap_or_else(|| Reach::subtree(Namespace::ROOT)); let (hits, beliefs) = join( - latest(engine, §ion.filter, want), + latest(engine, §ion.filter, want, keep), engine.beliefs(BeliefsRequest::new(reach, want)), ) .await; @@ -319,9 +324,10 @@ async fn fetch( query: &str, limit: usize, beliefs: usize, + keep: &dyn Fn(&Hit) -> bool, ) -> tinymemory_api::Result<(Vec, Vec)> { let Some(mode) = preferred_mode(engine) else { - return Ok((latest(engine, filter, limit).await?, Vec::new())); + return Ok((latest(engine, filter, limit, keep).await?, Vec::new())); }; let mut request = FetchRequest::new(query, mode, limit); request.filter = filter.clone(); @@ -332,10 +338,16 @@ async fn fetch( /// The newest hits, then the most confident, then the latest turn; ties /// keep the engine's order. +/// +/// Hits `keep` refuses (the request's exclusions: the prompt's own thread +/// window, ids already shown) are dropped *before* the cut to `limit`, so a +/// window holding more recent turns than the overfetch allowance cannot +/// crowd every older turn out of the section. async fn latest( engine: &dyn MemoryEngine, filter: &MetaFilter, limit: usize, + keep: &dyn Fn(&Hit) -> bool, ) -> tinymemory_api::Result> { let mut all: Vec = Vec::new(); let mut cursor: Option = None; @@ -343,7 +355,7 @@ async fn latest( let mut request = ListRequest::new(filter.clone(), LATEST_PAGE); request.cursor = cursor.take(); let page = engine.list(request).await?; - all.extend(page.items); + all.extend(page.items.into_iter().filter(|hit| keep(hit))); match page.next_cursor { Some(next) => cursor = Some(next), None => break,