1use alloc::format;
17use alloc::string::{String, ToString};
18use alloc::vec::Vec;
19
20use plugmem_arena::{
21 Arena, ArenaCfg, BlobHeap, BlobHeapCfg, BlobId, ChunkPool, ChunkPoolCfg, Interner, ListHandle,
22 ShardMode, Slot, TermId, key,
23};
24
25use crate::config::Config;
26use crate::error::Error;
27use crate::id::{EdgeId, EntityId, FactId, NONE_U32};
28use crate::index::IdListIndex;
29use crate::index::bm25::Bm25Index;
30use crate::index::hnsw::HnswGraph;
31use crate::index::vecpool::VecPool;
32use crate::journal::{JournalScan, Op, scan};
33use crate::model::{
34 EdgeHistorySlot, EdgeSlot, EntityByName, EntityRecord, FactAux, FactRecord, TemporalSlot,
35 VALID_TO_OPEN, close_edge_history_payload, edge_history_key, edge_key, fact_flags,
36};
37use crate::storage::Storage;
38use crate::tokenizer::Tokenizer;
39
40use maintain::TOKENIZER_INDEX_VERSION;
41
42const SIMILAR_CANDIDATE_CAP: usize = 32;
44
45mod maintain;
46mod migrations;
47mod persist;
48mod recall;
49mod reembed;
50mod shards;
51mod tags;
52
53pub use maintain::{MaintainReport, MaintenanceMode, MaintenanceOptions};
54pub use recall::{RecallQuery, RecallResult, RecallScratch, RecalledEdge, RecalledFact, source};
55pub use reembed::{ReembedError, ReembedReport};
56pub use shards::ShardLayout;
57pub use tags::{DEFAULT_TAG_PAGE_LIMIT, MAX_TAG_PAGE_LIMIT, TagPage, TagQuery, TagSummary};
58
59pub const MAX_VECTOR_SPACE_ID_BYTES: usize = 256;
61
62#[derive(Clone, Copy, Debug)]
64#[cfg_attr(feature = "serde", derive(serde::Serialize))]
65pub struct RememberInput<'a> {
66 pub now: u64,
68 pub text: &'a str,
70 pub entity: Option<&'a str>,
72 pub tags: &'a [&'a str],
74 pub links: &'a [(&'a str, &'a str)],
76 pub vector: Option<&'a [f32]>,
80 pub valid_from: Option<u64>,
82 pub metadata: Option<&'a [(&'a str, &'a str)]>,
87}
88
89impl<'a> RememberInput<'a> {
90 pub fn text(now: u64, text: &'a str) -> Self {
92 Self {
93 now,
94 text,
95 entity: None,
96 tags: &[],
97 links: &[],
98 vector: None,
99 valid_from: None,
100 metadata: None,
101 }
102 }
103}
104
105#[derive(Clone, Copy, Debug)]
107#[cfg_attr(feature = "serde", derive(serde::Serialize))]
108pub struct LinkInput<'a> {
109 pub now: u64,
111 pub src: &'a str,
113 pub rel: &'a str,
115 pub dst: &'a str,
117 pub provenance: Option<FactId>,
119}
120
121#[derive(Clone, Copy, Debug)]
123#[cfg_attr(feature = "serde", derive(serde::Serialize))]
124pub struct UnlinkInput<'a> {
125 pub now: u64,
127 pub src: &'a str,
129 pub rel: &'a str,
131 pub dst: &'a str,
133}
134
135#[derive(Clone, Debug, PartialEq)]
137#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
138pub struct RememberOutcome {
139 pub id: FactId,
141 pub entity: Option<EntityId>,
143 pub similar: Vec<Similar>,
148}
149
150#[derive(Clone, Debug, PartialEq)]
157#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
158#[cfg_attr(feature = "serde", serde(tag = "status", rename_all = "snake_case"))]
159pub enum GuardedRememberOutcome {
160 Stored {
163 outcome: RememberOutcome,
165 },
166 Blocked {
169 similar: Vec<Similar>,
172 },
173}
174
175#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
177#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
178pub struct RemoveTagReport {
179 pub affected: u32,
182}
183
184#[derive(Clone, Copy, Debug, PartialEq)]
186#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
187pub struct Similar {
188 pub id: FactId,
190 pub score: f32,
193 pub reason: SimilarReason,
195}
196
197#[derive(Clone, Copy, Debug, PartialEq, Eq)]
199#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
200pub enum SimilarReason {
201 LexicalOverlap,
204 VectorCosine,
207}
208
209#[derive(Debug, Default)]
213struct SimilarityScratch {
214 new_terms: Vec<u32>,
216 unknown: String,
219 unknown_ranges: Vec<(usize, usize)>,
221 candidate_terms: Vec<u32>,
223 vector: Vec<u8>,
225}
226
227#[derive(Clone, Copy)]
228enum SimilarVector<'s> {
229 None,
230 Stored(u32),
231 Encoded(&'s [u8]),
232}
233
234#[derive(Clone, Copy)]
235struct SimilarityQuery<'s> {
236 entity: EntityId,
237 exclude: Option<FactId>,
238 terms: &'s [u32],
239 term_count: usize,
240 vector: SimilarVector<'s>,
241}
242
243#[derive(Clone, Copy, Debug)]
245#[cfg_attr(feature = "serde", derive(serde::Serialize))]
246pub struct FactView<'a> {
247 pub record: FactRecord,
249 pub text: &'a str,
251}
252
253#[derive(Clone, Copy, Debug, PartialEq, Eq)]
258#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
259pub enum FactFault {
260 Text,
262 Vector,
265 Metadata,
268}
269
270#[derive(Clone, Copy, Debug, PartialEq, Eq)]
274#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
275#[non_exhaustive]
276pub struct Stats {
277 pub facts: usize,
280 pub entities: usize,
282 pub terms: usize,
284 pub edges: usize,
287 pub edge_versions: usize,
289 pub vectors: usize,
291 pub tombstones: usize,
293 pub hnsw_indexed: u32,
296 pub next_fact: u32,
299 pub next_entity: u32,
301 pub next_edge: u32,
303 pub db_uuid: u128,
306 pub pool_bytes: usize,
309 pub shards: ShardLayout,
315}
316
317#[derive(Clone, Debug, Default, PartialEq, Eq)]
319#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
320pub struct OpenReport {
321 pub replayed: usize,
323 pub skipped: usize,
325 pub truncated_tail: bool,
327}
328
329pub struct Memory<'a> {
336 cfg: Config,
337 facts: Arena<'a, FactRecord>,
339 fact_aux: Arena<'a, FactAux>,
340 entities: Arena<'a, EntityRecord>,
341 by_name: Arena<'a, EntityByName>,
342 edges_out: Arena<'a, EdgeSlot>,
343 edges_in: Arena<'a, EdgeSlot>,
344 edges_hist_out: Arena<'a, EdgeHistorySlot>,
345 edges_hist_in: Arena<'a, EdgeHistorySlot>,
346 temporal: Arena<'a, TemporalSlot>,
347 texts: BlobHeap<'a>,
349 metas: BlobHeap<'a>,
354 terms: Interner<'a>,
357 tag_lists: ChunkPool<'a>,
359 bm25: Bm25Index<'a>,
361 tags_idx: IdListIndex<'a>,
362 tag_catalog: tags::TagCatalog,
364 entity_facts: IdListIndex<'a>,
365 vecs: VecPool<'a>,
367 hnsw: HnswGraph<'a>,
371 vector_space: Option<String>,
375 next_fact: u32,
377 next_entity: u32,
378 next_edge: u32,
379 tombstones: usize,
381 bm25_tokenizer_version: u32,
382 tokenizer: Tokenizer,
384 tf_scratch: Vec<(u32, u8)>,
385 name_scratch: String,
386 similarity_scratch: SimilarityScratch,
387}
388
389impl<'a> Memory<'a> {
390 pub fn new(cfg: Config) -> Result<Self, Error> {
397 cfg.validate()?;
398 let uni =
399 |shards: usize| ArenaCfg::new(shards, ShardMode::Uniform).with_max_bytes(cfg.max_bytes);
400 let ord =
401 |shards: usize| ArenaCfg::new(shards, ShardMode::Ordered).with_max_bytes(cfg.max_bytes);
402 let blob = BlobHeapCfg::new()
403 .with_max_bytes(cfg.max_bytes)
404 .with_max_blob(cfg.max_blob);
405 Ok(Self {
406 facts: Arena::new(uni(cfg.shards_facts))?,
407 fact_aux: Arena::new(uni(cfg.shards_facts))?,
408 entities: Arena::new(uni(cfg.shards_entities))?,
409 by_name: Arena::new(ord(cfg.shards_entities))?,
410 edges_out: Arena::new(ord(cfg.shards_edges))?,
411 edges_in: Arena::new(ord(cfg.shards_edges))?,
412 edges_hist_out: Arena::new(ord(cfg.shards_edges))?,
413 edges_hist_in: Arena::new(ord(cfg.shards_edges))?,
414 temporal: Arena::new(ord(cfg.shards_temporal))?,
415 texts: BlobHeap::new(blob),
416 metas: BlobHeap::new(blob),
417 terms: Interner::new(blob),
418 tag_lists: ChunkPool::new(ChunkPoolCfg::new().with_max_bytes(cfg.max_bytes)),
419 bm25: Bm25Index::new(cfg.shards_postings, cfg.max_bytes)?,
420 tags_idx: IdListIndex::new(cfg.shards_postings, cfg.max_bytes)?,
421 tag_catalog: tags::TagCatalog::new(),
422 entity_facts: IdListIndex::new(cfg.shards_entities, cfg.max_bytes)?,
423 vecs: VecPool::new(cfg.dim, cfg.max_bytes),
424 hnsw: HnswGraph::new(cfg.hnsw_m, cfg.hnsw_m0, cfg.max_bytes)?,
425 vector_space: None,
426 next_fact: 0,
427 next_entity: 0,
428 next_edge: 0,
429 tombstones: 0,
430 bm25_tokenizer_version: maintain::TOKENIZER_INDEX_VERSION,
431 tokenizer: Tokenizer::new(),
432 tf_scratch: Vec::new(),
433 name_scratch: String::new(),
434 similarity_scratch: SimilarityScratch::default(),
435 cfg,
436 })
437 }
438
439 pub fn open<S: Storage>(store: &mut S, cfg: Config) -> Result<(Self, OpenReport), Error> {
443 let snapshot = store
444 .read_snapshot()
445 .map_err(|e| Error::Storage(format!("{e:?}")))?;
446 let journal = store
447 .read_journal()
448 .map_err(|e| Error::Storage(format!("{e:?}")))?;
449 Self::from_bytes(snapshot.as_deref(), &journal, cfg)
450 }
451
452 pub fn from_bytes(
454 snapshot: Option<&[u8]>,
455 journal: &[u8],
456 cfg: Config,
457 ) -> Result<(Self, OpenReport), Error> {
458 let mut mem = match snapshot {
459 Some(bytes) => Self::load_snapshot(bytes, cfg)?,
460 None => Self::new(cfg)?,
461 };
462 let report = mem.replay(journal)?;
463 Ok((mem, report))
464 }
465
466 pub fn from_bytes_borrowed(
477 snapshot: &'a [u8],
478 journal: &[u8],
479 cfg: Config,
480 ) -> Result<Self, Error> {
481 let mem = Self::load_snapshot_borrowed(snapshot, cfg)?;
482 let JournalScan { entries, .. } = scan(journal)?;
483 if !entries.is_empty() {
484 return Err(Error::Invalid(
485 "read-only open requires a checkpointed (empty) journal",
486 ));
487 }
488 Ok(mem)
489 }
490
491 pub fn from_bytes_overlay(
506 snapshot: &'a [u8],
507 journal: &[u8],
508 cfg: Config,
509 ) -> Result<(Self, OpenReport), Error> {
510 let mut mem = Self::load_snapshot_borrowed(snapshot, cfg)?;
511 let report = mem.replay(journal)?;
512 Ok((mem, report))
513 }
514
515 fn replay(&mut self, journal: &[u8]) -> Result<OpenReport, Error> {
520 let JournalScan {
521 entries,
522 truncated_tail,
523 } = scan(journal)?;
524 let mut report = OpenReport {
525 truncated_tail,
526 ..OpenReport::default()
527 };
528 for entry in entries {
529 let op = Op::decode(entry.op, entry.payload)?;
530 match op {
531 Op::Remember {
532 now,
533 valid_from,
534 entity,
535 text,
536 ref tags,
537 ref links,
538 ref vector,
539 ref metadata,
540 revises,
541 assigned,
542 } => {
543 if assigned.0 < self.next_fact {
544 report.skipped += 1;
545 continue;
546 }
547 if assigned.0 != self.next_fact {
548 return Err(Error::Corrupt("journal fact ids are not contiguous"));
549 }
550 if !vector.is_empty() && vector.len() != self.cfg.dim {
551 return Err(Error::Corrupt(
552 "journal vector dimension disagrees with dim",
553 ));
554 }
555 if let Some(target) = revises.some() {
556 self.check_revisable(target)
557 .map_err(|_| Error::Corrupt("journal revises an unrevisable fact"))?;
558 }
559 self.apply_remember(
560 &RememberInput {
561 now,
562 text,
563 entity,
564 tags: &tags.to_vec(),
565 links: &links.to_vec(),
566 vector: (!vector.is_empty()).then_some(vector.as_slice()),
567 valid_from: Some(valid_from),
568 metadata: (!metadata.is_empty()).then_some(metadata.as_slice()),
569 },
570 revises,
571 None,
572 )?;
573 if let Some(target) = revises.some() {
574 self.close_target(target, valid_from);
575 }
576 report.replayed += 1;
577 }
578 Op::Forget { fact, .. } => {
579 match self.apply_forget(fact) {
582 Ok(_) => report.replayed += 1,
583 Err(Error::NotFound(_)) => {
584 return Err(Error::Corrupt("journal forgets an unknown fact"));
585 }
586 Err(e) => return Err(e),
587 }
588 }
589 Op::Link {
590 now,
591 src,
592 rel,
593 dst,
594 provenance,
595 } => {
596 self.apply_link(now, src, rel, dst, provenance)?;
597 report.replayed += 1;
598 }
599 Op::Unlink { now, src, rel, dst } => {
600 self.apply_unlink(now, src, rel, dst)?;
601 report.replayed += 1;
602 }
603 Op::RemoveTag { now, tag } => {
604 self.apply_remove_tag(now, tag)?;
605 report.replayed += 1;
606 }
607 Op::SetVectorSpace { space } => {
608 self.apply_set_vector_space(space)?;
609 report.replayed += 1;
610 }
611 Op::Maintain {
612 mode,
613 max_hnsw_inserts,
614 ..
615 } => {
616 let options =
620 maintain::MaintenanceOptions::from_journal(mode, max_hnsw_inserts)?;
621 self.replay_maintain_with_options(options)?;
622 report.replayed += 1;
623 }
624 }
625 }
626 Ok(report)
627 }
628
629 pub fn vector_space(&self) -> Option<&str> {
634 self.vector_space.as_deref()
635 }
636
637 pub fn claim_vector_space<S: Storage>(
642 &mut self,
643 store: &mut S,
644 space: &str,
645 ) -> Result<bool, Error> {
646 let needs_claim = self.check_vector_space_claim(space)?;
647 if !needs_claim {
648 return Ok(false);
649 }
650 let mut entry = Vec::new();
651 Op::SetVectorSpace { space }.encode(&mut entry);
652 store
653 .append_journal(&entry)
654 .map_err(|e| Error::Storage(format!("{e:?}")))?;
655 self.vector_space = Some(space.into());
656 Ok(true)
657 }
658
659 fn check_vector_space_claim(&self, space: &str) -> Result<bool, Error> {
662 Self::validate_vector_space(space)?;
663 if let Some(stored) = &self.vector_space {
664 if stored == space {
665 return Ok(false);
666 }
667 if !self.vecs.is_empty() {
668 return Err(Error::VectorSpaceMismatch {
669 stored: stored.clone(),
670 requested: space.into(),
671 });
672 }
673 }
674 if !self.vecs.is_empty() {
675 return Err(Error::UntrackedVectorSpace);
676 }
677 Ok(true)
678 }
679
680 fn apply_set_vector_space(&mut self, space: &str) -> Result<(), Error> {
681 Self::validate_vector_space(space)?;
682 if let Some(stored) = &self.vector_space {
683 if stored == space {
684 return Ok(());
685 }
686 if !self.vecs.is_empty() {
687 return Err(Error::Corrupt(
688 "journal changes an established vector space",
689 ));
690 }
691 }
692 if !self.vecs.is_empty() {
693 return Err(Error::Corrupt(
694 "journal assigns a vector space after vector records",
695 ));
696 }
697 self.vector_space = Some(space.into());
698 Ok(())
699 }
700
701 pub(super) fn validate_vector_space(space: &str) -> Result<(), Error> {
702 if space.is_empty() {
703 return Err(Error::Invalid("vector space must not be empty"));
704 }
705 if space.len() > MAX_VECTOR_SPACE_ID_BYTES {
706 return Err(Error::TooLarge {
707 what: "vector space",
708 len: space.len(),
709 max: MAX_VECTOR_SPACE_ID_BYTES,
710 });
711 }
712 if space.bytes().any(|b| b < b' ' || b == 0x7f) {
713 return Err(Error::Invalid(
714 "vector space must not contain control bytes",
715 ));
716 }
717 Ok(())
718 }
719
720 pub fn remember<S: Storage>(
722 &mut self,
723 store: &mut S,
724 input: RememberInput<'_>,
725 ) -> Result<RememberOutcome, Error> {
726 self.validate_input(&input)?;
727 let mut outcome = self.apply_remember(&input, FactId::NONE, None)?;
728 self.find_similar(&mut outcome);
729 self.journal_remember(store, &input, FactId::NONE, outcome.id)?;
730 Ok(outcome)
731 }
732
733 pub fn remember_guarded<S: Storage>(
743 &mut self,
744 store: &mut S,
745 input: RememberInput<'_>,
746 ) -> Result<GuardedRememberOutcome, Error> {
747 self.remember_guarded_with_vector_space(store, input, None)
748 }
749
750 #[doc(hidden)]
755 pub fn remember_guarded_with_vector_space<S: Storage>(
756 &mut self,
757 store: &mut S,
758 input: RememberInput<'_>,
759 vector_space: Option<&str>,
760 ) -> Result<GuardedRememberOutcome, Error> {
761 self.validate_input(&input)?;
762 if let Some(space) = vector_space {
763 self.check_vector_space_claim(space)?;
764 }
765 let similar = self.find_similar_input(&input)?;
766 if !similar.is_empty() {
767 return Ok(GuardedRememberOutcome::Blocked { similar });
768 }
769 if let Some(space) = vector_space {
770 self.claim_vector_space(store, space)?;
771 }
772 let outcome = self.apply_remember(&input, FactId::NONE, None)?;
773 self.journal_remember(store, &input, FactId::NONE, outcome.id)?;
774 Ok(GuardedRememberOutcome::Stored { outcome })
775 }
776
777 pub fn remember_batch<S: Storage>(
785 &mut self,
786 store: &mut S,
787 inputs: &[RememberInput<'_>],
788 skip_similar: bool,
789 ) -> Result<Vec<RememberOutcome>, Error> {
790 let mut out = Vec::with_capacity(inputs.len());
791 for input in inputs {
792 self.validate_input(input)?;
793 let mut outcome = self.apply_remember(input, FactId::NONE, None)?;
794 if !skip_similar {
795 self.find_similar(&mut outcome);
796 }
797 self.journal_remember(store, input, FactId::NONE, outcome.id)?;
798 out.push(outcome);
799 }
800 Ok(out)
801 }
802
803 fn find_similar_input(&mut self, input: &RememberInput<'_>) -> Result<Vec<Similar>, Error> {
807 let Some(name) = input.entity else {
808 return Ok(Vec::new());
809 };
810 let Some(entity) = self.lookup_entity_name(name) else {
811 return Ok(Vec::new());
812 };
813
814 let mut scratch = core::mem::take(&mut self.similarity_scratch);
815 scratch.new_terms.clear();
816 scratch.unknown.clear();
817 scratch.unknown_ranges.clear();
818 scratch.candidate_terms.clear();
819 scratch.vector.clear();
820
821 let terms = &self.terms;
822 let new_terms = &mut scratch.new_terms;
823 let unknown = &mut scratch.unknown;
824 let unknown_ranges = &mut scratch.unknown_ranges;
825 self.tokenizer.tokenize(input.text, &mut |token| {
826 if let Some(term) = terms.lookup(token) {
827 if !new_terms.contains(&term.0) {
828 new_terms.push(term.0);
829 }
830 return;
831 }
832 if unknown_ranges
833 .iter()
834 .any(|&(start, end)| &unknown[start..end] == token)
835 {
836 return;
837 }
838 let start = unknown.len();
839 unknown.push_str(token);
840 unknown_ranges.push((start, unknown.len()));
841 });
842 let new_term_count = scratch.new_terms.len() + scratch.unknown_ranges.len();
843
844 let mut similar = Vec::new();
845 let result = (|| {
846 let encoded = match input.vector {
847 Some(vector) => {
848 self.vecs
849 .encode_slot_into(FactId::NONE, vector, &mut scratch.vector)?;
850 SimilarVector::Encoded(&scratch.vector)
851 }
852 None => SimilarVector::None,
853 };
854 self.scan_similar(
855 SimilarityQuery {
856 entity,
857 exclude: None,
858 terms: &scratch.new_terms,
859 term_count: new_term_count,
860 vector: encoded,
861 },
862 &mut scratch.candidate_terms,
863 &mut similar,
864 );
865 Ok::<(), Error>(())
866 })();
867 self.similarity_scratch = scratch;
868 result.map(|()| similar)
869 }
870
871 fn find_similar(&mut self, outcome: &mut RememberOutcome) {
875 let Some(entity) = outcome.entity else { return };
876 let new_vec = self
877 .fact(outcome.id)
878 .filter(|record| record.has_vector())
879 .map(|record| record.vector);
880 let mut scratch = core::mem::take(&mut self.similarity_scratch);
881 scratch.new_terms.clear();
882 scratch
883 .new_terms
884 .extend(self.tf_scratch.iter().map(|&(term, _)| term));
885 let new_term_count = scratch.new_terms.len();
886 self.scan_similar(
887 SimilarityQuery {
888 entity,
889 exclude: Some(outcome.id),
890 terms: &scratch.new_terms,
891 term_count: new_term_count,
892 vector: new_vec.map_or(SimilarVector::None, SimilarVector::Stored),
893 },
894 &mut scratch.candidate_terms,
895 &mut outcome.similar,
896 );
897 self.similarity_scratch = scratch;
898 }
899
900 fn scan_similar(
922 &mut self,
923 query: SimilarityQuery<'_>,
924 candidate_terms: &mut Vec<u32>,
925 out: &mut Vec<Similar>,
926 ) {
927 out.clear();
928 if query.term_count == 0 && matches!(query.vector, SimilarVector::None) {
929 return;
930 }
931 let mut ring = [FactId::NONE; SIMILAR_CANDIDATE_CAP];
933 let mut n = 0usize;
934 for (fact, _) in self.entity_facts.entries(query.entity.0) {
935 if query.exclude == Some(fact) {
936 continue;
937 }
938 ring[n % SIMILAR_CANDIDATE_CAP] = fact;
939 n += 1;
940 }
941 let summaries_trustworthy = self.bm25_tokenizer_version == TOKENIZER_INDEX_VERSION;
942 let lexical_only = matches!(query.vector, SimilarVector::None);
949 for &fact in ring.iter().take(n.min(SIMILAR_CANDIDATE_CAP)) {
950 let may_overlap = !query.terms.is_empty()
951 && self.overlap_possible(
952 fact,
953 query.terms,
954 query.term_count,
955 summaries_trustworthy,
956 );
957 if lexical_only && !may_overlap {
958 continue;
959 }
960 let Some(record) = self.fact(fact) else {
961 continue;
962 };
963 if record.is_tombstone() || record.is_closed() {
964 continue;
965 }
966
967 let mut lexical = None;
969 if may_overlap && let Ok(text) = core::str::from_utf8(self.texts.get(record.text)) {
972 candidate_terms.clear();
973 let terms = &self.terms;
974 let cand = &mut *candidate_terms;
975 self.tokenizer.tokenize(text, &mut |token| {
976 if let Some(term) = terms.lookup(token)
977 && !cand.contains(&term.0)
978 {
979 cand.push(term.0);
980 }
981 });
982 if !candidate_terms.is_empty() {
983 let both = candidate_terms
984 .iter()
985 .filter(|term| query.terms.contains(term))
986 .count();
987 let union = candidate_terms.len() + query.term_count - both;
988 let jaccard = both as f32 / union as f32;
989 if jaccard > self.cfg.similar_jaccard {
990 lexical = Some(jaccard);
991 }
992 }
993 }
994
995 let mut vector = None;
997 if record.has_vector() {
998 let cos = match query.vector {
999 SimilarVector::None => 0.0,
1000 SimilarVector::Stored(slot) => self.vecs.cosine_slots(slot, record.vector),
1001 SimilarVector::Encoded(encoded) => {
1002 self.vecs.cosine_encoded_slot(encoded, record.vector)
1003 }
1004 };
1005 if cos > self.cfg.similar_cos {
1006 vector = Some(cos);
1007 }
1008 }
1009
1010 let best = match (lexical, vector) {
1012 (Some(l), Some(v)) if v > l => Some((v, SimilarReason::VectorCosine)),
1013 (Some(l), _) => Some((l, SimilarReason::LexicalOverlap)),
1014 (None, Some(v)) => Some((v, SimilarReason::VectorCosine)),
1015 (None, None) => None,
1016 };
1017 if let Some((score, reason)) = best {
1018 out.push(Similar {
1019 id: fact,
1020 score,
1021 reason,
1022 });
1023 }
1024 }
1025 out.sort_unstable_by(|a, b| b.score.total_cmp(&a.score).then(a.id.cmp(&b.id)));
1026 out.truncate(8);
1027 }
1028
1029 fn overlap_possible(
1040 &self,
1041 candidate: FactId,
1042 new_terms: &[u32],
1043 new_term_count: usize,
1044 trust_summary: bool,
1045 ) -> bool {
1046 if !trust_summary {
1047 return true;
1048 }
1049 let Some(doc) = self.bm25.doc(candidate) else {
1050 return true;
1051 };
1052 if !doc.has_signature() {
1053 return true;
1054 }
1055 let bound = doc.overlap_bound(new_terms);
1056 debug_assert!(
1062 bound <= new_terms.len(),
1063 "the overlap bound counts query terms, so it cannot exceed them"
1064 );
1065 let union = (new_term_count - bound) + usize::from(doc.distinct);
1066 bound as f32 / union as f32 > self.cfg.similar_jaccard
1067 }
1068
1069 pub fn revise<S: Storage>(
1073 &mut self,
1074 store: &mut S,
1075 target: FactId,
1076 input: RememberInput<'_>,
1077 ) -> Result<RememberOutcome, Error> {
1078 self.validate_input(&input)?;
1079 self.check_revisable(target)?;
1083 let outcome = self.apply_remember(&input, target, None)?;
1084 let valid_from = input.valid_from.unwrap_or(input.now);
1085 self.close_target(target, valid_from);
1086 self.journal_remember(store, &input, target, outcome.id)?;
1087 Ok(outcome)
1088 }
1089
1090 pub fn forget<S: Storage>(
1094 &mut self,
1095 store: &mut S,
1096 now: u64,
1097 id: FactId,
1098 ) -> Result<bool, Error> {
1099 let fresh = self.apply_forget(id)?;
1100 let mut entry = Vec::new();
1101 Op::Forget { now, fact: id }.encode(&mut entry);
1102 store
1103 .append_journal(&entry)
1104 .map_err(|e| Error::Storage(format!("{e:?}")))?;
1105 Ok(fresh)
1106 }
1107
1108 pub fn remove_tag<S: Storage>(
1116 &mut self,
1117 store: &mut S,
1118 now: u64,
1119 tag: &str,
1120 ) -> Result<RemoveTagReport, Error> {
1121 if tag.is_empty() {
1122 return Err(Error::Invalid("empty tag"));
1123 }
1124 let report = self.apply_remove_tag(now, tag)?;
1125 if report.affected != 0 {
1126 let mut entry = Vec::new();
1127 Op::RemoveTag { now, tag }.encode(&mut entry);
1128 store
1129 .append_journal(&entry)
1130 .map_err(|e| Error::Storage(format!("{e:?}")))?;
1131 }
1132 Ok(report)
1133 }
1134
1135 pub fn link<S: Storage>(&mut self, store: &mut S, input: LinkInput<'_>) -> Result<(), Error> {
1138 self.apply_link(
1139 input.now,
1140 input.src,
1141 input.rel,
1142 input.dst,
1143 FactId::from_opt(input.provenance),
1144 )?;
1145 let mut entry = Vec::new();
1146 Op::Link {
1147 now: input.now,
1148 src: input.src,
1149 rel: input.rel,
1150 dst: input.dst,
1151 provenance: FactId::from_opt(input.provenance),
1152 }
1153 .encode(&mut entry);
1154 store
1155 .append_journal(&entry)
1156 .map_err(|e| Error::Storage(format!("{e:?}")))?;
1157 Ok(())
1158 }
1159
1160 pub fn unlink<S: Storage>(
1163 &mut self,
1164 store: &mut S,
1165 input: UnlinkInput<'_>,
1166 ) -> Result<bool, Error> {
1167 let fresh = self.apply_unlink(input.now, input.src, input.rel, input.dst)?;
1168 let mut entry = Vec::new();
1169 Op::Unlink {
1170 now: input.now,
1171 src: input.src,
1172 rel: input.rel,
1173 dst: input.dst,
1174 }
1175 .encode(&mut entry);
1176 store
1177 .append_journal(&entry)
1178 .map_err(|e| Error::Storage(format!("{e:?}")))?;
1179 Ok(fresh)
1180 }
1181
1182 pub fn get(&self, id: FactId) -> Option<FactView<'_>> {
1185 let record = self.fact(id)?;
1186 if record.is_tombstone() {
1187 return None;
1188 }
1189 let text = core::str::from_utf8(self.texts.get(record.text)).ok()?;
1193 Some(FactView { record, text })
1194 }
1195
1196 pub fn tags_of(&self, id: FactId, out: &mut Vec<TermId>) {
1199 let Some(record) = self.fact(id) else { return };
1200 if record.is_tombstone() {
1201 return;
1202 }
1203 let Some(aux) = self.fact_aux.get(&id.0.to_be_bytes()) else {
1204 return;
1205 };
1206 for chunk in self.tag_lists.iter(&aux.tags) {
1207 for raw in chunk.chunks_exact(4) {
1208 out.push(TermId(u32::from_be_bytes(raw.try_into().unwrap())));
1209 }
1210 }
1211 }
1212
1213 pub fn list_tags(&self, query: TagQuery<'_>) -> Result<TagPage, Error> {
1220 self.tag_catalog.page(&self.terms, self.cfg.db_uuid, query)
1221 }
1222
1223 pub fn metadata_of<'s>(&'s self, id: FactId, out: &mut Vec<(&'s str, &'s str)>) -> bool {
1234 out.clear();
1235 let Some(record) = self.fact(id) else {
1236 return false;
1237 };
1238 if record.is_tombstone() {
1239 return false;
1240 }
1241 let Some(aux) = self.fact_aux.get(&id.0.to_be_bytes()) else {
1242 return false;
1243 };
1244 if aux.meta.0 == NONE_U32 || aux.meta.0 >= self.metas.len() as u32 {
1245 return false;
1246 }
1247 crate::metadata::decode(self.metas.get(aux.meta), out).is_ok() && !out.is_empty()
1248 }
1249
1250 pub fn entity(&mut self, name: &str) -> Option<EntityId> {
1252 let mut norm = core::mem::take(&mut self.name_scratch);
1253 normalize_name(&mut self.tokenizer, name, &mut norm);
1254 let found = if norm.is_empty() {
1255 None
1256 } else {
1257 self.lookup_entity_by_norm(&norm)
1259 };
1260 self.name_scratch = norm;
1261 found
1262 }
1263
1264 pub fn term(&self, id: TermId) -> &str {
1266 self.terms.resolve(id)
1267 }
1268
1269 pub fn entity_name(&self, id: EntityId) -> Option<&str> {
1274 let record = self.entities.get(&id.0.to_be_bytes())?;
1275 core::str::from_utf8(self.texts.get(record.name)).ok()
1278 }
1279
1280 pub fn edges_each(&self, mut visit: impl FnMut(&str, &str, &str, FactId) -> bool) {
1298 for slot in self.edges_out.iter() {
1299 let (Some(src), Some(dst)) = (self.entity_name(slot.a), self.entity_name(slot.b))
1300 else {
1301 continue;
1302 };
1303 if !visit(src, self.term(slot.rel), dst, slot.fact) {
1304 return;
1305 }
1306 }
1307 }
1308
1309 pub fn facts_len(&self) -> usize {
1313 self.facts.len()
1314 }
1315
1316 pub fn entities_len(&self) -> usize {
1318 self.entities.len()
1319 }
1320
1321 pub fn cfg(&self) -> &Config {
1323 &self.cfg
1324 }
1325
1326 pub fn stats(&self) -> Stats {
1328 Stats {
1329 facts: self.facts.len(),
1330 entities: self.entities.len(),
1331 terms: self.terms.len(),
1332 edges: self.edges_out.len(),
1333 edge_versions: self.edges_hist_out.len(),
1334 vectors: self.vecs.len(),
1335 tombstones: self.tombstones,
1336 hnsw_indexed: self.hnsw.indexed(),
1337 next_fact: self.next_fact,
1338 next_entity: self.next_entity,
1339 next_edge: self.next_edge,
1340 db_uuid: self.cfg.db_uuid,
1341 pool_bytes: self.facts.pool_bytes()
1342 + self.fact_aux.pool_bytes()
1343 + self.entities.pool_bytes()
1344 + self.by_name.pool_bytes()
1345 + self.edges_out.pool_bytes()
1346 + self.edges_in.pool_bytes()
1347 + self.edges_hist_out.pool_bytes()
1348 + self.edges_hist_in.pool_bytes()
1349 + self.temporal.pool_bytes()
1350 + self.texts.pool_bytes()
1351 + self.terms.pool_bytes()
1352 + self.tag_lists.pool_bytes()
1353 + self.bm25.pool_bytes()
1354 + self.tags_idx.pool_bytes()
1355 + self.tag_catalog.pool_bytes()
1356 + self.entity_facts.pool_bytes()
1357 + self.vecs.pool_bytes()
1358 + self.hnsw.pool_bytes(),
1359 shards: shards::ShardLayout::of_config(&self.cfg),
1360 }
1361 }
1362
1363 fn fact(&self, id: FactId) -> Option<FactRecord> {
1366 self.facts.get(&id.0.to_be_bytes())
1367 }
1368
1369 fn validate_input(&self, input: &RememberInput<'_>) -> Result<(), Error> {
1371 if input.text.len() > self.cfg.max_text {
1372 return Err(Error::TooLarge {
1373 what: "text",
1374 len: input.text.len(),
1375 max: self.cfg.max_text,
1376 });
1377 }
1378 if input.tags.len() > 32 {
1379 return Err(Error::TooLarge {
1380 what: "tags",
1381 len: input.tags.len(),
1382 max: 32,
1383 });
1384 }
1385 if input.links.len() > 16 {
1386 return Err(Error::TooLarge {
1387 what: "links",
1388 len: input.links.len(),
1389 max: 16,
1390 });
1391 }
1392 if input.tags.iter().any(|t| t.is_empty()) {
1393 return Err(Error::Invalid("empty tag"));
1394 }
1395 if !input.links.is_empty() && input.entity.is_none() {
1396 return Err(Error::Invalid("links require a subject entity"));
1397 }
1398 if let Some(v) = input.vector {
1399 if self.cfg.dim == 0 {
1400 return Err(Error::Invalid("vector given but dim is 0"));
1401 }
1402 if v.len() != self.cfg.dim {
1403 return Err(Error::DimMismatch {
1404 got: v.len(),
1405 want: self.cfg.dim,
1406 });
1407 }
1408 }
1409 Ok(())
1410 }
1411
1412 fn check_revisable(&self, target: FactId) -> Result<(), Error> {
1414 let record = self.fact(target).ok_or(Error::NotFound(target))?;
1415 if record.is_tombstone() {
1416 return Err(Error::NotFound(target));
1417 }
1418 if record.is_closed() {
1419 return Err(Error::AlreadyClosed(target));
1420 }
1421 Ok(())
1422 }
1423
1424 fn close_target(&mut self, target: FactId, valid_to: u64) {
1427 let record = self.fact(target).expect("checked revisable");
1428 let payload = self
1429 .facts
1430 .payload_mut(&target.0.to_be_bytes())
1431 .expect("record fetched above");
1432 let flags = record.flags | fact_flags::CLOSED;
1435 payload[4..6].copy_from_slice(&flags.to_be_bytes());
1436 payload[36..44].copy_from_slice(&valid_to.to_be_bytes());
1437 self.change_catalog_for_fact(target, -1);
1438 }
1439
1440 fn apply_remember(
1443 &mut self,
1444 input: &RememberInput<'_>,
1445 revises: FactId,
1446 copy_vector: Option<u32>,
1447 ) -> Result<RememberOutcome, Error> {
1448 let id = FactId(self.next_fact);
1449 let entity = match input.entity {
1450 Some(name) => Some(self.resolve_or_create_entity(name, input.now)?),
1451 None => None,
1452 };
1453 let text_id = self.texts.push(input.text.as_bytes())?;
1454
1455 let mut tfs = core::mem::take(&mut self.tf_scratch);
1457 tfs.clear();
1458 let terms = &mut self.terms;
1459 let mut intern_err = None;
1460 self.tokenizer.tokenize(input.text, &mut |token| {
1461 if intern_err.is_some() {
1462 return;
1463 }
1464 match terms.intern(token) {
1465 Ok(term) => match tfs.iter_mut().find(|(t, _)| *t == term.0) {
1466 Some((_, tf)) => *tf = tf.saturating_add(1),
1467 None => tfs.push((term.0, 1)),
1468 },
1469 Err(e) => intern_err = Some(e),
1470 }
1471 });
1472 if let Some(e) = intern_err {
1473 self.tf_scratch = tfs;
1474 return Err(Error::Arena(e));
1475 }
1476 self.bm25.index_doc(id, &tfs)?;
1477 self.tf_scratch = tfs;
1478
1479 let meta = match input.metadata {
1483 Some(pairs) if !pairs.is_empty() => {
1484 self.metas.push(&crate::metadata::encode(pairs)?)?
1485 }
1486 _ => BlobId(NONE_U32),
1487 };
1488
1489 let mut aux = FactAux {
1492 id,
1493 tags: ListHandle::EMPTY,
1494 meta,
1495 };
1496 let mut seen_tags: [u32; 32] = [NONE_U32; 32];
1497 let mut seen_cnt = 0usize;
1498 for tag in input.tags {
1499 let term = self.terms.intern(tag)?;
1500 if seen_tags[..seen_cnt].contains(&term.0) {
1501 continue;
1502 }
1503 seen_tags[seen_cnt] = term.0;
1504 seen_cnt += 1;
1505 self.tag_lists.push(&mut aux.tags, &term.0.to_be_bytes())?;
1506 self.tags_idx.push(term.0, id, 0)?;
1507 }
1508 self.fact_aux.insert(&aux)?;
1509
1510 if let Some(src) = entity {
1512 for &(rel, dst_name) in input.links {
1513 let dst = self.resolve_or_create_entity(dst_name, input.now)?;
1514 let rel = self.terms.intern(rel)?;
1515 self.open_edge(input.now, src, rel, dst, id)?;
1516 }
1517 self.entity_facts.push(src.0, id, 0)?;
1518 }
1519
1520 let (vector, flags) = match (input.vector, copy_vector) {
1524 (Some(v), None) => (self.vecs.push(id, v)?, fact_flags::HAS_VECTOR),
1525 (None, Some(source)) => (
1526 self.vecs.clone_slot_for_fact(id, source)?,
1527 fact_flags::HAS_VECTOR,
1528 ),
1529 (None, None) => (NONE_U32, 0),
1530 (Some(_), Some(_)) => unreachable!("retag does not provide a raw vector"),
1531 };
1532
1533 let recorded_at = input.now;
1534 let valid_from = input.valid_from.unwrap_or(input.now);
1535 self.facts.insert(&FactRecord {
1536 id,
1537 entity: EntityId::from_opt(entity),
1538 flags,
1539 kind: 0,
1540 text: text_id,
1541 vector,
1542 revises,
1543 recorded_at,
1544 valid_from,
1545 valid_to: VALID_TO_OPEN,
1546 })?;
1547 self.temporal.insert(&TemporalSlot {
1548 recorded_at,
1549 fact: id,
1550 })?;
1551 for &term in &seen_tags[..seen_cnt] {
1552 self.tag_catalog.change(&self.terms, TermId(term), 1);
1553 }
1554 self.next_fact += 1;
1555 Ok(RememberOutcome {
1556 id,
1557 entity,
1558 similar: Vec::new(),
1559 })
1560 }
1561
1562 fn apply_forget(&mut self, id: FactId) -> Result<bool, Error> {
1563 let record = self.fact(id).ok_or(Error::NotFound(id))?;
1564 if record.is_tombstone() {
1565 return Ok(false);
1566 }
1567 let payload = self
1568 .facts
1569 .payload_mut(&id.0.to_be_bytes())
1570 .expect("record fetched above");
1571 let flags = record.flags | fact_flags::TOMBSTONE;
1572 payload[4..6].copy_from_slice(&flags.to_be_bytes());
1573 self.tombstones += 1;
1574 if !record.is_closed() {
1575 self.change_catalog_for_fact(id, -1);
1576 }
1577 Ok(true)
1578 }
1579
1580 fn apply_remove_tag(&mut self, now: u64, tag: &str) -> Result<RemoveTagReport, Error> {
1581 let Some(term) = self.terms.lookup(tag) else {
1582 return Ok(RemoveTagReport::default());
1583 };
1584 let targets: Vec<FactId> = self
1585 .tags_idx
1586 .entries(term.0)
1587 .map(|(id, _)| id)
1588 .filter(|&id| {
1589 self.fact(id)
1590 .is_some_and(|record| !record.is_tombstone() && !record.is_closed())
1591 })
1592 .collect();
1593 let mut affected = 0u32;
1594 for target in targets {
1595 self.retag_without(now, target, tag)?;
1596 affected = affected.saturating_add(1);
1597 }
1598 Ok(RemoveTagReport { affected })
1599 }
1600
1601 fn retag_without(&mut self, now: u64, target: FactId, removed: &str) -> Result<(), Error> {
1602 let record = self.fact(target).ok_or(Error::NotFound(target))?;
1603 let view = self.get(target).ok_or(Error::NotFound(target))?;
1604 let text = view.text.to_string();
1605 let entity = record
1606 .entity
1607 .some()
1608 .and_then(|id| self.entity_name(id))
1609 .map(ToString::to_string);
1610
1611 let mut tag_terms = Vec::new();
1612 self.tags_of(target, &mut tag_terms);
1613 let tags: Vec<String> = tag_terms
1614 .into_iter()
1615 .map(|term| self.term(term))
1616 .filter(|name| *name != removed)
1617 .map(ToString::to_string)
1618 .collect();
1619 let tag_refs: Vec<&str> = tags.iter().map(String::as_str).collect();
1620
1621 let mut metadata = Vec::new();
1622 self.metadata_of(target, &mut metadata);
1623 let metadata: Vec<(String, String)> = metadata
1624 .into_iter()
1625 .map(|(key, value)| (key.to_string(), value.to_string()))
1626 .collect();
1627 let metadata_refs: Vec<(&str, &str)> = metadata
1628 .iter()
1629 .map(|(key, value)| (key.as_str(), value.as_str()))
1630 .collect();
1631 let vector = (record.flags & fact_flags::HAS_VECTOR != 0).then_some(record.vector);
1632 let input = RememberInput {
1633 now,
1634 text: &text,
1635 entity: entity.as_deref(),
1636 tags: &tag_refs,
1637 links: &[],
1638 vector: None,
1639 valid_from: Some(now),
1640 metadata: (!metadata_refs.is_empty()).then_some(metadata_refs.as_slice()),
1641 };
1642 self.apply_remember(&input, target, vector)?;
1643 self.close_target(target, now);
1644 Ok(())
1645 }
1646
1647 fn change_catalog_for_fact(&mut self, id: FactId, delta: i32) {
1651 let Some(aux) = self.fact_aux.get(&id.0.to_be_bytes()) else {
1652 return;
1653 };
1654 let mut ids = [NONE_U32; 32];
1655 let mut len = 0usize;
1656 for chunk in self.tag_lists.iter(&aux.tags) {
1657 for raw in chunk.chunks_exact(4) {
1658 if len == ids.len() {
1659 break;
1660 }
1661 ids[len] = u32::from_be_bytes(raw.try_into().unwrap());
1662 len += 1;
1663 }
1664 }
1665 for &term in &ids[..len] {
1666 self.tag_catalog.change(&self.terms, TermId(term), delta);
1667 }
1668 }
1669
1670 fn apply_link(
1671 &mut self,
1672 now: u64,
1673 src: &str,
1674 rel: &str,
1675 dst: &str,
1676 provenance: FactId,
1677 ) -> Result<(), Error> {
1678 let src = self.resolve_or_create_entity(src, now)?;
1679 let dst = self.resolve_or_create_entity(dst, now)?;
1680 let rel = self.terms.intern(rel)?;
1681 self.open_edge(now, src, rel, dst, provenance)
1682 }
1683
1684 fn apply_unlink(&mut self, now: u64, src: &str, rel: &str, dst: &str) -> Result<bool, Error> {
1685 let Some(src) = self.lookup_entity_name(src) else {
1686 return Ok(false);
1687 };
1688 let Some(dst) = self.lookup_entity_name(dst) else {
1689 return Ok(false);
1690 };
1691 let Some(rel) = self.terms.lookup(rel) else {
1692 return Ok(false);
1693 };
1694 self.close_current_edge(now, src, rel, dst)
1695 }
1696
1697 fn open_edge(
1698 &mut self,
1699 now: u64,
1700 src: EntityId,
1701 rel: TermId,
1702 dst: EntityId,
1703 fact: FactId,
1704 ) -> Result<(), Error> {
1705 if let Some(current) = self.current_edge(src, rel, dst) {
1706 if current.fact == fact {
1707 return Ok(());
1708 }
1709 self.close_current_edge(now, src, rel, dst)?;
1710 }
1711 let edge = EdgeId(self.next_edge);
1712 let history = EdgeHistorySlot {
1713 a: src,
1714 rel,
1715 b: dst,
1716 edge,
1717 fact,
1718 flags: 0,
1719 kind: 0,
1720 recorded_at: now,
1721 valid_from: now,
1722 valid_to: VALID_TO_OPEN,
1723 };
1724 self.insert_history_edge(history)?;
1725 self.insert_current_edge(src, rel, dst, fact, edge, now)?;
1726 self.next_edge += 1;
1727 Ok(())
1728 }
1729
1730 fn insert_current_edge(
1731 &mut self,
1732 src: EntityId,
1733 rel: TermId,
1734 dst: EntityId,
1735 fact: FactId,
1736 edge: EdgeId,
1737 valid_from: u64,
1738 ) -> Result<(), Error> {
1739 for (arena, a, b) in [
1740 (&mut self.edges_out, src, dst),
1741 (&mut self.edges_in, dst, src),
1742 ] {
1743 let slot = EdgeSlot {
1744 a,
1745 rel,
1746 b,
1747 fact,
1748 edge,
1749 valid_from,
1750 };
1751 if !arena.insert(&slot)? {
1752 let payload = arena
1753 .payload_mut(&edge_key(a, rel, b))
1754 .expect("insert reported a duplicate");
1755 let mut full = [0u8; EdgeSlot::SIZE];
1756 slot.write(&mut full);
1757 payload.copy_from_slice(&full[EdgeSlot::KEY_LEN..]);
1758 }
1759 }
1760 Ok(())
1761 }
1762
1763 fn insert_history_edge(&mut self, edge: EdgeHistorySlot) -> Result<(), Error> {
1764 self.edges_hist_out.insert(&edge)?;
1765 self.edges_hist_in.insert(&EdgeHistorySlot {
1766 a: edge.b,
1767 b: edge.a,
1768 ..edge
1769 })?;
1770 Ok(())
1771 }
1772
1773 fn close_current_edge(
1779 &mut self,
1780 now: u64,
1781 src: EntityId,
1782 rel: TermId,
1783 dst: EntityId,
1784 ) -> Result<bool, Error> {
1785 let Some(current) = self.current_edge(src, rel, dst) else {
1786 return Ok(false);
1787 };
1788 let close_at = now.max(current.valid_from);
1791 let out_key = edge_history_key(src, current.valid_from, current.edge);
1792 let in_key = edge_history_key(dst, current.valid_from, current.edge);
1793 close_edge_history_payload(
1794 self.edges_hist_out
1795 .payload_mut(&out_key)
1796 .ok_or(Error::Corrupt("missing outgoing edge history"))?,
1797 close_at,
1798 );
1799 close_edge_history_payload(
1800 self.edges_hist_in
1801 .payload_mut(&in_key)
1802 .ok_or(Error::Corrupt("missing incoming edge history"))?,
1803 close_at,
1804 );
1805 let out_removed = self.edges_out.remove(&edge_key(src, rel, dst));
1806 let in_removed = self.edges_in.remove(&edge_key(dst, rel, src));
1807 if out_removed != in_removed {
1808 return Err(Error::Corrupt("edge mirrors disagree"));
1809 }
1810 Ok(out_removed)
1811 }
1812
1813 fn current_edge(&self, src: EntityId, rel: TermId, dst: EntityId) -> Option<EdgeSlot> {
1814 self.edges_out
1815 .get_slot(&edge_key(src, rel, dst))
1816 .map(EdgeSlot::read)
1817 }
1818
1819 fn lookup_entity_by_norm(&self, norm: &str) -> Option<EntityId> {
1822 let term = self.terms.lookup(norm)?;
1823 let mut from = [0u8; 8];
1824 key::write_u32(&mut from, term.0);
1825 let mut to = [0u8; 8];
1826 key::write_u32(&mut to, term.0);
1827 to[4..].copy_from_slice(&u32::MAX.to_be_bytes());
1828 self.by_name.range(&from, &to).next().map(|e| e.id)
1829 }
1830
1831 fn lookup_entity_name(&mut self, name: &str) -> Option<EntityId> {
1832 let mut norm = core::mem::take(&mut self.name_scratch);
1833 normalize_name(&mut self.tokenizer, name, &mut norm);
1834 let result = (!norm.is_empty())
1835 .then(|| self.lookup_entity_by_norm(&norm))
1836 .flatten();
1837 self.name_scratch = norm;
1838 result
1839 }
1840
1841 fn resolve_or_create_entity(&mut self, name: &str, now: u64) -> Result<EntityId, Error> {
1842 let mut norm = core::mem::take(&mut self.name_scratch);
1843 normalize_name(&mut self.tokenizer, name, &mut norm);
1844 if norm.is_empty() {
1845 self.name_scratch = norm;
1846 return Err(Error::Invalid("entity name has no indexable characters"));
1847 }
1848 let result = (|| {
1849 if let Some(found) = self.lookup_entity_by_norm(&norm) {
1850 return Ok(found);
1851 }
1852 let term = self.terms.intern(&norm)?;
1853 let id = EntityId(self.next_entity);
1854 let name_id = self.texts.push(name.as_bytes())?;
1855 self.entities.insert(&EntityRecord {
1856 id,
1857 name: name_id,
1858 name_term: term,
1859 created_at: now,
1860 flags: 0,
1861 })?;
1862 self.by_name.insert(&EntityByName {
1863 name_term: term,
1864 id,
1865 })?;
1866 self.next_entity += 1;
1867 Ok(id)
1868 })();
1869 self.name_scratch = norm;
1870 result
1871 }
1872
1873 fn journal_remember<S: Storage>(
1874 &mut self,
1875 store: &mut S,
1876 input: &RememberInput<'_>,
1877 revises: FactId,
1878 assigned: FactId,
1879 ) -> Result<(), Error> {
1880 let mut entry = Vec::new();
1881 Op::Remember {
1882 now: input.now,
1883 valid_from: input.valid_from.unwrap_or(input.now),
1884 entity: input.entity,
1885 text: input.text,
1886 tags: input.tags.to_vec(),
1887 links: input.links.to_vec(),
1888 vector: input.vector.map(<[f32]>::to_vec).unwrap_or_default(),
1889 metadata: input.metadata.map(<[_]>::to_vec).unwrap_or_default(),
1890 revises,
1891 assigned,
1892 }
1893 .encode(&mut entry);
1894 store
1895 .append_journal(&entry)
1896 .map_err(|e| Error::Storage(format!("{e:?}")))
1897 }
1898}
1899
1900impl core::fmt::Debug for Memory<'_> {
1901 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
1904 f.debug_struct("Memory")
1905 .field("facts", &self.facts.len())
1906 .field("entities", &self.entities.len())
1907 .field("terms", &self.terms.len())
1908 .finish()
1909 }
1910}
1911
1912fn normalize_name(tokenizer: &mut Tokenizer, name: &str, out: &mut String) {
1916 out.clear();
1917 tokenizer.tokenize(name, &mut |token| {
1918 if !out.is_empty() {
1919 out.push(' ');
1920 }
1921 out.push_str(token);
1922 });
1923}