1use crate::signals::is_positive_signal;
6use crate::{Engine, EngineResult, RetrieveInput, RetrieveOutput};
7use hippmem_core::hash::stable_hash64;
8use hippmem_core::ids::MemoryId;
9use hippmem_core::model::links::{ActivationStep, RecallChannel, RetrievalResult};
10use hippmem_core::model::unit::{MemoryLifecycle, MemoryUnit};
11use hippmem_core::time::Clock;
12use hippmem_model::deterministic::extract::DeterministicExtractor;
13use hippmem_model::lang::active_locales;
14use hippmem_retrieval::explain::deduce_dimensions;
15use hippmem_retrieval::seeds::{multi_channel_seeds, rrf_fuse};
16use hippmem_retrieval::spreading::spread_multi_hop_fused;
17use hippmem_retrieval::warnings::check_warnings;
18use hippmem_store::activation_log::ActivationLogger;
19use hippmem_store::kv::InvertedIndex;
20use hippmem_store::semantic::vector_index::BinaryIndex;
21use hippmem_store::semantic::vector_index::VectorIndex;
22use std::collections::HashMap;
23
24impl Engine {
25 pub fn retrieve(&self, input: RetrieveInput) -> EngineResult<RetrieveOutput> {
27 let start = std::time::Instant::now();
28 let params = self.params.read();
29
30 let extractor = DeterministicExtractor;
32 let query_content = hippmem_core::model::unit::MemoryContent {
33 raw: input.query.clone(),
34 summary: None,
35 normalized: None,
36 language: hippmem_core::model::unit::Language::Zh,
37 content_type: hippmem_core::model::enums::ContentType::UserStatement,
38 };
39 let understanding = extractor
40 .extract_sync_immediate(&query_content)
41 .unwrap_or_else(|_| hippmem_model::traits::ImmediateExtraction {
42 entities: vec![],
43 topics: vec![],
44 explicit_causals: vec![],
45 language: hippmem_core::model::unit::Language::Zh,
46 content_type: None,
47 importance: hippmem_core::score::UnitScore::new(0.0),
48 });
49
50 let inverted = InvertedIndex::new(self.store.db_arc());
52
53 let entity_hits: Vec<(MemoryId, f32)> = understanding
55 .entities
56 .iter()
57 .filter_map(|em| {
58 let key = hippmem_core::hash::stable_hash64(&em.canonical);
59 inverted.get_entity(&key).ok().map(|ids| {
60 ids.into_iter()
61 .map(|id| (MemoryId(id), 0.2f32))
62 .collect::<Vec<_>>()
63 })
64 })
65 .flatten()
66 .collect();
67
68 let topic_hits: Vec<(MemoryId, f32)> = understanding
70 .topics
71 .iter()
72 .filter_map(|t| {
73 let key = hippmem_core::hash::stable_hash64(&t.label);
74 inverted.get_topic(&key).ok().map(|ids| {
75 ids.into_iter()
76 .map(|id| (MemoryId(id), 0.15f32))
77 .collect::<Vec<_>>()
78 })
79 })
80 .flatten()
81 .collect();
82
83 let now = hippmem_core::time::SystemClock.now();
85 let temporal_keys = temporal_bucket_keys(now);
86 let mut temporal_hit_ids = std::collections::HashSet::new();
87 for tk in &temporal_keys {
88 if let Ok(ids) = inverted.get_temporal(tk) {
89 for id in ids {
90 temporal_hit_ids.insert(MemoryId(id));
91 }
92 }
93 }
94 let temporal_hits: Vec<(MemoryId, bool)> =
95 temporal_hit_ids.into_iter().map(|id| (id, true)).collect();
96
97 let bm25_hits: Vec<(MemoryId, f32)> = self
99 .fulltext_index
100 .lock()
101 .search(&input.query, params.seed_per_channel as usize)
102 .unwrap_or_default()
103 .into_iter()
104 .map(|(id, score)| {
105 let norm = (score / params.bm25_norm_factor).tanh();
106 (MemoryId(id), norm)
107 })
108 .collect();
109
110 let semantic_hits: Vec<(MemoryId, f32)> = {
112 let query_texts = vec![input.query.clone()];
113 self.embedder
114 .embed_sync(&query_texts)
115 .ok()
116 .and_then(|vectors| vectors.first().cloned())
117 .map(|query_vec| {
118 let idx = self.dense_vector_index.lock();
119 idx.search(&query_vec, params.seed_per_channel as usize)
120 .unwrap_or_default()
121 .into_iter()
122 .map(|(id, l2_dist)| {
123 let cos_sim = 1.0 / (1.0 + l2_dist);
125 (MemoryId(id), cos_sim)
126 })
127 .filter(|(_, sim)| *sim > 0.0)
128 .collect()
129 })
130 .unwrap_or_default()
131 };
132
133 let binary_hits: Vec<(MemoryId, f32)> = {
135 let query_bc = query_binary_code(&input.query);
136 let idx = self.binary_code_index.lock();
137 idx.search(&query_bc, params.seed_per_channel as usize)
138 .unwrap_or_default()
139 .into_iter()
140 .map(|(id, hamming)| {
141 let sim = 1.0 - (hamming as f32 / 128.0);
142 (MemoryId(id), sim.max(0.0))
143 })
144 .filter(|(_, sim)| *sim > 0.0)
145 .collect()
146 };
147
148 let query_goals = extract_query_goals(&input.query);
150 let goal_hits: Vec<(MemoryId, usize)> = query_goals
151 .iter()
152 .filter_map(|goal| {
153 let key = stable_hash64(goal);
154 inverted.get_goal(&key).ok().map(|ids| {
155 ids.into_iter()
156 .map(|id| (MemoryId(id), 1))
157 .collect::<Vec<_>>()
158 })
159 })
160 .flatten()
161 .collect();
162
163 let query_events = extract_query_events(&input.query);
165 let event_hits: Vec<(MemoryId, usize)> = query_events
166 .iter()
167 .filter_map(|event| {
168 let key = stable_hash64(event);
169 inverted.get_event(&key).ok().map(|ids| {
170 ids.into_iter()
171 .map(|id| (MemoryId(id), 1))
172 .collect::<Vec<_>>()
173 })
174 })
175 .flatten()
176 .collect();
177
178 let causal_hits: Vec<(MemoryId, usize)> = understanding
180 .explicit_causals
181 .iter()
182 .filter_map(|c| {
183 let causal_str = format!("{} -> {}", c.cause, c.effect);
184 let key = stable_hash64(&causal_str);
185 inverted.get_causal(&key).ok().map(|ids| {
186 ids.into_iter()
187 .map(|id| (MemoryId(id), 1))
188 .collect::<Vec<_>>()
189 })
190 })
191 .flatten()
192 .collect();
193
194 let recent_hits: Vec<(MemoryId, f32)> = {
196 let mut recent_map: HashMap<MemoryId, f32> = HashMap::new();
197
198 for mid in &input.context.recent_memory_ids {
200 recent_map
201 .entry(*mid)
202 .and_modify(|s| *s = (*s + 0.3).min(1.0))
203 .or_insert(0.3);
204 }
205
206 let graph = hippmem_store::graph::GraphStore::new(self.store.db_arc());
208 for mid in &input.context.recent_memory_ids {
209 if let Ok(links) = graph.get_outgoing(mid) {
210 for link in links.iter().take(8) {
211 recent_map
212 .entry(link.target_id)
213 .and_modify(|s| *s = (*s + 0.15).min(1.0))
214 .or_insert(0.15);
215 }
216 }
217 }
218
219 let act_log = ActivationLogger::new(self.store.db_arc());
221 if let Ok(records) = act_log.read_all() {
222 let mut freq: HashMap<MemoryId, u32> = HashMap::new();
223 for rec in records.iter() {
224 if !is_positive_signal(&rec.signal) {
226 continue;
227 }
228 for &mid in &rec.used_memory_ids {
229 *freq.entry(MemoryId(mid)).or_default() += 1;
231 }
232 }
233 let max_freq = freq.values().max().copied().unwrap_or(1) as f32;
234 for (mid, count) in freq {
235 let score = (count as f32 / max_freq) * 0.25;
236 recent_map
237 .entry(mid)
238 .and_modify(|s| *s = (*s + score).min(1.0))
239 .or_insert(score);
240 }
241 }
242
243 let mut hits: Vec<(MemoryId, f32)> = recent_map.into_iter().collect();
244 hits.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
245 hits.truncate(params.seed_per_channel as usize);
246 hits
247 };
248
249 let seed_result = multi_channel_seeds(
250 &input.query,
251 &entity_hits,
252 &temporal_hits,
253 &semantic_hits,
254 &topic_hits,
255 &bm25_hits,
256 &binary_hits,
257 &goal_hits,
258 &event_hits,
259 &causal_hits,
260 &recent_hits,
261 params.seed_per_channel as usize,
262 );
263
264 let fused_scores: HashMap<MemoryId, (f32, RecallChannel)> = if seed_result.seeds.is_empty()
266 {
267 let fallback = load_limited_units(self.store.db_arc(), 50)
270 .into_iter()
271 .filter(|u| !matches!(u.lifecycle, MemoryLifecycle::Compressed { .. }))
272 .collect::<Vec<_>>();
273 fallback
274 .into_iter()
275 .map(|u| (u.id, (0.3_f32, RecallChannel::RecentActivation)))
276 .collect()
277 } else {
278 rrf_fuse(&seed_result.seeds, ¶ms)
279 };
280
281 let seed_ids: Vec<MemoryId> = fused_scores.keys().cloned().collect();
283 let mut unit_map: HashMap<MemoryId, MemoryUnit> = HashMap::new();
284 for unit in load_units_by_ids(self.store.db_arc(), &seed_ids) {
285 unit_map.insert(unit.id, unit);
286 }
287
288 let importance_map: HashMap<MemoryId, f32> = unit_map
290 .iter()
291 .map(|(id, unit)| (*id, unit.understanding.importance.value()))
292 .collect();
293 let usage_map: HashMap<MemoryId, f32> = unit_map
295 .iter()
296 .map(|(id, unit)| (*id, unit.activation.usage_score.value()))
297 .collect();
298
299 let graph = hippmem_store::graph::GraphStore::new(self.store.db_arc());
300 let mut links_map: HashMap<MemoryId, Vec<hippmem_core::model::links::AssociationLink>> =
301 HashMap::new();
302
303 for sid in &seed_ids {
305 if let Ok(links) = graph.get_outgoing(sid) {
306 links_map.insert(*sid, links);
307 }
308 }
309
310 let neighbor_ids: Vec<MemoryId> = links_map
312 .values()
313 .flatten()
314 .map(|l| l.target_id)
315 .filter(|tid| !links_map.contains_key(tid))
316 .collect();
317 for nid in &neighbor_ids {
318 if let Ok(links) = graph.get_outgoing(nid) {
319 links_map.insert(*nid, links);
320 }
321 }
322 for unit in load_units_by_ids(self.store.db_arc(), &neighbor_ids) {
324 unit_map.entry(unit.id).or_insert(unit);
325 }
326
327 let (activated, merged_count) = spread_multi_hop_fused(
329 &fused_scores,
330 &links_map,
331 ¶ms,
332 &importance_map,
333 &usage_map,
334 input.max_hops.map(|h| h as u32),
335 );
336 let max_k = input.top_k.min(activated.len());
337
338 let extra_ids: Vec<MemoryId> = activated
340 .iter()
341 .map(|(id, _, _)| *id)
342 .filter(|id| !unit_map.contains_key(id))
343 .collect();
344 for unit in load_units_by_ids(self.store.db_arc(), &extra_ids) {
345 unit_map.insert(unit.id, unit);
346 }
347
348 let loaded_units: Vec<MemoryUnit> = activated
350 .iter()
351 .filter_map(|(id, _, _)| unit_map.get(id).cloned())
352 .collect();
353 let mut reranked = hippmem_retrieval::rerank::rerank_by_energy(&activated, &loaded_units);
354
355 apply_question_aware_boost(&input.query, &mut reranked, ¶ms);
358
359 reranked.retain(|(_, _, _, unit)| {
361 !matches!(unit.lifecycle, MemoryLifecycle::Compressed { .. })
362 });
363
364 let results: Vec<RetrievalResult> = reranked
366 .iter()
367 .take(max_k)
368 .map(|(_id, energy, trace, unit)| {
369 let matched = deduce_dimensions(trace);
370 let warns = check_warnings(unit, *energy);
371 RetrievalResult {
372 memory: unit.clone(),
373 final_score: *energy,
374 activation_trace: trace.clone(),
375 matched_dimensions: matched,
376 warnings: warns,
377 }
378 })
379 .collect();
380
381 let channel_contributions: Vec<(RecallChannel, u32)> = {
383 let mut map: HashMap<RecallChannel, u32> = HashMap::new();
384 for seed in &seed_result.seeds {
385 *map.entry(seed.channel).or_default() += 1;
386 }
387 map.into_iter().collect()
388 };
389
390 let retrieval_id = {
393 let act_log = ActivationLogger::new(self.store.db_arc());
394 let used_ids: Vec<u128> = results.iter().map(|r| r.memory.id.0).collect();
396 let now_ms =
397 if let Ok(t) = std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH) {
398 t.as_millis() as i64
399 } else {
400 0
401 };
402 let _ = act_log.record(&hippmem_store::activation_log::ActivationRecord {
403 retrieval_id: now_ms as u64,
404 used_memory_ids: used_ids,
405 signal: "retrieve".into(),
406 recorded_at_ms: now_ms,
407 });
408 now_ms as u64
409 };
410
411 Ok(RetrieveOutput {
412 retrieval_id,
413 results,
414 trace: crate::RetrievalTrace {
415 seeds: seed_result
416 .seeds
417 .iter()
418 .map(|s| crate::SeedRecord {
419 id: s.id,
420 channel: s.channel,
421 initial_energy: s.score,
422 rank_in_channel: s.rank_in_channel,
423 })
424 .collect(),
425 steps: activated
426 .iter()
427 .flat_map(|(_, _, trace)| trace.clone())
428 .collect(),
429 hops_used: activated
430 .iter()
431 .flat_map(|(_, _, trace)| trace.iter())
432 .map(|s| s.hop)
433 .max()
434 .unwrap_or(0),
435 merged_count,
436 },
437 diagnostics: crate::RetrievalDiagnostics {
438 channel_contributions,
439 reranked: true,
440 pruned_branches: 0,
441 backend_used: crate::BackendUsage {
442 embedder: self.embedder.backend_id().to_string(),
443 reranker: Some("rule".into()),
444 },
445 latency_ms: start.elapsed().as_millis() as u32,
446 },
447 })
448 }
449}
450
451#[derive(Debug, Clone, Copy, PartialEq)]
457enum QuestionType {
458 Why,
460 How,
462 What,
464 Correction,
466 Preference,
468 None,
470}
471
472fn detect_question_type(query: &str) -> QuestionType {
478 let q = query.to_lowercase();
479
480 for lang in active_locales() {
482 if let Some((before, after)) = lang.change_pair {
483 if q.contains(before) && q.contains(after) {
484 return QuestionType::Correction;
485 }
486 }
487 }
488
489 for lang in active_locales() {
493 for keyword in lang.q_correction {
494 if q.contains(keyword) {
495 return QuestionType::Correction;
496 }
497 }
498 }
499 for lang in active_locales() {
500 for keyword in lang.q_preference {
501 if q.contains(keyword) {
502 return QuestionType::Preference;
503 }
504 }
505 }
506 for lang in active_locales() {
507 for keyword in lang.q_why {
508 if q.contains(keyword) {
509 return QuestionType::Why;
510 }
511 }
512 }
513 for lang in active_locales() {
514 for keyword in lang.q_how {
515 if q.contains(keyword) {
516 return QuestionType::How;
517 }
518 }
519 }
520 for lang in active_locales() {
521 for keyword in lang.q_what {
522 if q.contains(keyword) {
523 return QuestionType::What;
524 }
525 }
526 }
527 QuestionType::None
528}
529
530fn explanatory_pattern_score(text: &str) -> f32 {
532 let mut score = 0.0f32;
533 for lang in active_locales() {
534 for (pattern, boost) in lang.explanatory {
535 if text.contains(pattern) {
536 score += boost;
537 }
538 }
539 }
540 score.min(0.20) }
542
543fn content_type_boost(query: &str) -> Vec<(hippmem_core::model::unit::ContentType, f32)> {
553 let qt = detect_question_type(query);
554 let mut boosts = Vec::new();
555
556 match qt {
557 QuestionType::Correction => {
558 boosts.push((hippmem_core::model::unit::ContentType::Correction, 0.12));
560 }
561 QuestionType::Preference => {
562 boosts.push((hippmem_core::model::unit::ContentType::Preference, 0.08));
564 boosts.push((hippmem_core::model::unit::ContentType::Decision, 0.04));
566 }
567 QuestionType::Why => {
568 boosts.push((hippmem_core::model::unit::ContentType::Decision, 0.08));
570 boosts.push((hippmem_core::model::unit::ContentType::TaskState, 0.08));
571 }
572 QuestionType::How => {
573 boosts.push((hippmem_core::model::unit::ContentType::TaskState, 0.08));
575 }
576 QuestionType::What => {
577 boosts.push((
583 hippmem_core::model::unit::ContentType::ProjectKnowledge,
584 0.15,
585 ));
586 }
587 QuestionType::None => {
588 }
590 }
591
592 if qt != QuestionType::Correction {
594 let q = query.to_lowercase();
595 let has_correction_signal = active_locales().iter().any(|lang| {
596 lang.q_correction.iter().any(|kw| q.contains(kw))
597 || lang
598 .change_pair
599 .is_some_and(|(b, a)| q.contains(b) && q.contains(a))
600 });
601 if has_correction_signal {
602 boosts.push((hippmem_core::model::unit::ContentType::Correction, 0.10));
603 }
604 }
605
606 boosts
607}
608
609fn apply_question_aware_boost(
619 query: &str,
620 reranked: &mut [(MemoryId, f32, Vec<ActivationStep>, MemoryUnit)],
621 params: &hippmem_core::config::AlgoParams,
622) {
623 let qt = detect_question_type(query);
624 let ct_boosts = content_type_boost(query);
625 let cap = params.seed_energy_cap;
626 let what_subject: Option<String> = if qt == QuestionType::What {
628 extract_subject_for_what_query(query)
629 } else {
630 None
631 };
632
633 match qt {
635 QuestionType::Why => {
636 for (_, energy, _, unit) in reranked.iter_mut() {
637 let boost = explanatory_pattern_score(&unit.content.raw);
638 if boost > 0.0 {
639 *energy = (*energy + boost).min(cap);
640 }
641 }
642 }
643 QuestionType::Correction
644 | QuestionType::Preference
645 | QuestionType::How
646 | QuestionType::What
647 | QuestionType::None => {
648 }
650 }
651
652 if !ct_boosts.is_empty() {
656 for (_, energy, _, unit) in reranked.iter_mut() {
657 for (ct, boost) in &ct_boosts {
658 if unit.content.content_type != *ct {
659 continue;
660 }
661 if qt == QuestionType::What
663 && *ct == hippmem_core::model::unit::ContentType::ProjectKnowledge
664 {
665 if let Some(ref subject) = what_subject {
666 let content_lower = unit.content.raw.to_lowercase();
667 if !content_lower.contains(&subject.to_lowercase()) {
668 break; }
670 }
671 }
672 *energy = (*energy + boost).min(cap);
673 break; }
675 }
676 }
677
678 let keywords = extract_discriminative_keywords(query);
684 if !keywords.is_empty() {
685 for (_, energy, _, unit) in reranked.iter_mut() {
686 let mut kw_bonus = 0.0f32;
687 let content_lower = unit.content.raw.to_lowercase();
688 for kw in &keywords {
689 if content_lower.contains(&kw.to_lowercase()) {
690 kw_bonus += 0.04;
691 }
692 }
693 if kw_bonus > 0.0 {
694 *energy = (*energy + kw_bonus.min(0.08)).min(cap);
695 }
696 }
697 }
698
699 if qt == QuestionType::What {
704 if let Some(ref subject) = extract_subject_for_what_query(query) {
705 let subject_lower = subject.to_lowercase();
706 for (_, energy, _, unit) in reranked.iter_mut() {
707 let content_lower = unit.content.raw.to_lowercase();
708 let has_definition = active_locales().iter().any(|lang| {
709 lang.definition_patterns
710 .iter()
711 .any(|pat| content_lower.contains(&format!("{} {pat}", subject_lower)))
712 });
713 if has_definition {
714 *energy = (*energy + 0.05).min(cap);
715 }
716 }
717 }
718 }
719
720 reranked.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
722}
723
724fn extract_subject_for_what_query(query: &str) -> Option<String> {
733 let q = query.to_lowercase();
734 for lang in active_locales() {
735 for delimiter in lang.what_delimiters {
736 if let Some(pos) = q.find(delimiter) {
737 let prefix = &q[..pos];
738 let subject = if let Some(particle) = lang.possessive_particle {
739 prefix
742 .rsplit(particle)
743 .next()
744 .unwrap_or("")
745 .rsplit(|c: char| c.is_whitespace() || c == '?' || c == '?')
746 .next()
747 .unwrap_or("")
748 .trim()
749 .to_string()
750 } else {
751 prefix
752 .rsplit(|c: char| c.is_whitespace() || c == '?' || c == '?')
753 .next()
754 .unwrap_or("")
755 .trim()
756 .to_string()
757 };
758 if subject.len() >= 2 {
759 return Some(subject);
760 }
761 return None;
762 }
763 }
764 }
765 None
766}
767
768fn extract_discriminative_keywords(query: &str) -> Vec<String> {
776 let stop_words: Vec<&str> = active_locales()
778 .iter()
779 .flat_map(|lang| lang.stop_words.iter().copied())
780 .collect();
781
782 let mut keywords: Vec<String> = Vec::new();
783 let mut seen = std::collections::HashSet::new();
784
785 for word in query.split(|c: char| !c.is_alphanumeric()) {
787 let is_keyword = (word.len() >= 2 && word.chars().any(|c| c.is_uppercase()))
788 || (word.chars().all(|c| c.is_ascii_alphabetic()) && word.len() >= 3);
789 if is_keyword
790 && !stop_words.contains(&word.to_lowercase().as_str())
791 && seen.insert(word.to_string())
792 {
793 keywords.push(word.to_string());
794 }
795 }
796
797 for word in query
799 .split(|c: char| c.is_whitespace() || c.is_ascii_punctuation() || c == '?' || c == '?')
800 {
801 let trimmed = word.trim();
802 if trimmed.chars().count() >= 2
803 && trimmed.chars().all(|c| c as u32 > 0x2E80) && !stop_words.contains(&trimmed)
805 && seen.insert(trimmed.to_string())
806 {
807 keywords.push(trimmed.to_string());
808 }
809 }
810
811 keywords.truncate(5); keywords
813}
814
815fn query_binary_code(text: &str) -> [u8; 16] {
817 let bc0 = stable_hash64(&format!("bc_0_{}", text));
818 let bc1 = stable_hash64(&format!("bc_1_{}", text));
819 let mut bytes = [0u8; 16];
820 bytes[..8].copy_from_slice(&bc0.to_le_bytes());
821 bytes[8..].copy_from_slice(&bc1.to_le_bytes());
822 bytes
823}
824
825fn temporal_bucket_keys(ts: hippmem_core::time::Timestamp) -> Vec<u32> {
827 let ms = ts.0;
828 vec![
829 (ms / 3_600_000) as u32, (ms / 86_400_000) as u32, (ms / 604_800_000) as u32, ]
833}
834
835pub(crate) fn load_all_units(db: std::sync::Arc<redb::Database>) -> Vec<MemoryUnit> {
836 use redb::ReadableDatabase;
837 use redb::ReadableTable;
838 let mut units = Vec::new();
839 let read_txn = db.begin_read().expect("read transaction should succeed");
840 let table = read_txn
841 .open_table(hippmem_store::store::MEMORY_KV)
842 .expect("memory_kv table should exist");
843 let iter = table.iter().expect("iter should succeed");
844 for entry in iter.flatten() {
845 let (_key, value) = entry;
846 if let Ok((unit, _)) = bincode::serde::decode_from_slice::<MemoryUnit, _>(
847 value.value(),
848 bincode::config::standard(),
849 ) {
850 units.push(unit);
851 }
852 }
853 units
854}
855
856fn load_units_by_ids(db: std::sync::Arc<redb::Database>, ids: &[MemoryId]) -> Vec<MemoryUnit> {
858 if ids.is_empty() {
859 return vec![];
860 }
861 use redb::ReadableDatabase;
862 let mut units = Vec::new();
863 let read_txn = db.begin_read().expect("read transaction should succeed");
864 let table = read_txn
865 .open_table(hippmem_store::store::MEMORY_KV)
866 .expect("memory_kv table should exist");
867 for id in ids {
868 if let Some(value) = table.get(id.0).expect("get should succeed") {
869 if let Ok((unit, _)) = bincode::serde::decode_from_slice::<MemoryUnit, _>(
870 value.value(),
871 bincode::config::standard(),
872 ) {
873 units.push(unit);
874 }
875 }
876 }
877 units
878}
879
880fn extract_query_goals(text: &str) -> Vec<String> {
882 let mut goals = Vec::new();
883 for lang in active_locales() {
884 for m in lang.goal_markers {
885 if text.contains(m) {
886 goals.push(format!("goal_marker:{m}"));
887 }
888 }
889 }
890 goals
891}
892
893fn extract_query_events(text: &str) -> Vec<String> {
895 let mut events = Vec::new();
896 for lang in active_locales() {
897 for m in lang.event_markers {
898 if text.contains(m) {
899 events.push(format!("event_marker:{m}"));
900 }
901 }
902 }
903 events
904}
905
906fn load_limited_units(db: std::sync::Arc<redb::Database>, limit: usize) -> Vec<MemoryUnit> {
908 use redb::ReadableDatabase;
909 use redb::ReadableTable;
910 let mut units = Vec::new();
911 let read_txn = db.begin_read().expect("read transaction should succeed");
912 let table = read_txn
913 .open_table(hippmem_store::store::MEMORY_KV)
914 .expect("memory_kv table should exist");
915 let iter = table.iter().expect("iter should succeed");
916 for entry in iter.flatten().take(limit) {
917 let (_key, value) = entry;
918 if let Ok((unit, _)) = bincode::serde::decode_from_slice::<MemoryUnit, _>(
919 value.value(),
920 bincode::config::standard(),
921 ) {
922 units.push(unit);
923 }
924 }
925 units
926}