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)]
163#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
164#[cfg_attr(feature = "serde", serde(tag = "status", rename_all = "snake_case"))]
165pub enum GuardedRememberOutcome {
166 Stored {
169 outcome: RememberOutcome,
171 checked: bool,
192 },
193 Blocked {
196 similar: Vec<Similar>,
199 },
200}
201
202#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
204#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
205pub struct RemoveTagReport {
206 pub affected: u32,
209}
210
211#[derive(Clone, Copy, Debug, PartialEq)]
213#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
214pub struct Similar {
215 pub id: FactId,
217 pub score: f32,
220 pub reason: SimilarReason,
222}
223
224#[derive(Clone, Copy, Debug, PartialEq, Eq)]
226#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
227pub enum SimilarReason {
228 LexicalOverlap,
231 VectorCosine,
234}
235
236#[derive(Debug, Default)]
240struct SimilarityScratch {
241 new_terms: Vec<u32>,
243 unknown: String,
246 unknown_ranges: Vec<(usize, usize)>,
248 candidate_terms: Vec<u32>,
250 vector: Vec<u8>,
252}
253
254#[derive(Clone, Copy)]
255enum SimilarVector<'s> {
256 None,
257 Stored(u32),
258 Encoded(&'s [u8]),
259}
260
261#[derive(Clone, Copy)]
262struct SimilarityQuery<'s> {
263 entity: EntityId,
264 exclude: Option<FactId>,
265 terms: &'s [u32],
266 term_count: usize,
267 vector: SimilarVector<'s>,
268}
269
270#[derive(Clone, Copy, Debug)]
272#[cfg_attr(feature = "serde", derive(serde::Serialize))]
273pub struct FactView<'a> {
274 pub record: FactRecord,
276 pub text: &'a str,
278}
279
280#[derive(Clone, Copy, Debug, PartialEq, Eq)]
285#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
286pub enum FactFault {
287 Text,
289 Vector,
292 Metadata,
295}
296
297#[derive(Clone, Copy, Debug, PartialEq, Eq)]
301#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
302#[non_exhaustive]
303pub struct Stats {
304 pub facts: usize,
307 pub entities: usize,
309 pub terms: usize,
311 pub edges: usize,
314 pub edge_versions: usize,
316 pub vectors: usize,
318 pub tombstones: usize,
320 pub hnsw_indexed: u32,
323 pub next_fact: u32,
326 pub next_entity: u32,
328 pub next_edge: u32,
330 pub db_uuid: u128,
333 pub pool_bytes: usize,
336 pub shards: ShardLayout,
342}
343
344#[derive(Clone, Debug, Default, PartialEq, Eq)]
346#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
347pub struct OpenReport {
348 pub replayed: usize,
350 pub skipped: usize,
352 pub truncated_tail: bool,
354}
355
356pub struct Memory<'a> {
363 cfg: Config,
364 facts: Arena<'a, FactRecord>,
366 fact_aux: Arena<'a, FactAux>,
367 entities: Arena<'a, EntityRecord>,
368 by_name: Arena<'a, EntityByName>,
369 edges_out: Arena<'a, EdgeSlot>,
370 edges_in: Arena<'a, EdgeSlot>,
371 edges_hist_out: Arena<'a, EdgeHistorySlot>,
372 edges_hist_in: Arena<'a, EdgeHistorySlot>,
373 temporal: Arena<'a, TemporalSlot>,
374 texts: BlobHeap<'a>,
376 metas: BlobHeap<'a>,
381 terms: Interner<'a>,
384 tag_lists: ChunkPool<'a>,
386 bm25: Bm25Index<'a>,
388 tags_idx: IdListIndex<'a>,
389 tag_catalog: tags::TagCatalog,
391 entity_facts: IdListIndex<'a>,
392 vecs: VecPool<'a>,
394 hnsw: HnswGraph<'a>,
398 vector_space: Option<String>,
402 next_fact: u32,
404 next_entity: u32,
405 next_edge: u32,
406 tombstones: usize,
408 bm25_tokenizer_version: u32,
409 tokenizer: Tokenizer,
411 tf_scratch: Vec<(u32, u8)>,
412 name_scratch: String,
413 similarity_scratch: SimilarityScratch,
414}
415
416impl<'a> Memory<'a> {
417 pub fn new(cfg: Config) -> Result<Self, Error> {
424 cfg.validate()?;
425 let uni =
426 |shards: usize| ArenaCfg::new(shards, ShardMode::Uniform).with_max_bytes(cfg.max_bytes);
427 let ord =
428 |shards: usize| ArenaCfg::new(shards, ShardMode::Ordered).with_max_bytes(cfg.max_bytes);
429 let blob = BlobHeapCfg::new()
430 .with_max_bytes(cfg.max_bytes)
431 .with_max_blob(cfg.max_blob);
432 Ok(Self {
433 facts: Arena::new(uni(cfg.shards_facts))?,
434 fact_aux: Arena::new(uni(cfg.shards_facts))?,
435 entities: Arena::new(uni(cfg.shards_entities))?,
436 by_name: Arena::new(ord(cfg.shards_entities))?,
437 edges_out: Arena::new(ord(cfg.shards_edges))?,
438 edges_in: Arena::new(ord(cfg.shards_edges))?,
439 edges_hist_out: Arena::new(ord(cfg.shards_edges))?,
440 edges_hist_in: Arena::new(ord(cfg.shards_edges))?,
441 temporal: Arena::new(ord(cfg.shards_temporal))?,
442 texts: BlobHeap::new(blob),
443 metas: BlobHeap::new(blob),
444 terms: Interner::new(blob),
445 tag_lists: ChunkPool::new(ChunkPoolCfg::new().with_max_bytes(cfg.max_bytes)),
446 bm25: Bm25Index::new(cfg.shards_postings, cfg.max_bytes)?,
447 tags_idx: IdListIndex::new(cfg.shards_postings, cfg.max_bytes)?,
448 tag_catalog: tags::TagCatalog::new(),
449 entity_facts: IdListIndex::new(cfg.shards_entities, cfg.max_bytes)?,
450 vecs: VecPool::new(cfg.dim, cfg.max_bytes),
451 hnsw: HnswGraph::new(cfg.hnsw_m, cfg.hnsw_m0, cfg.max_bytes)?,
452 vector_space: None,
453 next_fact: 0,
454 next_entity: 0,
455 next_edge: 0,
456 tombstones: 0,
457 bm25_tokenizer_version: maintain::TOKENIZER_INDEX_VERSION,
458 tokenizer: Tokenizer::new(),
459 tf_scratch: Vec::new(),
460 name_scratch: String::new(),
461 similarity_scratch: SimilarityScratch::default(),
462 cfg,
463 })
464 }
465
466 pub fn open<S: Storage>(store: &mut S, cfg: Config) -> Result<(Self, OpenReport), Error> {
470 let snapshot = store
471 .read_snapshot()
472 .map_err(|e| Error::Storage(format!("{e:?}")))?;
473 let journal = store
474 .read_journal()
475 .map_err(|e| Error::Storage(format!("{e:?}")))?;
476 Self::from_bytes(snapshot.as_deref(), &journal, cfg)
477 }
478
479 pub fn from_bytes(
481 snapshot: Option<&[u8]>,
482 journal: &[u8],
483 cfg: Config,
484 ) -> Result<(Self, OpenReport), Error> {
485 let mut mem = match snapshot {
486 Some(bytes) => Self::load_snapshot(bytes, cfg)?,
487 None => Self::new(cfg)?,
488 };
489 let report = mem.replay(journal)?;
490 Ok((mem, report))
491 }
492
493 pub fn from_bytes_borrowed(
504 snapshot: &'a [u8],
505 journal: &[u8],
506 cfg: Config,
507 ) -> Result<Self, Error> {
508 let mem = Self::load_snapshot_borrowed(snapshot, cfg)?;
509 let JournalScan { entries, .. } = scan(journal)?;
510 if !entries.is_empty() {
511 return Err(Error::Invalid(
512 "read-only open requires a checkpointed (empty) journal",
513 ));
514 }
515 Ok(mem)
516 }
517
518 pub fn from_bytes_overlay(
533 snapshot: &'a [u8],
534 journal: &[u8],
535 cfg: Config,
536 ) -> Result<(Self, OpenReport), Error> {
537 let mut mem = Self::load_snapshot_borrowed(snapshot, cfg)?;
538 let report = mem.replay(journal)?;
539 Ok((mem, report))
540 }
541
542 fn replay(&mut self, journal: &[u8]) -> Result<OpenReport, Error> {
547 let JournalScan {
548 entries,
549 truncated_tail,
550 } = scan(journal)?;
551 let mut report = OpenReport {
552 truncated_tail,
553 ..OpenReport::default()
554 };
555 for entry in entries {
556 let op = Op::decode(entry.op, entry.payload)?;
557 match op {
558 Op::Remember {
559 now,
560 valid_from,
561 entity,
562 text,
563 ref tags,
564 ref links,
565 ref vector,
566 ref metadata,
567 revises,
568 assigned,
569 } => {
570 if assigned.0 < self.next_fact {
571 report.skipped += 1;
572 continue;
573 }
574 if assigned.0 != self.next_fact {
575 return Err(Error::Corrupt("journal fact ids are not contiguous"));
576 }
577 if !vector.is_empty() && vector.len() != self.cfg.dim {
578 return Err(Error::Corrupt(
579 "journal vector dimension disagrees with dim",
580 ));
581 }
582 if let Some(target) = revises.some() {
583 self.check_revisable(target)
584 .map_err(|_| Error::Corrupt("journal revises an unrevisable fact"))?;
585 }
586 self.apply_remember(
587 &RememberInput {
588 now,
589 text,
590 entity,
591 tags: &tags.to_vec(),
592 links: &links.to_vec(),
593 vector: (!vector.is_empty()).then_some(vector.as_slice()),
594 valid_from: Some(valid_from),
595 metadata: (!metadata.is_empty()).then_some(metadata.as_slice()),
596 },
597 revises,
598 None,
599 )?;
600 if let Some(target) = revises.some() {
601 self.close_target(target, valid_from);
602 }
603 report.replayed += 1;
604 }
605 Op::Forget { fact, .. } => {
606 match self.apply_forget(fact) {
609 Ok(_) => report.replayed += 1,
610 Err(Error::NotFound(_)) => {
611 return Err(Error::Corrupt("journal forgets an unknown fact"));
612 }
613 Err(e) => return Err(e),
614 }
615 }
616 Op::Link {
617 now,
618 src,
619 rel,
620 dst,
621 provenance,
622 } => {
623 self.apply_link(now, src, rel, dst, provenance)?;
624 report.replayed += 1;
625 }
626 Op::Unlink { now, src, rel, dst } => {
627 self.apply_unlink(now, src, rel, dst)?;
628 report.replayed += 1;
629 }
630 Op::RemoveTag { now, tag } => {
631 self.apply_remove_tag(now, tag)?;
632 report.replayed += 1;
633 }
634 Op::SetVectorSpace { space } => {
635 self.apply_set_vector_space(space)?;
636 report.replayed += 1;
637 }
638 Op::Maintain {
639 mode,
640 max_hnsw_inserts,
641 ..
642 } => {
643 let options =
647 maintain::MaintenanceOptions::from_journal(mode, max_hnsw_inserts)?;
648 self.replay_maintain_with_options(options)?;
649 report.replayed += 1;
650 }
651 }
652 }
653 Ok(report)
654 }
655
656 pub fn vector_space(&self) -> Option<&str> {
661 self.vector_space.as_deref()
662 }
663
664 pub fn claim_vector_space<S: Storage>(
669 &mut self,
670 store: &mut S,
671 space: &str,
672 ) -> Result<bool, Error> {
673 let needs_claim = self.check_vector_space_claim(space)?;
674 if !needs_claim {
675 return Ok(false);
676 }
677 let mut entry = Vec::new();
678 Op::SetVectorSpace { space }.encode(&mut entry);
679 store
680 .append_journal(&entry)
681 .map_err(|e| Error::Storage(format!("{e:?}")))?;
682 self.vector_space = Some(space.into());
683 Ok(true)
684 }
685
686 fn check_vector_space_claim(&self, space: &str) -> Result<bool, Error> {
689 Self::validate_vector_space(space)?;
690 if let Some(stored) = &self.vector_space {
691 if stored == space {
692 return Ok(false);
693 }
694 if !self.vecs.is_empty() {
695 return Err(Error::VectorSpaceMismatch {
696 stored: stored.clone(),
697 requested: space.into(),
698 });
699 }
700 }
701 if !self.vecs.is_empty() {
702 return Err(Error::UntrackedVectorSpace);
703 }
704 Ok(true)
705 }
706
707 fn apply_set_vector_space(&mut self, space: &str) -> Result<(), Error> {
708 Self::validate_vector_space(space)?;
709 if let Some(stored) = &self.vector_space {
710 if stored == space {
711 return Ok(());
712 }
713 if !self.vecs.is_empty() {
714 return Err(Error::Corrupt(
715 "journal changes an established vector space",
716 ));
717 }
718 }
719 if !self.vecs.is_empty() {
720 return Err(Error::Corrupt(
721 "journal assigns a vector space after vector records",
722 ));
723 }
724 self.vector_space = Some(space.into());
725 Ok(())
726 }
727
728 pub(super) fn validate_vector_space(space: &str) -> Result<(), Error> {
729 if space.is_empty() {
730 return Err(Error::Invalid("vector space must not be empty"));
731 }
732 if space.len() > MAX_VECTOR_SPACE_ID_BYTES {
733 return Err(Error::TooLarge {
734 what: "vector space",
735 len: space.len(),
736 max: MAX_VECTOR_SPACE_ID_BYTES,
737 });
738 }
739 if space.bytes().any(|b| b < b' ' || b == 0x7f) {
740 return Err(Error::Invalid(
741 "vector space must not contain control bytes",
742 ));
743 }
744 Ok(())
745 }
746
747 pub fn remember<S: Storage>(
749 &mut self,
750 store: &mut S,
751 input: RememberInput<'_>,
752 ) -> Result<RememberOutcome, Error> {
753 self.validate_input(&input)?;
754 let mut outcome = self.apply_remember(&input, FactId::NONE, None)?;
755 self.find_similar(&mut outcome);
756 self.journal_remember(store, &input, FactId::NONE, outcome.id)?;
757 Ok(outcome)
758 }
759
760 pub fn remember_guarded<S: Storage>(
770 &mut self,
771 store: &mut S,
772 input: RememberInput<'_>,
773 ) -> Result<GuardedRememberOutcome, Error> {
774 self.remember_guarded_with_vector_space(store, input, None)
775 }
776
777 #[doc(hidden)]
782 pub fn remember_guarded_with_vector_space<S: Storage>(
783 &mut self,
784 store: &mut S,
785 input: RememberInput<'_>,
786 vector_space: Option<&str>,
787 ) -> Result<GuardedRememberOutcome, Error> {
788 self.validate_input(&input)?;
789 if let Some(space) = vector_space {
790 self.check_vector_space_claim(space)?;
791 }
792 let checked = input.entity.is_some();
800 let similar = self.find_similar_input(&input)?;
801 if !similar.is_empty() {
802 return Ok(GuardedRememberOutcome::Blocked { similar });
803 }
804 if let Some(space) = vector_space {
805 self.claim_vector_space(store, space)?;
806 }
807 let outcome = self.apply_remember(&input, FactId::NONE, None)?;
808 self.journal_remember(store, &input, FactId::NONE, outcome.id)?;
809 Ok(GuardedRememberOutcome::Stored { outcome, checked })
810 }
811
812 pub fn remember_batch<S: Storage>(
820 &mut self,
821 store: &mut S,
822 inputs: &[RememberInput<'_>],
823 skip_similar: bool,
824 ) -> Result<Vec<RememberOutcome>, Error> {
825 let mut out = Vec::with_capacity(inputs.len());
826 for input in inputs {
827 self.validate_input(input)?;
828 let mut outcome = self.apply_remember(input, FactId::NONE, None)?;
829 if !skip_similar {
830 self.find_similar(&mut outcome);
831 }
832 self.journal_remember(store, input, FactId::NONE, outcome.id)?;
833 out.push(outcome);
834 }
835 Ok(out)
836 }
837
838 fn find_similar_input(&mut self, input: &RememberInput<'_>) -> Result<Vec<Similar>, Error> {
842 let Some(name) = input.entity else {
843 return Ok(Vec::new());
844 };
845 let Some(entity) = self.lookup_entity_name(name) else {
846 return Ok(Vec::new());
847 };
848
849 let mut scratch = core::mem::take(&mut self.similarity_scratch);
850 scratch.new_terms.clear();
851 scratch.unknown.clear();
852 scratch.unknown_ranges.clear();
853 scratch.candidate_terms.clear();
854 scratch.vector.clear();
855
856 let terms = &self.terms;
857 let new_terms = &mut scratch.new_terms;
858 let unknown = &mut scratch.unknown;
859 let unknown_ranges = &mut scratch.unknown_ranges;
860 self.tokenizer.tokenize(input.text, &mut |token| {
861 if let Some(term) = terms.lookup(token) {
862 if !new_terms.contains(&term.0) {
863 new_terms.push(term.0);
864 }
865 return;
866 }
867 if unknown_ranges
868 .iter()
869 .any(|&(start, end)| &unknown[start..end] == token)
870 {
871 return;
872 }
873 let start = unknown.len();
874 unknown.push_str(token);
875 unknown_ranges.push((start, unknown.len()));
876 });
877 let new_term_count = scratch.new_terms.len() + scratch.unknown_ranges.len();
878
879 let mut similar = Vec::new();
880 let result = (|| {
881 let encoded = match input.vector {
882 Some(vector) => {
883 self.vecs
884 .encode_slot_into(FactId::NONE, vector, &mut scratch.vector)?;
885 SimilarVector::Encoded(&scratch.vector)
886 }
887 None => SimilarVector::None,
888 };
889 self.scan_similar(
890 SimilarityQuery {
891 entity,
892 exclude: None,
893 terms: &scratch.new_terms,
894 term_count: new_term_count,
895 vector: encoded,
896 },
897 &mut scratch.candidate_terms,
898 &mut similar,
899 );
900 Ok::<(), Error>(())
901 })();
902 self.similarity_scratch = scratch;
903 result.map(|()| similar)
904 }
905
906 fn find_similar(&mut self, outcome: &mut RememberOutcome) {
910 let Some(entity) = outcome.entity else { return };
911 let new_vec = self
912 .fact(outcome.id)
913 .filter(|record| record.has_vector())
914 .map(|record| record.vector);
915 let mut scratch = core::mem::take(&mut self.similarity_scratch);
916 scratch.new_terms.clear();
917 scratch
918 .new_terms
919 .extend(self.tf_scratch.iter().map(|&(term, _)| term));
920 let new_term_count = scratch.new_terms.len();
921 self.scan_similar(
922 SimilarityQuery {
923 entity,
924 exclude: Some(outcome.id),
925 terms: &scratch.new_terms,
926 term_count: new_term_count,
927 vector: new_vec.map_or(SimilarVector::None, SimilarVector::Stored),
928 },
929 &mut scratch.candidate_terms,
930 &mut outcome.similar,
931 );
932 self.similarity_scratch = scratch;
933 }
934
935 fn scan_similar(
957 &mut self,
958 query: SimilarityQuery<'_>,
959 candidate_terms: &mut Vec<u32>,
960 out: &mut Vec<Similar>,
961 ) {
962 out.clear();
963 if query.term_count == 0 && matches!(query.vector, SimilarVector::None) {
964 return;
965 }
966 let mut ring = [FactId::NONE; SIMILAR_CANDIDATE_CAP];
968 let mut n = 0usize;
969 for (fact, _) in self.entity_facts.entries(query.entity.0) {
970 if query.exclude == Some(fact) {
971 continue;
972 }
973 ring[n % SIMILAR_CANDIDATE_CAP] = fact;
974 n += 1;
975 }
976 let summaries_trustworthy = self.bm25_tokenizer_version == TOKENIZER_INDEX_VERSION;
977 let lexical_only = matches!(query.vector, SimilarVector::None);
984 for &fact in ring.iter().take(n.min(SIMILAR_CANDIDATE_CAP)) {
985 let may_overlap = !query.terms.is_empty()
986 && self.overlap_possible(
987 fact,
988 query.terms,
989 query.term_count,
990 summaries_trustworthy,
991 );
992 if lexical_only && !may_overlap {
993 continue;
994 }
995 let Some(record) = self.fact(fact) else {
996 continue;
997 };
998 if record.is_tombstone() || record.is_closed() {
999 continue;
1000 }
1001
1002 let mut lexical = None;
1004 if may_overlap && let Ok(text) = core::str::from_utf8(self.texts.get(record.text)) {
1007 candidate_terms.clear();
1008 let terms = &self.terms;
1009 let cand = &mut *candidate_terms;
1010 self.tokenizer.tokenize(text, &mut |token| {
1011 if let Some(term) = terms.lookup(token)
1012 && !cand.contains(&term.0)
1013 {
1014 cand.push(term.0);
1015 }
1016 });
1017 if !candidate_terms.is_empty() {
1018 let both = candidate_terms
1019 .iter()
1020 .filter(|term| query.terms.contains(term))
1021 .count();
1022 let union = candidate_terms.len() + query.term_count - both;
1023 let jaccard = both as f32 / union as f32;
1024 if jaccard > self.cfg.similar_jaccard {
1025 lexical = Some(jaccard);
1026 }
1027 }
1028 }
1029
1030 let mut vector = None;
1032 if record.has_vector() {
1033 let cos = match query.vector {
1034 SimilarVector::None => 0.0,
1035 SimilarVector::Stored(slot) => self.vecs.cosine_slots(slot, record.vector),
1036 SimilarVector::Encoded(encoded) => {
1037 self.vecs.cosine_encoded_slot(encoded, record.vector)
1038 }
1039 };
1040 if cos > self.cfg.similar_cos {
1041 vector = Some(cos);
1042 }
1043 }
1044
1045 let best = match (lexical, vector) {
1047 (Some(l), Some(v)) if v > l => Some((v, SimilarReason::VectorCosine)),
1048 (Some(l), _) => Some((l, SimilarReason::LexicalOverlap)),
1049 (None, Some(v)) => Some((v, SimilarReason::VectorCosine)),
1050 (None, None) => None,
1051 };
1052 if let Some((score, reason)) = best {
1053 out.push(Similar {
1054 id: fact,
1055 score,
1056 reason,
1057 });
1058 }
1059 }
1060 out.sort_unstable_by(|a, b| b.score.total_cmp(&a.score).then(a.id.cmp(&b.id)));
1061 out.truncate(8);
1062 }
1063
1064 fn overlap_possible(
1075 &self,
1076 candidate: FactId,
1077 new_terms: &[u32],
1078 new_term_count: usize,
1079 trust_summary: bool,
1080 ) -> bool {
1081 if !trust_summary {
1082 return true;
1083 }
1084 let Some(doc) = self.bm25.doc(candidate) else {
1085 return true;
1086 };
1087 if !doc.has_signature() {
1088 return true;
1089 }
1090 let bound = doc.overlap_bound(new_terms);
1091 debug_assert!(
1097 bound <= new_terms.len(),
1098 "the overlap bound counts query terms, so it cannot exceed them"
1099 );
1100 let union = (new_term_count - bound) + usize::from(doc.distinct);
1101 bound as f32 / union as f32 > self.cfg.similar_jaccard
1102 }
1103
1104 pub fn revise<S: Storage>(
1108 &mut self,
1109 store: &mut S,
1110 target: FactId,
1111 input: RememberInput<'_>,
1112 ) -> Result<RememberOutcome, Error> {
1113 self.validate_input(&input)?;
1114 self.check_revisable(target)?;
1118 let outcome = self.apply_remember(&input, target, None)?;
1119 let valid_from = input.valid_from.unwrap_or(input.now);
1120 self.close_target(target, valid_from);
1121 self.journal_remember(store, &input, target, outcome.id)?;
1122 Ok(outcome)
1123 }
1124
1125 pub fn forget<S: Storage>(
1129 &mut self,
1130 store: &mut S,
1131 now: u64,
1132 id: FactId,
1133 ) -> Result<bool, Error> {
1134 let fresh = self.apply_forget(id)?;
1135 let mut entry = Vec::new();
1136 Op::Forget { now, fact: id }.encode(&mut entry);
1137 store
1138 .append_journal(&entry)
1139 .map_err(|e| Error::Storage(format!("{e:?}")))?;
1140 Ok(fresh)
1141 }
1142
1143 pub fn remove_tag<S: Storage>(
1151 &mut self,
1152 store: &mut S,
1153 now: u64,
1154 tag: &str,
1155 ) -> Result<RemoveTagReport, Error> {
1156 if tag.is_empty() {
1157 return Err(Error::Invalid("empty tag"));
1158 }
1159 let report = self.apply_remove_tag(now, tag)?;
1160 if report.affected != 0 {
1161 let mut entry = Vec::new();
1162 Op::RemoveTag { now, tag }.encode(&mut entry);
1163 store
1164 .append_journal(&entry)
1165 .map_err(|e| Error::Storage(format!("{e:?}")))?;
1166 }
1167 Ok(report)
1168 }
1169
1170 pub fn link<S: Storage>(&mut self, store: &mut S, input: LinkInput<'_>) -> Result<(), Error> {
1173 self.apply_link(
1174 input.now,
1175 input.src,
1176 input.rel,
1177 input.dst,
1178 FactId::from_opt(input.provenance),
1179 )?;
1180 let mut entry = Vec::new();
1181 Op::Link {
1182 now: input.now,
1183 src: input.src,
1184 rel: input.rel,
1185 dst: input.dst,
1186 provenance: FactId::from_opt(input.provenance),
1187 }
1188 .encode(&mut entry);
1189 store
1190 .append_journal(&entry)
1191 .map_err(|e| Error::Storage(format!("{e:?}")))?;
1192 Ok(())
1193 }
1194
1195 pub fn unlink<S: Storage>(
1198 &mut self,
1199 store: &mut S,
1200 input: UnlinkInput<'_>,
1201 ) -> Result<bool, Error> {
1202 let fresh = self.apply_unlink(input.now, input.src, input.rel, input.dst)?;
1203 let mut entry = Vec::new();
1204 Op::Unlink {
1205 now: input.now,
1206 src: input.src,
1207 rel: input.rel,
1208 dst: input.dst,
1209 }
1210 .encode(&mut entry);
1211 store
1212 .append_journal(&entry)
1213 .map_err(|e| Error::Storage(format!("{e:?}")))?;
1214 Ok(fresh)
1215 }
1216
1217 pub fn get(&self, id: FactId) -> Option<FactView<'_>> {
1220 let record = self.fact(id)?;
1221 if record.is_tombstone() {
1222 return None;
1223 }
1224 let text = core::str::from_utf8(self.texts.get(record.text)).ok()?;
1228 Some(FactView { record, text })
1229 }
1230
1231 pub fn tags_of(&self, id: FactId, out: &mut Vec<TermId>) {
1234 let Some(record) = self.fact(id) else { return };
1235 if record.is_tombstone() {
1236 return;
1237 }
1238 let Some(aux) = self.fact_aux.get(&id.0.to_be_bytes()) else {
1239 return;
1240 };
1241 for chunk in self.tag_lists.iter(&aux.tags) {
1242 for raw in chunk.chunks_exact(4) {
1243 out.push(TermId(u32::from_be_bytes(raw.try_into().unwrap())));
1244 }
1245 }
1246 }
1247
1248 pub fn list_tags(&self, query: TagQuery<'_>) -> Result<TagPage, Error> {
1255 self.tag_catalog.page(&self.terms, self.cfg.db_uuid, query)
1256 }
1257
1258 pub fn metadata_of<'s>(&'s self, id: FactId, out: &mut Vec<(&'s str, &'s str)>) -> bool {
1269 out.clear();
1270 let Some(record) = self.fact(id) else {
1271 return false;
1272 };
1273 if record.is_tombstone() {
1274 return false;
1275 }
1276 let Some(aux) = self.fact_aux.get(&id.0.to_be_bytes()) else {
1277 return false;
1278 };
1279 if aux.meta.0 == NONE_U32 || aux.meta.0 >= self.metas.len() as u32 {
1280 return false;
1281 }
1282 crate::metadata::decode(self.metas.get(aux.meta), out).is_ok() && !out.is_empty()
1283 }
1284
1285 pub fn entity(&mut self, name: &str) -> Option<EntityId> {
1287 let mut norm = core::mem::take(&mut self.name_scratch);
1288 normalize_name(&mut self.tokenizer, name, &mut norm);
1289 let found = if norm.is_empty() {
1290 None
1291 } else {
1292 self.lookup_entity_by_norm(&norm)
1294 };
1295 self.name_scratch = norm;
1296 found
1297 }
1298
1299 pub fn term(&self, id: TermId) -> &str {
1301 self.terms.resolve(id)
1302 }
1303
1304 pub fn entity_name(&self, id: EntityId) -> Option<&str> {
1309 let record = self.entities.get(&id.0.to_be_bytes())?;
1310 core::str::from_utf8(self.texts.get(record.name)).ok()
1313 }
1314
1315 pub fn edges_each(&self, mut visit: impl FnMut(&str, &str, &str, FactId) -> bool) {
1333 for slot in self.edges_out.iter() {
1334 let (Some(src), Some(dst)) = (self.entity_name(slot.a), self.entity_name(slot.b))
1335 else {
1336 continue;
1337 };
1338 if !visit(src, self.term(slot.rel), dst, slot.fact) {
1339 return;
1340 }
1341 }
1342 }
1343
1344 pub fn facts_len(&self) -> usize {
1348 self.facts.len()
1349 }
1350
1351 pub fn entities_len(&self) -> usize {
1353 self.entities.len()
1354 }
1355
1356 pub fn cfg(&self) -> &Config {
1358 &self.cfg
1359 }
1360
1361 pub fn stats(&self) -> Stats {
1363 Stats {
1364 facts: self.facts.len(),
1365 entities: self.entities.len(),
1366 terms: self.terms.len(),
1367 edges: self.edges_out.len(),
1368 edge_versions: self.edges_hist_out.len(),
1369 vectors: self.vecs.len(),
1370 tombstones: self.tombstones,
1371 hnsw_indexed: self.hnsw.indexed(),
1372 next_fact: self.next_fact,
1373 next_entity: self.next_entity,
1374 next_edge: self.next_edge,
1375 db_uuid: self.cfg.db_uuid,
1376 pool_bytes: self.facts.pool_bytes()
1377 + self.fact_aux.pool_bytes()
1378 + self.entities.pool_bytes()
1379 + self.by_name.pool_bytes()
1380 + self.edges_out.pool_bytes()
1381 + self.edges_in.pool_bytes()
1382 + self.edges_hist_out.pool_bytes()
1383 + self.edges_hist_in.pool_bytes()
1384 + self.temporal.pool_bytes()
1385 + self.texts.pool_bytes()
1386 + self.terms.pool_bytes()
1387 + self.tag_lists.pool_bytes()
1388 + self.bm25.pool_bytes()
1389 + self.tags_idx.pool_bytes()
1390 + self.tag_catalog.pool_bytes()
1391 + self.entity_facts.pool_bytes()
1392 + self.vecs.pool_bytes()
1393 + self.hnsw.pool_bytes(),
1394 shards: shards::ShardLayout::of_config(&self.cfg),
1395 }
1396 }
1397
1398 fn fact(&self, id: FactId) -> Option<FactRecord> {
1401 self.facts.get(&id.0.to_be_bytes())
1402 }
1403
1404 fn validate_input(&self, input: &RememberInput<'_>) -> Result<(), Error> {
1406 if input.text.len() > self.cfg.max_text {
1407 return Err(Error::TooLarge {
1408 what: "text",
1409 len: input.text.len(),
1410 max: self.cfg.max_text,
1411 });
1412 }
1413 if input.tags.len() > 32 {
1414 return Err(Error::TooLarge {
1415 what: "tags",
1416 len: input.tags.len(),
1417 max: 32,
1418 });
1419 }
1420 if input.links.len() > 16 {
1421 return Err(Error::TooLarge {
1422 what: "links",
1423 len: input.links.len(),
1424 max: 16,
1425 });
1426 }
1427 if input.tags.iter().any(|t| t.is_empty()) {
1428 return Err(Error::Invalid("empty tag"));
1429 }
1430 if !input.links.is_empty() && input.entity.is_none() {
1431 return Err(Error::Invalid("links require a subject entity"));
1432 }
1433 if let Some(v) = input.vector {
1434 if self.cfg.dim == 0 {
1435 return Err(Error::Invalid("vector given but dim is 0"));
1436 }
1437 if v.len() != self.cfg.dim {
1438 return Err(Error::DimMismatch {
1439 got: v.len(),
1440 want: self.cfg.dim,
1441 });
1442 }
1443 }
1444 Ok(())
1445 }
1446
1447 fn check_revisable(&self, target: FactId) -> Result<(), Error> {
1449 let record = self.fact(target).ok_or(Error::NotFound(target))?;
1450 if record.is_tombstone() {
1451 return Err(Error::NotFound(target));
1452 }
1453 if record.is_closed() {
1454 return Err(Error::AlreadyClosed(target));
1455 }
1456 Ok(())
1457 }
1458
1459 fn close_target(&mut self, target: FactId, valid_to: u64) {
1462 let record = self.fact(target).expect("checked revisable");
1463 let payload = self
1464 .facts
1465 .payload_mut(&target.0.to_be_bytes())
1466 .expect("record fetched above");
1467 let flags = record.flags | fact_flags::CLOSED;
1470 payload[4..6].copy_from_slice(&flags.to_be_bytes());
1471 payload[36..44].copy_from_slice(&valid_to.to_be_bytes());
1472 self.change_catalog_for_fact(target, -1);
1473 }
1474
1475 fn apply_remember(
1478 &mut self,
1479 input: &RememberInput<'_>,
1480 revises: FactId,
1481 copy_vector: Option<u32>,
1482 ) -> Result<RememberOutcome, Error> {
1483 let id = FactId(self.next_fact);
1484 let entity = match input.entity {
1485 Some(name) => Some(self.resolve_or_create_entity(name, input.now)?),
1486 None => None,
1487 };
1488 let text_id = self.texts.push(input.text.as_bytes())?;
1489
1490 let mut tfs = core::mem::take(&mut self.tf_scratch);
1492 tfs.clear();
1493 let terms = &mut self.terms;
1494 let mut intern_err = None;
1495 self.tokenizer.tokenize(input.text, &mut |token| {
1496 if intern_err.is_some() {
1497 return;
1498 }
1499 match terms.intern(token) {
1500 Ok(term) => match tfs.iter_mut().find(|(t, _)| *t == term.0) {
1501 Some((_, tf)) => *tf = tf.saturating_add(1),
1502 None => tfs.push((term.0, 1)),
1503 },
1504 Err(e) => intern_err = Some(e),
1505 }
1506 });
1507 if let Some(e) = intern_err {
1508 self.tf_scratch = tfs;
1509 return Err(Error::Arena(e));
1510 }
1511 self.bm25.index_doc(id, &tfs)?;
1512 self.tf_scratch = tfs;
1513
1514 let meta = match input.metadata {
1518 Some(pairs) if !pairs.is_empty() => {
1519 self.metas.push(&crate::metadata::encode(pairs)?)?
1520 }
1521 _ => BlobId(NONE_U32),
1522 };
1523
1524 let mut aux = FactAux {
1527 id,
1528 tags: ListHandle::EMPTY,
1529 meta,
1530 };
1531 let mut seen_tags: [u32; 32] = [NONE_U32; 32];
1532 let mut seen_cnt = 0usize;
1533 for tag in input.tags {
1534 let term = self.terms.intern(tag)?;
1535 if seen_tags[..seen_cnt].contains(&term.0) {
1536 continue;
1537 }
1538 seen_tags[seen_cnt] = term.0;
1539 seen_cnt += 1;
1540 self.tag_lists.push(&mut aux.tags, &term.0.to_be_bytes())?;
1541 self.tags_idx.push(term.0, id, 0)?;
1542 }
1543 self.fact_aux.insert(&aux)?;
1544
1545 if let Some(src) = entity {
1547 for &(rel, dst_name) in input.links {
1548 let dst = self.resolve_or_create_entity(dst_name, input.now)?;
1549 let rel = self.terms.intern(rel)?;
1550 self.open_edge(input.now, src, rel, dst, id)?;
1551 }
1552 self.entity_facts.push(src.0, id, 0)?;
1553 }
1554
1555 let (vector, flags) = match (input.vector, copy_vector) {
1559 (Some(v), None) => (self.vecs.push(id, v)?, fact_flags::HAS_VECTOR),
1560 (None, Some(source)) => (
1561 self.vecs.clone_slot_for_fact(id, source)?,
1562 fact_flags::HAS_VECTOR,
1563 ),
1564 (None, None) => (NONE_U32, 0),
1565 (Some(_), Some(_)) => unreachable!("retag does not provide a raw vector"),
1566 };
1567
1568 let recorded_at = input.now;
1569 let valid_from = input.valid_from.unwrap_or(input.now);
1570 self.facts.insert(&FactRecord {
1571 id,
1572 entity: EntityId::from_opt(entity),
1573 flags,
1574 kind: 0,
1575 text: text_id,
1576 vector,
1577 revises,
1578 recorded_at,
1579 valid_from,
1580 valid_to: VALID_TO_OPEN,
1581 })?;
1582 self.temporal.insert(&TemporalSlot {
1583 recorded_at,
1584 fact: id,
1585 })?;
1586 for &term in &seen_tags[..seen_cnt] {
1587 self.tag_catalog.change(&self.terms, TermId(term), 1);
1588 }
1589 self.next_fact += 1;
1590 Ok(RememberOutcome {
1591 id,
1592 entity,
1593 similar: Vec::new(),
1594 })
1595 }
1596
1597 fn apply_forget(&mut self, id: FactId) -> Result<bool, Error> {
1598 let record = self.fact(id).ok_or(Error::NotFound(id))?;
1599 if record.is_tombstone() {
1600 return Ok(false);
1601 }
1602 let payload = self
1603 .facts
1604 .payload_mut(&id.0.to_be_bytes())
1605 .expect("record fetched above");
1606 let flags = record.flags | fact_flags::TOMBSTONE;
1607 payload[4..6].copy_from_slice(&flags.to_be_bytes());
1608 self.tombstones += 1;
1609 if !record.is_closed() {
1610 self.change_catalog_for_fact(id, -1);
1611 }
1612 Ok(true)
1613 }
1614
1615 fn apply_remove_tag(&mut self, now: u64, tag: &str) -> Result<RemoveTagReport, Error> {
1616 let Some(term) = self.terms.lookup(tag) else {
1617 return Ok(RemoveTagReport::default());
1618 };
1619 let targets: Vec<FactId> = self
1620 .tags_idx
1621 .entries(term.0)
1622 .map(|(id, _)| id)
1623 .filter(|&id| {
1624 self.fact(id)
1625 .is_some_and(|record| !record.is_tombstone() && !record.is_closed())
1626 })
1627 .collect();
1628 let mut affected = 0u32;
1629 for target in targets {
1630 self.retag_without(now, target, tag)?;
1631 affected = affected.saturating_add(1);
1632 }
1633 Ok(RemoveTagReport { affected })
1634 }
1635
1636 fn retag_without(&mut self, now: u64, target: FactId, removed: &str) -> Result<(), Error> {
1637 let record = self.fact(target).ok_or(Error::NotFound(target))?;
1638 let view = self.get(target).ok_or(Error::NotFound(target))?;
1639 let text = view.text.to_string();
1640 let entity = record
1641 .entity
1642 .some()
1643 .and_then(|id| self.entity_name(id))
1644 .map(ToString::to_string);
1645
1646 let mut tag_terms = Vec::new();
1647 self.tags_of(target, &mut tag_terms);
1648 let tags: Vec<String> = tag_terms
1649 .into_iter()
1650 .map(|term| self.term(term))
1651 .filter(|name| *name != removed)
1652 .map(ToString::to_string)
1653 .collect();
1654 let tag_refs: Vec<&str> = tags.iter().map(String::as_str).collect();
1655
1656 let mut metadata = Vec::new();
1657 self.metadata_of(target, &mut metadata);
1658 let metadata: Vec<(String, String)> = metadata
1659 .into_iter()
1660 .map(|(key, value)| (key.to_string(), value.to_string()))
1661 .collect();
1662 let metadata_refs: Vec<(&str, &str)> = metadata
1663 .iter()
1664 .map(|(key, value)| (key.as_str(), value.as_str()))
1665 .collect();
1666 let vector = (record.flags & fact_flags::HAS_VECTOR != 0).then_some(record.vector);
1667 let input = RememberInput {
1668 now,
1669 text: &text,
1670 entity: entity.as_deref(),
1671 tags: &tag_refs,
1672 links: &[],
1673 vector: None,
1674 valid_from: Some(now),
1675 metadata: (!metadata_refs.is_empty()).then_some(metadata_refs.as_slice()),
1676 };
1677 self.apply_remember(&input, target, vector)?;
1678 self.close_target(target, now);
1679 Ok(())
1680 }
1681
1682 fn change_catalog_for_fact(&mut self, id: FactId, delta: i32) {
1686 let Some(aux) = self.fact_aux.get(&id.0.to_be_bytes()) else {
1687 return;
1688 };
1689 let mut ids = [NONE_U32; 32];
1690 let mut len = 0usize;
1691 for chunk in self.tag_lists.iter(&aux.tags) {
1692 for raw in chunk.chunks_exact(4) {
1693 if len == ids.len() {
1694 break;
1695 }
1696 ids[len] = u32::from_be_bytes(raw.try_into().unwrap());
1697 len += 1;
1698 }
1699 }
1700 for &term in &ids[..len] {
1701 self.tag_catalog.change(&self.terms, TermId(term), delta);
1702 }
1703 }
1704
1705 fn apply_link(
1706 &mut self,
1707 now: u64,
1708 src: &str,
1709 rel: &str,
1710 dst: &str,
1711 provenance: FactId,
1712 ) -> Result<(), Error> {
1713 let src = self.resolve_or_create_entity(src, now)?;
1714 let dst = self.resolve_or_create_entity(dst, now)?;
1715 let rel = self.terms.intern(rel)?;
1716 self.open_edge(now, src, rel, dst, provenance)
1717 }
1718
1719 fn apply_unlink(&mut self, now: u64, src: &str, rel: &str, dst: &str) -> Result<bool, Error> {
1720 let Some(src) = self.lookup_entity_name(src) else {
1721 return Ok(false);
1722 };
1723 let Some(dst) = self.lookup_entity_name(dst) else {
1724 return Ok(false);
1725 };
1726 let Some(rel) = self.terms.lookup(rel) else {
1727 return Ok(false);
1728 };
1729 self.close_current_edge(now, src, rel, dst)
1730 }
1731
1732 fn open_edge(
1733 &mut self,
1734 now: u64,
1735 src: EntityId,
1736 rel: TermId,
1737 dst: EntityId,
1738 fact: FactId,
1739 ) -> Result<(), Error> {
1740 if let Some(current) = self.current_edge(src, rel, dst) {
1741 if current.fact == fact {
1742 return Ok(());
1743 }
1744 self.close_current_edge(now, src, rel, dst)?;
1745 }
1746 let edge = EdgeId(self.next_edge);
1747 let history = EdgeHistorySlot {
1748 a: src,
1749 rel,
1750 b: dst,
1751 edge,
1752 fact,
1753 flags: 0,
1754 kind: 0,
1755 recorded_at: now,
1756 valid_from: now,
1757 valid_to: VALID_TO_OPEN,
1758 };
1759 self.insert_history_edge(history)?;
1760 self.insert_current_edge(src, rel, dst, fact, edge, now)?;
1761 self.next_edge += 1;
1762 Ok(())
1763 }
1764
1765 fn insert_current_edge(
1766 &mut self,
1767 src: EntityId,
1768 rel: TermId,
1769 dst: EntityId,
1770 fact: FactId,
1771 edge: EdgeId,
1772 valid_from: u64,
1773 ) -> Result<(), Error> {
1774 for (arena, a, b) in [
1775 (&mut self.edges_out, src, dst),
1776 (&mut self.edges_in, dst, src),
1777 ] {
1778 let slot = EdgeSlot {
1779 a,
1780 rel,
1781 b,
1782 fact,
1783 edge,
1784 valid_from,
1785 };
1786 if !arena.insert(&slot)? {
1787 let payload = arena
1788 .payload_mut(&edge_key(a, rel, b))
1789 .expect("insert reported a duplicate");
1790 let mut full = [0u8; EdgeSlot::SIZE];
1791 slot.write(&mut full);
1792 payload.copy_from_slice(&full[EdgeSlot::KEY_LEN..]);
1793 }
1794 }
1795 Ok(())
1796 }
1797
1798 fn insert_history_edge(&mut self, edge: EdgeHistorySlot) -> Result<(), Error> {
1799 self.edges_hist_out.insert(&edge)?;
1800 self.edges_hist_in.insert(&EdgeHistorySlot {
1801 a: edge.b,
1802 b: edge.a,
1803 ..edge
1804 })?;
1805 Ok(())
1806 }
1807
1808 fn close_current_edge(
1814 &mut self,
1815 now: u64,
1816 src: EntityId,
1817 rel: TermId,
1818 dst: EntityId,
1819 ) -> Result<bool, Error> {
1820 let Some(current) = self.current_edge(src, rel, dst) else {
1821 return Ok(false);
1822 };
1823 let close_at = now.max(current.valid_from);
1826 let out_key = edge_history_key(src, current.valid_from, current.edge);
1827 let in_key = edge_history_key(dst, current.valid_from, current.edge);
1828 close_edge_history_payload(
1829 self.edges_hist_out
1830 .payload_mut(&out_key)
1831 .ok_or(Error::Corrupt("missing outgoing edge history"))?,
1832 close_at,
1833 );
1834 close_edge_history_payload(
1835 self.edges_hist_in
1836 .payload_mut(&in_key)
1837 .ok_or(Error::Corrupt("missing incoming edge history"))?,
1838 close_at,
1839 );
1840 let out_removed = self.edges_out.remove(&edge_key(src, rel, dst));
1841 let in_removed = self.edges_in.remove(&edge_key(dst, rel, src));
1842 if out_removed != in_removed {
1843 return Err(Error::Corrupt("edge mirrors disagree"));
1844 }
1845 Ok(out_removed)
1846 }
1847
1848 fn current_edge(&self, src: EntityId, rel: TermId, dst: EntityId) -> Option<EdgeSlot> {
1849 self.edges_out
1850 .get_slot(&edge_key(src, rel, dst))
1851 .map(EdgeSlot::read)
1852 }
1853
1854 fn lookup_entity_by_norm(&self, norm: &str) -> Option<EntityId> {
1857 let term = self.terms.lookup(norm)?;
1858 let mut from = [0u8; 8];
1859 key::write_u32(&mut from, term.0);
1860 let mut to = [0u8; 8];
1861 key::write_u32(&mut to, term.0);
1862 to[4..].copy_from_slice(&u32::MAX.to_be_bytes());
1863 self.by_name.range(&from, &to).next().map(|e| e.id)
1864 }
1865
1866 fn lookup_entity_name(&mut self, name: &str) -> Option<EntityId> {
1867 let mut norm = core::mem::take(&mut self.name_scratch);
1868 normalize_name(&mut self.tokenizer, name, &mut norm);
1869 let result = (!norm.is_empty())
1870 .then(|| self.lookup_entity_by_norm(&norm))
1871 .flatten();
1872 self.name_scratch = norm;
1873 result
1874 }
1875
1876 fn resolve_or_create_entity(&mut self, name: &str, now: u64) -> Result<EntityId, Error> {
1877 let mut norm = core::mem::take(&mut self.name_scratch);
1878 normalize_name(&mut self.tokenizer, name, &mut norm);
1879 if norm.is_empty() {
1880 self.name_scratch = norm;
1881 return Err(Error::Invalid("entity name has no indexable characters"));
1882 }
1883 let result = (|| {
1884 if let Some(found) = self.lookup_entity_by_norm(&norm) {
1885 return Ok(found);
1886 }
1887 let term = self.terms.intern(&norm)?;
1888 let id = EntityId(self.next_entity);
1889 let name_id = self.texts.push(name.as_bytes())?;
1890 self.entities.insert(&EntityRecord {
1891 id,
1892 name: name_id,
1893 name_term: term,
1894 created_at: now,
1895 flags: 0,
1896 })?;
1897 self.by_name.insert(&EntityByName {
1898 name_term: term,
1899 id,
1900 })?;
1901 self.next_entity += 1;
1902 Ok(id)
1903 })();
1904 self.name_scratch = norm;
1905 result
1906 }
1907
1908 fn journal_remember<S: Storage>(
1909 &mut self,
1910 store: &mut S,
1911 input: &RememberInput<'_>,
1912 revises: FactId,
1913 assigned: FactId,
1914 ) -> Result<(), Error> {
1915 let mut entry = Vec::new();
1916 Op::Remember {
1917 now: input.now,
1918 valid_from: input.valid_from.unwrap_or(input.now),
1919 entity: input.entity,
1920 text: input.text,
1921 tags: input.tags.to_vec(),
1922 links: input.links.to_vec(),
1923 vector: input.vector.map(<[f32]>::to_vec).unwrap_or_default(),
1924 metadata: input.metadata.map(<[_]>::to_vec).unwrap_or_default(),
1925 revises,
1926 assigned,
1927 }
1928 .encode(&mut entry);
1929 store
1930 .append_journal(&entry)
1931 .map_err(|e| Error::Storage(format!("{e:?}")))
1932 }
1933}
1934
1935impl core::fmt::Debug for Memory<'_> {
1936 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
1939 f.debug_struct("Memory")
1940 .field("facts", &self.facts.len())
1941 .field("entities", &self.entities.len())
1942 .field("terms", &self.terms.len())
1943 .finish()
1944 }
1945}
1946
1947fn normalize_name(tokenizer: &mut Tokenizer, name: &str, out: &mut String) {
1951 out.clear();
1952 tokenizer.tokenize(name, &mut |token| {
1953 if !out.is_empty() {
1954 out.push(' ');
1955 }
1956 out.push_str(token);
1957 });
1958}