Danchi17 commited on
Commit
46d778f
Β·
verified Β·
1 Parent(s): 0a51a38

Upload mnemo.py with huggingface_hub

Browse files
Files changed (1) hide show
  1. mnemo.py +1165 -0
mnemo.py ADDED
@@ -0,0 +1,1165 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """
2
+ mnemo β€” a memory layer for AI agents. (brand: Mnemosyne)
3
+
4
+ The memory that runs an autonomous research OS over ~5,800 notes, distilled to a single file with
5
+ no required dependencies. It does the four things agent memory actually needs, the way that held up
6
+ in production:
7
+
8
+ remember(text) append-only raw capture, stamped with an ABSOLUTE time (never rewritten)
9
+ recall(query, k) value-ranked retrieval: relevance Γ— the memory's accrued value, not just
10
+ cosine similarity β€” the high-value memories surface first
11
+ consolidate(cap) the "dream" pass: value-rank under a keep-budget, link near-duplicates, mark
12
+ stale/superseded β€” it only ADDS a derived layer, it never edits the raw note
13
+ contradictions() flag mutually-incompatible memories for REVIEW (never auto-delete)
14
+
15
+ Design rules that are not optional (each one cost us to learn):
16
+ β€’ Raw capture is immutable. Consolidation adds links/markers; it never overwrites the source β€”
17
+ that is what stops the slow accuracy drift of LLM-rewritten memory.
18
+ β€’ Absolute timestamps at write time. Relative/derived times rot the moment they're consolidated.
19
+ β€’ Value-ranked, capacity-aware consolidation. The payoff from ranking *what to keep* scales
20
+ super-linearly as the budget shrinks (measured), so retention tracks value, not recency β€” and
21
+ NOT access-frequency: decaying on reads keeps *popular* memories, but popularity != value, so a
22
+ pure access-reset policy starves the rarely-read-but-load-bearing fact (measured: it retains
23
+ ~3x less total value than a value blend under a tight budget). Forgetting blends value + recency.
24
+ β€’ Report value at the COHORT level (tag / time-block), never per-memory: per-item value at n-of-1
25
+ is statistical noise; cohorts are where the signal lives.
26
+ β€’ Contradictions are flagged for review, not auto-resolved. Silent rewrites destroy trust.
27
+
28
+ Bring your own embedder for semantic recall (any text->vector fn); with none, mnemo falls back to a
29
+ lexical token overlap so it runs anywhere, today.
30
+
31
+ from mnemo import Mnemo
32
+ m = Mnemo("memory.json") # or Mnemo("memory.json", embed=my_embedder)
33
+ m.remember("Pre-trend tests catch only ~31% of fatal DiD bias.", tags=["causal"], value=3)
34
+ m.recall("difference in differences", k=5)
35
+ m.consolidate(keep=200)
36
+ m.contradictions()
37
+
38
+ MIT-licensed. Part of Agora (https://github.com/DanceNitra/agora).
39
+ """
40
+ from __future__ import annotations
41
+
42
+ import hashlib
43
+ import json
44
+ import math
45
+ import os
46
+ import re
47
+ import time
48
+ import uuid
49
+ from pathlib import Path
50
+
51
+ try: # OPTIONAL: numpy only ACCELERATES semantic recall at scale.
52
+ import numpy as _np # mnemo still runs (pure-Python cosine) with no numpy installed.
53
+ except Exception:
54
+ _np = None
55
+
56
+ try: # OPTIONAL: only needed to SIGN write receipts (see receipts=...).
57
+ from cryptography.hazmat.primitives.asymmetric.ed25519 import (
58
+ Ed25519PrivateKey as _Ed25519SK, Ed25519PublicKey as _Ed25519PK)
59
+ from cryptography.hazmat.primitives import serialization as _ser
60
+ _HAVE_ED = True
61
+ except Exception:
62
+ _HAVE_ED = False
63
+
64
+ _GENESIS = "0" * 64
65
+
66
+
67
+ def _canon(obj) -> bytes:
68
+ return json.dumps(obj, sort_keys=True, separators=(",", ":"), ensure_ascii=False).encode("utf-8")
69
+
70
+
71
+ def _sha256_hex(b: bytes) -> str:
72
+ return hashlib.sha256(b).hexdigest()
73
+
74
+
75
+ def new_receipt_keypair():
76
+ """Return (private_key_hex, public_key_hex) for signing mnemo write receipts. Needs `cryptography`."""
77
+ if not _HAVE_ED:
78
+ raise RuntimeError("signing write receipts needs the `cryptography` package (pip install cryptography)")
79
+ sk = _Ed25519SK.generate()
80
+ return (sk.private_bytes(_ser.Encoding.Raw, _ser.PrivateFormat.Raw, _ser.NoEncryption()).hex(),
81
+ sk.public_key().public_bytes(_ser.Encoding.Raw, _ser.PublicFormat.Raw).hex())
82
+
83
+ __version__ = "0.4.0"
84
+ _WORD = re.compile(r"[a-z0-9][a-z0-9\-']{2,}")
85
+ _STOP = frozenset("the a an of for to in on and or is are was were be been with this that it its as "
86
+ "by at from into our we us you your he she they them his her their not no".split())
87
+
88
+
89
+ def _stem(w: str) -> str:
90
+ return w[:-1] if (w.endswith("s") and len(w) > 4) else w # crude plural/3rd-person fold
91
+
92
+
93
+ def _tokens(text: str) -> set:
94
+ return {_stem(w) for w in _WORD.findall((text or "").lower()) if w not in _STOP}
95
+
96
+
97
+ def _token_counts(text: str) -> dict:
98
+ """Term-frequency map with the SAME tokenization as _tokens (stem + stopword filter). BM25 needs TF;
99
+ _tokens loses it by returning a set."""
100
+ d: dict = {}
101
+ for w in _WORD.findall((text or "").lower()):
102
+ if w in _STOP:
103
+ continue
104
+ s = _stem(w)
105
+ d[s] = d.get(s, 0) + 1
106
+ return d
107
+
108
+
109
+ def _cosine(a, b) -> float:
110
+ dot = sum(x * y for x, y in zip(a, b))
111
+ na = math.sqrt(sum(x * x for x in a)) or 1.0
112
+ nb = math.sqrt(sum(y * y for y in b)) or 1.0
113
+ return dot / (na * nb)
114
+
115
+
116
+ class Mnemo:
117
+ def __init__(self, path: str | None = None, embed=None, receipts: bool = False,
118
+ receipt_key: str | None = None, receipt_pubkey: str | None = None):
119
+ """path: optional JSON file to persist to. embed: optional fn(str)->list[float] for semantic
120
+ recall; if omitted, recall uses lexical token overlap (zero dependencies).
121
+
122
+ receipts/receipt_key (OPT-IN, default OFF -> identical legacy behavior): when enabled, every
123
+ remember() appends a tamper-evident, hash-chained WRITE RECEIPT committing to the memory's
124
+ content hash, persisted to a sidecar "<path>.receipts.json" (the main store format is unchanged).
125
+ verify_writes() then proves the write history wasn't altered out-of-band β€” something an
126
+ append-only store alone can't, because anyone who can edit the store file can rewrite a stored
127
+ memory and the store would serve the altered text as original. The hash chain is zero-dependency;
128
+ pass receipt_key (+ receipt_pubkey) from new_receipt_keypair() to also Ed25519-SIGN each receipt
129
+ so a third party can verify it with the public key only. (Standalone version: agora-agent-receipts.)"""
130
+ self.path = Path(path) if path else None
131
+ self.embed = embed
132
+ self.items: list[dict] = []
133
+ self._tok_cache: dict[str, set] = {} # id -> token set, so recall doesn't re-tokenize
134
+ self._tc_cache: dict[str, dict] = {} # id -> term-frequency map, for the BM25 hybrid channel
135
+ # recall auto-mode: below this many active memories lexical is as good and free; above it the
136
+ # lexical+semantic HYBRID (RRF) pays β€” measured to beat either channel alone on agent memory. Tunable.
137
+ self.semantic_threshold = 300
138
+ self._last_mode = "lexical" # which mode the most recent recall() actually used
139
+ self._mat = None # cached L2-normalized matrix of memory vectors (numpy)
140
+ self._vec_rowof: dict[str, int] = {} # memory id -> its row in self._mat
141
+ self._mat_built_n = -1 # item count when the matrix was built (rebuild on change)
142
+ self._vec_mean = None # corpus mean vector (for anisotropy centering of semantic recall)
143
+ # Anisotropy centering: subtract the corpus mean before cosine. Many embedders (e.g. nomic) are
144
+ # anisotropic β€” all cosines compress into a narrow band β€” so semantic recall under-separates.
145
+ # MEASURED on real LoCoMo (419 turns): centering lifts single-hop full-evidence recall@k by
146
+ # +0.04..+0.07 (k=5/10/20) and is neutral on multi-hop. Reversible: set center_embeddings=False.
147
+ self.center_embeddings = True
148
+ # Two-tier keep-budget: when consolidate(keep) must drop surplus, PROTECT the top protect_frac
149
+ # of the budget by RAW value (recency-immune) and fill the REST by EFFECTIVE (decay-weighted)
150
+ # value β€” so a freshly-useful memory isn't evicted by a stale high-value one. A pure top-N-by-raw
151
+ # prune keeps old high-value items forever and starves a drifting working set. MEASURED on a
152
+ # simulation of mnemo's own value-accrual + per-type decay: locality served-hit 0.22 -> 0.78,
153
+ # neutral on rare-critical + poison-flood. Reversible: two_tier_keep=False -> legacy top-N-by-raw.
154
+ self.two_tier_keep = True
155
+ self.protect_frac = 0.30
156
+ # Fast-novelty channel guard (OPT-IN, default OFF). mnemo's state-toggle supersedes a standing
157
+ # fact the moment a single similar+contradicting memory arrives β€” correct + fast for a TRUSTED
158
+ # single source (configs/preferences: latest assertion wins), but a single-shot poison flip
159
+ # (AgentPoison / MINJA) can then override a true fact. With this ON, a contradiction supersedes
160
+ # only when CORROBORATED (earned credit, or >=2 corroborating links) β€” the same bar as graduation;
161
+ # an uncorroborated single contradiction is recorded as a link but does NOT supersede. This is the
162
+ # two-channel capstone's latency-floor tradeoff made explicit: robustness to single-shot poison at
163
+ # the cost of lagging an uncorroborated single legitimate change. Leave OFF for trusted-source
164
+ # stores; turn ON for adversarial / multi-tenant ingestion.
165
+ self.supersede_requires_corroboration = False
166
+ # Persistence supersession (OPT-IN, default OFF; set to an int >= 2 to enable). A standing fact is
167
+ # superseded only when the contradicting NEW state is asserted by >= this many INDEPENDENT records β€”
168
+ # i.e. the change must PERSIST/accumulate, not arrive once. This is the sequential-change-detection
169
+ # (CUSUM) escape applied to memory: an isolated single-shot poison flip never crosses the threshold
170
+ # and is rejected, while a genuinely sustained value change is adopted once `supersede_persistence`
171
+ # corroborating records exist. The integer IS the Adaptation-Corruption law's detection-latency floor
172
+ # d* made explicit β€” set it to your stream's corruption-vs-change ratio. Unlike
173
+ # supersede_requires_corroboration this needs NO external credit(): it adopts a genuine change purely
174
+ # from repeated independent assertions, where the corroboration guard would lag one forever. MEASURED
175
+ # (lab fea933, mnemo's real consolidate() path): isolated-poison false-supersede 1 -> 0 while a
176
+ # 3-record sustained change is still adopted; it Pareto-dominates both the naive (poison-fooled) and
177
+ # corroboration-only (change-lagging) rules β€” see the Adaptation-Corruption Separation Law (lab f490d8).
178
+ # Reversible: 0 or 1 -> legacy fast supersession.
179
+ self.supersede_persistence = 0
180
+ # _save() THROTTLE: serializing the whole store (json.dumps of every item) is O(store size); doing
181
+ # it on EVERY recall/remember froze callers once the store grew (recall mutates access value, so it
182
+ # used to re-serialize everything each call). Coalesce disk writes to at most once / _save_min_s;
183
+ # at most _save_min_s of access-metadata is lost on a hard crash (working memory β€” acceptable).
184
+ self._save_min_s = 5.0
185
+ self._last_save = 0.0
186
+ self._dirty = False
187
+ if self.path and self.path.exists():
188
+ try:
189
+ self.items = json.loads(self.path.read_text(encoding="utf-8"))
190
+ except Exception:
191
+ self.items = []
192
+ # OPT-IN write receipts (default OFF -> zero behavior change; no sidecar created)
193
+ self.receipts_enabled = bool(receipts or receipt_key)
194
+ self._receipt_sk = receipt_key
195
+ self.receipt_pubkey = receipt_pubkey
196
+ self._receipts: list[dict] = []
197
+ self._receipts_path = (self.path.parent / (self.path.name + ".receipts.json")) if self.path else None
198
+ if self.receipts_enabled and self._receipts_path and self._receipts_path.exists():
199
+ try:
200
+ self._receipts = json.loads(self._receipts_path.read_text(encoding="utf-8"))
201
+ except Exception:
202
+ self._receipts = []
203
+
204
+ # ── capture ──────────────────────────────────────────────────────────────
205
+ def remember(self, text: str, tags=None, value: float = 1.0, meta: dict | None = None,
206
+ mtype: str | None = None, valid_from: float | None = None,
207
+ source: dict | None = None, key: str | None = None) -> str:
208
+ """Append-only raw capture. Stamped with an absolute UTC time; never edited afterward.
209
+ mtype in {episodic, semantic, procedural} sets the decay prior (episodic fades fast,
210
+ semantic slow, procedural barely); inferred from the text if not given. Pass it explicitly
211
+ when the caller knows the kind β€” inference defaults to episodic (the conservative, fast-decay
212
+ choice) and only promotes on clear markers.
213
+
214
+ `key` (OPT-IN) is a deterministic supersession key β€” typically a (subject, relation) identifier,
215
+ e.g. "billing-api::auth-method". When set, remembering a new value RETIRES every active record
216
+ sharing the same key (status -> superseded), with NO similarity threshold and NO LLM call. This
217
+ closes the 'supersession blind spot': cosine similarity cannot tell a contradicted fact from its
218
+ replacement (we measured AUROC ~0.61, near chance β€” a contradiction is often MORE embedding-similar
219
+ to the original than a rephrase is), so a similarity-based store silently serves the stale value
220
+ (~42% of the time in our test). A deterministic (subject, relation, object) ledger drives that to
221
+ ~0%. Bi-temporal: a back-filled record (earlier valid_from) does NOT overwrite a genuinely newer
222
+ same-key value β€” the stale-on-arrival record is the one retired."""
223
+ mid = uuid.uuid4().hex[:10]
224
+ now = time.time()
225
+ rec = {"id": mid, "text": text, "tags": list(tags or []), "value": float(value),
226
+ "ts": now, "iso": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),
227
+ "valid_from": float(valid_from) if valid_from is not None else now, # event-time (bi-temporal); defaults to ingest-time
228
+ "source": dict(source) if source else None, # re-checkable origin (e.g. {"doc": id, "span": [start, end]}) so a recalled fact can be traced back, not trusted blind
229
+ "mtype": mtype or _infer_type(text), "last_access": now,
230
+ "status": "active", "links": [], "meta": dict(meta or {})}
231
+ if key is not None:
232
+ rec["key"] = str(key)
233
+ if self.embed:
234
+ try:
235
+ rec["vec"] = list(self.embed(text))
236
+ except Exception:
237
+ rec["vec"] = None
238
+ self.items.append(rec)
239
+ if key is not None:
240
+ self._supersede_by_key(rec) # deterministic SRO supersession (no embedding, no threshold)
241
+ self._save(force=True) # a new memory is real content - persist immediately, not throttled
242
+ if self.receipts_enabled:
243
+ self._emit_write_receipt(rec)
244
+ return mid
245
+
246
+ # ── write receipts (OPT-IN: tamper-evident write history) ─────────────────
247
+ def _write_commit(self, rec: dict) -> dict:
248
+ """What a receipt commits to for a stored memory: its id + a hash of its content-bearing fields."""
249
+ return {"id": rec["id"],
250
+ "content_sha256": _sha256_hex(_canon({"text": rec.get("text"), "key": rec.get("key"),
251
+ "mtype": rec.get("mtype")}))}
252
+
253
+ def _emit_write_receipt(self, rec: dict) -> dict:
254
+ prev = self._receipts[-1]["hash"] if self._receipts else _GENESIS
255
+ r = {"seq": len(self._receipts), "ts": rec.get("ts"), "memory_id": rec["id"],
256
+ "commit": self._write_commit(rec), "prev": prev}
257
+ r["hash"] = _sha256_hex(_canon({k: r[k] for k in ("seq", "ts", "memory_id", "commit", "prev")}))
258
+ if self._receipt_sk and _HAVE_ED:
259
+ sk = _Ed25519SK.from_private_bytes(bytes.fromhex(self._receipt_sk))
260
+ r["pubkey"] = self.receipt_pubkey
261
+ r["sig"] = sk.sign(bytes.fromhex(r["hash"])).hex()
262
+ self._receipts.append(r)
263
+ if self._receipts_path:
264
+ try:
265
+ self._receipts_path.write_text(json.dumps(self._receipts, indent=2, ensure_ascii=False),
266
+ encoding="utf-8")
267
+ except Exception:
268
+ pass
269
+ return r
270
+
271
+ def verify_writes(self, expected_pubkey: str | None = None) -> tuple[bool, list[str]]:
272
+ """Verify the write-receipt chain AND that each stored memory still matches its write receipt.
273
+ Returns (ok, problems). Catches out-of-band edits to the store the normal flow can't see.
274
+ Requires receipts to have been enabled at write time."""
275
+ problems: list[str] = []
276
+ prev = _GENESIS
277
+ by_id = {it["id"]: it for it in self.items}
278
+ for i, r in enumerate(self._receipts):
279
+ core = {k: r.get(k) for k in ("seq", "ts", "memory_id", "commit", "prev")}
280
+ if r.get("prev") != prev:
281
+ problems.append(f"receipt {i}: broken chain link (a prior receipt was altered/removed)")
282
+ if _sha256_hex(_canon(core)) != r.get("hash"):
283
+ problems.append(f"receipt {i}: receipt tampered (hash mismatch)")
284
+ if "sig" in r and _HAVE_ED:
285
+ try:
286
+ _Ed25519PK.from_public_bytes(bytes.fromhex(r["pubkey"])).verify(
287
+ bytes.fromhex(r["sig"]), bytes.fromhex(r["hash"]))
288
+ if expected_pubkey and r.get("pubkey") != expected_pubkey:
289
+ problems.append(f"receipt {i}: signed by an unexpected key")
290
+ except Exception:
291
+ problems.append(f"receipt {i}: invalid signature")
292
+ elif expected_pubkey:
293
+ problems.append(f"receipt {i}: unsigned, but a signature was required")
294
+ cur = by_id.get(r["memory_id"])
295
+ if cur is None:
296
+ problems.append(f"memory {r['memory_id']}: written but missing from the store (deleted out-of-band)")
297
+ elif self._write_commit(cur) != r["commit"]:
298
+ problems.append(f"memory {r['memory_id']}: stored content no longer matches its write receipt (edited after write)")
299
+ prev = r.get("hash")
300
+ return (len(problems) == 0, problems)
301
+
302
+ def _supersede_by_key(self, rec: dict) -> None:
303
+ """Deterministic (subject, relation, object) supersession: retire active records that share
304
+ rec['key']. No similarity threshold, no LLM call β€” the fix our Crucible replication validated
305
+ (stale-fact recall 41.7% -> 0.0%, where cosine-based detection is near chance at AUROC ~0.61).
306
+ Bi-temporal: only same-key records with valid_from <= rec's are retired; if an active same-key
307
+ record is genuinely newer (later valid_from), the INCOMING rec is the stale one and is retired
308
+ instead β€” a back-filled value never overwrites the current one. recall() hides superseded records
309
+ by default, so a keyed store never surfaces a stale fact."""
310
+ k = rec.get("key")
311
+ if not k:
312
+ return
313
+ vf_new = rec.get("valid_from", rec["ts"])
314
+ for r in self.items:
315
+ if r is rec or r.get("status") != "active" or r.get("key") != k:
316
+ continue
317
+ vf_r = r.get("valid_from", r["ts"])
318
+ if vf_r <= vf_new: # r is the older value -> retire it
319
+ r["status"] = "superseded"
320
+ r["superseded_ts"] = time.time()
321
+ r["invalidated_at"] = vf_new
322
+ r.setdefault("meta", {})["superseded_by_toggle"] = rec["id"]
323
+ else: # an active same-key value is newer -> incoming is stale-on-arrival
324
+ rec["status"] = "superseded"
325
+ rec["superseded_ts"] = time.time()
326
+ rec["invalidated_at"] = vf_r
327
+ rec.setdefault("meta", {})["superseded_by_toggle"] = r["id"]
328
+
329
+ def remember_dedup(self, text: str, tags=None, value: float = 1.0, meta: dict | None = None,
330
+ mtype: str | None = None, dup_threshold: float = 0.95) -> str:
331
+ """OPT-IN write that skips redundant appends. If an active memory is near-identical (similarity >=
332
+ dup_threshold) AND carries the SAME value(s) (no numeric clash), this returns that memory's id WITHOUT
333
+ appending a duplicate raw row -- cutting raw-store bloat from repeated identical writes. A near-identical
334
+ text with a DIFFERENT number (a value UPDATE) is NOT a duplicate: it appends, so the consolidation pass can
335
+ supersede the stale value. Default `remember()` stays strictly append-only (the 'zero rewrites' contract);
336
+ this is a separate opt-in path for high-duplicate ingest."""
337
+ hits = self.recall(text, k=1)
338
+ if hits:
339
+ h = hits[0]
340
+ s = self._similarity(text, h, self._qvec(text) if self.embed else None)
341
+ if s >= dup_threshold and not _value_clash(text, h["text"]):
342
+ return h["id"] # NO-OP: near-identical, same value -> skip the redundant append
343
+ return self.remember(text, tags=tags, value=value, meta=meta, mtype=mtype)
344
+
345
+ def forget(self, ids=None, where=None, redact_links: bool = True) -> dict:
346
+ """HARD-DELETE memories β€” the one operation that genuinely REMOVES content. mnemo is otherwise
347
+ append-only: supersession / invalidation only DEMOTE a record (it still exists, recallable with
348
+ include_superseded). forget() is for the cases where demotion is not enough: a right-to-be-forgotten
349
+ / erasure request, a poisoned or libellous memory, or a hard correction.
350
+
351
+ Select by `ids` (a single id or an iterable) and/or `where` (a predicate fn(record)->bool; e.g.
352
+ lambda r: 'secret' in r['text']). VERIFIED FORGETTING: the matched records are deleted AND their ids
353
+ are scrubbed from every surviving record's `links` and toggle-supersession pointers, and the cached
354
+ vec matrix + token caches are dropped β€” so a forgotten memory cannot resurface via recall, via a
355
+ consolidation link, or via a stale derived-summary pointer. This is complete because consolidation
356
+ never copies raw text into other records (it only links ids and toggles status) β€” there is no merged
357
+ blob left holding the forgotten content. Returns {forgotten, ids, scrubbed_links}."""
358
+ target = set()
359
+ if ids is not None:
360
+ target |= ({ids} if isinstance(ids, str) else set(ids))
361
+ if where is not None:
362
+ for r in self.items:
363
+ try:
364
+ if where(r):
365
+ target.add(r["id"])
366
+ except Exception:
367
+ pass
368
+ target &= {r["id"] for r in self.items} # ignore ids not actually present
369
+ if not target:
370
+ return {"forgotten": 0, "ids": [], "scrubbed_links": 0}
371
+ self.items = [r for r in self.items if r["id"] not in target]
372
+ scrubbed = 0
373
+ if redact_links:
374
+ for r in self.items:
375
+ if r.get("links"):
376
+ before = len(r["links"])
377
+ r["links"] = [l for l in r["links"] if l not in target]
378
+ scrubbed += before - len(r["links"])
379
+ meta = r.get("meta")
380
+ if meta and meta.get("superseded_by_toggle") in target:
381
+ meta.pop("superseded_by_toggle", None) # drop dangling toggle pointer (no ghost stale-derived)
382
+ for tid in target:
383
+ self._tok_cache.pop(tid, None)
384
+ self._mat = None; self._mat_built_n = -1 # force vec-matrix rebuild (drops forgotten rows)
385
+ self._save(force=True) # a deletion is real content change β€” persist now
386
+ return {"forgotten": len(target), "ids": sorted(target), "scrubbed_links": scrubbed}
387
+
388
+ # ── retrieval (value-ranked) ──────────────────────────────────────────────
389
+ def _qvec(self, query: str):
390
+ """Embed a query ONCE per scan, or None (no embedder / failure). Callers pass the result
391
+ into _similarity so a recall over N memories costs 1 embedding, not N."""
392
+ if not self.embed:
393
+ return None
394
+ try:
395
+ return self.embed(query)
396
+ except Exception:
397
+ return None
398
+
399
+ def _vec_matrix(self):
400
+ """Cached L2-normalized matrix (numpy) of every memory that carries a vec β€” so a semantic
401
+ recall is ONE matmul, not an O(NΒ·d) pure-Python cosine loop. Rebuilt only when the item count
402
+ changes (remember / bulk load); status changes (consolidate) don't touch the vectors."""
403
+ if _np is None:
404
+ return None
405
+ if self._mat is None or self._mat_built_n != len(self.items):
406
+ rows, ids = [], []
407
+ for r in self.items:
408
+ if r.get("vec"):
409
+ rows.append(r["vec"]); ids.append(r["id"])
410
+ if rows:
411
+ M = _np.asarray(rows, dtype=_np.float32)
412
+ self._vec_mean = M.mean(axis=0) # corpus mean (computed regardless; used to center)
413
+ if self.center_embeddings:
414
+ M = M - self._vec_mean # de-anisotropise: remove the common component
415
+ M /= (_np.linalg.norm(M, axis=1, keepdims=True) + 1e-9)
416
+ self._mat = M
417
+ self._vec_rowof = {i: k for k, i in enumerate(ids)}
418
+ else:
419
+ self._mat, self._vec_rowof, self._vec_mean = None, {}, None
420
+ self._mat_built_n = len(self.items)
421
+ return self._mat
422
+
423
+ def _rec_tokens(self, rec: dict) -> set:
424
+ """Token set for a memory, cached by id β€” recall over N memories shouldn't re-tokenize."""
425
+ rid = rec.get("id") or id(rec)
426
+ t = self._tok_cache.get(rid)
427
+ if t is None:
428
+ t = _tokens(rec["text"]); self._tok_cache[rid] = t
429
+ return t
430
+
431
+ def _rec_tokcount(self, rec: dict) -> dict:
432
+ """Term-frequency map for a memory, cached by id (for the BM25 hybrid channel)."""
433
+ rid = rec.get("id") or id(rec)
434
+ c = self._tc_cache.get(rid)
435
+ if c is None:
436
+ c = _token_counts(rec["text"]); self._tc_cache[rid] = c
437
+ return c
438
+
439
+ def _bm25_scores(self, qtok: set, pool: list, k1: float = 1.5, b: float = 0.75) -> list:
440
+ """Okapi BM25 score of `query` (token set) against every record in `pool` β€” the strong lexical
441
+ channel for the hybrid. df/avgdl are computed over the pool (the live corpus). Returns a list of
442
+ scores aligned to `pool`. Pure-Python, zero-dependency. We MEASURED BM25 (not token-overlap) as the
443
+ lexical channel that makes the hybrid beat either alone (mnemo/probes/locomo_retrieval_map.py)."""
444
+ N = len(pool)
445
+ if N == 0:
446
+ return []
447
+ counts = [self._rec_tokcount(r) for r in pool]
448
+ dl = [sum(c.values()) for c in counts]
449
+ avgdl = (sum(dl) / N) or 1.0
450
+ df: dict = {}
451
+ for c in counts:
452
+ for t in c:
453
+ df[t] = df.get(t, 0) + 1
454
+ idf = {t: math.log(1 + (N - n + 0.5) / (n + 0.5)) for t, n in df.items()}
455
+ out = []
456
+ for c, L in zip(counts, dl):
457
+ s = 0.0
458
+ for t in qtok:
459
+ f = c.get(t, 0)
460
+ if f:
461
+ s += idf.get(t, 0.0) * (f * (k1 + 1)) / (f + k1 * (1 - b + b * L / avgdl))
462
+ out.append(s)
463
+ return out
464
+
465
+ def _similarity(self, query: str, rec: dict, qvec=None, qtok: set | None = None) -> float:
466
+ if qvec is not None and rec.get("vec"):
467
+ return max(0.0, _cosine(qvec, rec["vec"]))
468
+ q = qtok if qtok is not None else _tokens(query)
469
+ t = self._rec_tokens(rec)
470
+ if not q or not t:
471
+ return 0.0
472
+ return len(q & t) / min(len(q), len(t)) # overlap coefficient β€” forgiving without an embedder
473
+
474
+ def recall(self, query: str, k: int = 6, include_superseded: bool = False,
475
+ include_hubs: bool = False, mode: str = "auto", min_relevance: float = 0.0,
476
+ scope: str | None = None, as_of: float | None = None,
477
+ where: dict | None = None, influence_only: bool = False) -> list[dict]:
478
+ """Top-k memories by RELEVANCE Γ— VALUE β€” high-value memories outrank merely-similar ones.
479
+ Memories the dream pass flagged as hubs (universal matchers) are skipped unless include_hubs.
480
+
481
+ mode: 'auto' (default) uses LEXICAL token overlap while the store is small (< semantic_threshold
482
+ active memories) and a LEXICAL+SEMANTIC HYBRID (Reciprocal Rank Fusion) once it grows past that β€”
483
+ the hybrid robustly beat either channel alone in our agent-memory benchmark (details in the recall
484
+ body / mnemo/probes/locomo_retrieval_map.py). Force a single channel with mode='lexical' /
485
+ 'semantic', or the fusion explicitly with mode='hybrid'. Semantic/hybrid need an embedder (set on
486
+ the store); without one, or if embedding fails, recall falls back to lexical automatically.
487
+
488
+ where: an OPT-IN metadata pre-filter applied to the candidate pool BEFORE ranking β€” the cheap
489
+ 'filter before you rank' lever (measured on LoCoMo: a metadata pre-filter can beat retriever choice;
490
+ mnemo/probes/locomo_metadata_prefilter.py). A dict of field -> condition; a record must match ALL
491
+ fields (AND). Each field is matched against the record's top-level attributes first, then its meta
492
+ dict, so both `valid_from`/`mtype`/`key` and any `meta` key work. A condition is either a scalar
493
+ (equality), a list/tuple/set (membership), or a dict of operators:
494
+ {"$gte","$lte","$gt","$lt","$in","$nin","$ne","$contains"} β€” e.g. a time range
495
+ where={"valid_from": {"$gte": t0, "$lte": t1}} (hard-filter the SOLVED half), or a closed-set
496
+ entity where={"speaker": {"$in": ["Caroline","Mel"]}}. NOTE: this is a HARD filter β€” a record that
497
+ doesn't match is removed, so on lossy/predicted extraction prefer a broad/loose filter (or rerank)
498
+ over an aggressive one, since a wrong filter hard-deletes the answer (measured harm mode).
499
+
500
+ influence_only (OPT-IN, default False -> zero behavior change): restrict the result to CORROBORATED
501
+ memories β€” those that meet the same bar mnemo uses for episodic->semantic GRADUATION (an EARNED
502
+ net-positive outcome via credit() [good>0 and good>=bad], OR already-graduated 'semantic' type, OR
503
+ >=2 DISTINCT-canonical-source corroborating links). This is the retrieve-then-INFLUENCE split: recall
504
+ freely for context, but call with influence_only=True for the set that is allowed to DRIVE an action.
505
+ MEASURED (mnemo/probes/agentpoison_influence_gate*.py) against a real AgentPoison-style single-
506
+ instance retrieval-poisoning attack (Chen et al., NeurIPS 2024, arXiv:2407.12784; PoisonedRAG, Zou
507
+ et al., arXiv:2402.07867): a natural-sentence trigger hijacks RAW top-1 retrieval 88-100% and is
508
+ scale-invariant (60->10k memories), and retrieval-time / embedding-geometry defenses do NOT
509
+ generalize across encoders β€” but influence_only drops the single-instance poison's rank-1 hijack to
510
+ 0% on all three tested retrievers (MiniLM/BGE/Contriever) and all scales, because an injected poison
511
+ never earns corroboration while legitimate memories earn it through use. It GENERALIZES precisely
512
+ because it lives in provenance metadata, not embedding geometry. HONEST COST (calibration tradeoff):
513
+ a rare-but-true memory that has not yet earned corroboration is filtered too (measured recall 1.00
514
+ corroborated vs 0.08 uncorroborated) β€” so this is for adversarial / untrusted-ingestion use, where a
515
+ recalled-but-uncorroborated memory should inform but not unilaterally drive an action. It RAISES
516
+ attacker cost (a single free injection is filtered; defeating it needs >=3 coordinated records with
517
+ >=2 independent forged provenances) rather than making poisoning impossible. Reversible: default
518
+ False = legacy recall."""
519
+ def _eligible(r: dict) -> bool:
520
+ s = r["status"]
521
+ if as_of is not None:
522
+ # Bi-temporal "as of T": a memory counts if it was VALID at time T β€” valid_from <= T and not yet
523
+ # invalidated by T β€” INCLUDING records now superseded (they were current back then). Records
524
+ # superseded by the pre-bitemporal pass carry no invalidated_at; treat them as still-valid here.
525
+ vf = r.get("valid_from", r["ts"])
526
+ inv = r.get("invalidated_at")
527
+ if vf > as_of or (inv is not None and inv <= as_of):
528
+ return False
529
+ return include_hubs if s == "hub" else True
530
+ if s == "active":
531
+ return True
532
+ if s == "hub":
533
+ return include_hubs
534
+ return include_superseded # superseded / other non-active
535
+ pool = [r for r in self.items if _eligible(r)]
536
+ # Scope/namespace isolation: when a scope is requested, recall ONLY sees memories tagged with that scope
537
+ # (meta['scope']) BEFORE ranking β€” a shared store (e.g. many agents / tenants in one Mnemo) cannot bleed
538
+ # one scope's memories into another's recall. scope=None (default) sees everything (legacy behavior).
539
+ if scope is not None:
540
+ pool = [r for r in pool if (r.get("meta") or {}).get("scope") == scope]
541
+ # Metadata pre-filter (the 'filter before you rank' lever): keep only records matching ALL `where`
542
+ # conditions, matched against top-level fields then meta. Deterministic, no embedder, O(pool).
543
+ if where:
544
+ def _match(r: dict) -> bool:
545
+ meta = r.get("meta") or {}
546
+ for field, cond in where.items():
547
+ val = r[field] if field in r else meta.get(field)
548
+ if isinstance(cond, dict):
549
+ for op, cv in cond.items():
550
+ if op in ("$eq", "eq"):
551
+ if val != cv: return False
552
+ elif op in ("$ne", "ne"):
553
+ if val == cv: return False
554
+ elif op in ("$in", "in"):
555
+ if val not in cv: return False
556
+ elif op in ("$nin", "nin"):
557
+ if val in cv: return False
558
+ elif op in ("$gte", "gte"):
559
+ if val is None or val < cv: return False
560
+ elif op in ("$lte", "lte"):
561
+ if val is None or val > cv: return False
562
+ elif op in ("$gt", "gt"):
563
+ if val is None or val <= cv: return False
564
+ elif op in ("$lt", "lt"):
565
+ if val is None or val >= cv: return False
566
+ elif op in ("$contains", "contains"):
567
+ if val is None or cv not in val: return False
568
+ else:
569
+ raise ValueError(f"recall(where=): unknown operator {op!r}")
570
+ elif isinstance(cond, (list, tuple, set)):
571
+ if val not in cond: return False
572
+ else:
573
+ if val != cond: return False
574
+ return True
575
+ pool = [r for r in pool if _match(r)]
576
+ # Influence gate (retrieve-then-influence split): keep only CORROBORATED memories in the set that is
577
+ # allowed to drive an action. Same bar as episodic->semantic graduation; embedder-independent, so it
578
+ # generalizes across retrievers where geometry-based poison defenses do not (see the docstring).
579
+ if influence_only:
580
+ _byid = {x["id"]: x for x in self.items}
581
+ pool = [r for r in pool if self._is_corroborated(r, _byid)]
582
+ # Mode selection. 'hybrid' = lexical (token overlap) + semantic (embedding) fused with Reciprocal
583
+ # Rank Fusion. We MEASURED hybrid robustly beating EITHER channel alone for agent memory on LoCoMo
584
+ # (recall@20 0.61 hybrid vs 0.55 lexical vs 0.53 semantic; +0.057 over the best single channel,
585
+ # 9/10 conversations, conversation-level bootstrap CI excludes 0). So 'auto' now fuses (was: switch
586
+ # lexical->semantic at the threshold). Receipt: mnemo/probes/locomo_retrieval_map.py. RRF needs no
587
+ # tuning and no extra dependency. Force a single channel with mode='lexical'/'semantic'.
588
+ has_embed = self.embed is not None
589
+ if mode == "lexical" or not has_embed:
590
+ sel = "lexical"
591
+ elif mode in ("semantic", "hybrid"):
592
+ sel = mode
593
+ else: # 'auto': fuse once the store is worth it
594
+ sel = "hybrid" if len(pool) >= self.semantic_threshold else "lexical"
595
+ qvec = self._qvec(query) if sel in ("semantic", "hybrid") else None
596
+ if qvec is None and sel != "lexical":
597
+ sel = "lexical" # embedder absent or failed -> graceful fallback
598
+ self._last_mode = sel
599
+ qtok = _tokens(query) # tokenize the query once (lexical + fallback)
600
+ # Vectorized semantic fast-path: one matmul gives the cosine to every vec-bearing memory.
601
+ sims_vec = None
602
+ if qvec is not None and _np is not None:
603
+ M = self._vec_matrix()
604
+ if M is not None:
605
+ qv = _np.asarray(qvec, dtype=_np.float32)
606
+ if self.center_embeddings and self._vec_mean is not None:
607
+ qv = qv - self._vec_mean # center the query the SAME way as the matrix
608
+ sims_vec = M @ (qv / (float(_np.linalg.norm(qv)) or 1.0))
609
+ _now = time.time() # for per-type decay of the ranking value
610
+ _by_id = {x["id"]: x for x in self.items} # for provenance lookups (source-episode status)
611
+ def _semsim(r) -> float:
612
+ if sims_vec is not None and r.get("vec") and r["id"] in self._vec_rowof:
613
+ return max(0.0, float(sims_vec[self._vec_rowof[r["id"]]]))
614
+ return max(0.0, _cosine(qvec, r["vec"])) if (qvec is not None and r.get("vec")) else 0.0
615
+ def _lexsim(r) -> float:
616
+ t = self._rec_tokens(r)
617
+ return (len(qtok & t) / min(len(qtok), len(t))) if (qtok and t) else 0.0
618
+ def _candrec(r, sim): # provenance gate + value, shared by all modes
619
+ # Provenance gate: a memory that absorbed near-duplicates (links) is STALE-DERIVED if any of
620
+ # those sources was later CONTRADICTED (state-toggle supersession) β€” the merged summary
621
+ # outlived a fact it summarized. Demote it (don't drop β€” flag for re-consolidation), so a
622
+ # consolidated claim can't quietly outrank the fresh memory that overturned its source.
623
+ stale = bool(r.get("links")) and any(
624
+ (_by_id.get(lid, {}).get("meta") or {}).get("superseded_by_toggle") for lid in r["links"])
625
+ r["_stale_derived"] = stale # surfaced in the returned record
626
+ return (sim, 0.5 if stale else 1.0, self._effective_value(r, _now), r)
627
+ cands = [] # (sim, prov, eff_value, r), sim in [0,1]
628
+ # Relevance-floor ABSTENTION: drop candidates below an absolute similarity floor; if the WHOLE top-k
629
+ # falls below it, recall() returns [] ("not in memory") instead of padding context with a weak false
630
+ # match. min_relevance=0.0 (default) keeps legacy behavior (only sim<=0 dropped). In hybrid the floor
631
+ # is applied to the stronger of the two raw channels, then the FUSED rank score becomes the relevance.
632
+ if sel == "hybrid":
633
+ bm = self._bm25_scores(qtok, pool) # strong BM25 lexical channel over the live corpus
634
+ scn = [] # (r, sem, bm25) candidates above the floor
635
+ for r, bx in zip(pool, bm):
636
+ sem = _semsim(r)
637
+ if (sem <= 0 or sem < min_relevance) and bx <= 0:
638
+ continue # abstain only when BOTH channels are empty/below floor
639
+ scn.append((r, sem, bx))
640
+ if scn:
641
+ order_sem = sorted(range(len(scn)), key=lambda i: -scn[i][1])
642
+ order_bm = sorted(range(len(scn)), key=lambda i: -scn[i][2])
643
+ rrf = [0.0] * len(scn)
644
+ for rank, i in enumerate(order_sem): rrf[i] += 1.0 / (60 + rank)
645
+ for rank, i in enumerate(order_bm): rrf[i] += 1.0 / (60 + rank)
646
+ mx = max(rrf) or 1.0
647
+ for i, (r, sem, bx) in enumerate(scn):
648
+ cands.append(_candrec(r, rrf[i] / mx)) # normalize the fused rank score to a [0,1] relevance
649
+ else:
650
+ for r in pool:
651
+ sim = _semsim(r) if sel == "semantic" else _lexsim(r)
652
+ if sim <= 0 or sim < min_relevance:
653
+ continue
654
+ cands.append(_candrec(r, sim))
655
+ # Calibration WAS-IT-RIGHT: a per-memory Beta(good,bad) posterior nudges the score by track record.
656
+ # cal_mode controls how the outcome-credit channel is allowed to act (our measured signal-reliability
657
+ # law: a selection signal only beats relevance once reliability p > the no-signal floor 1/(1+D)):
658
+ # 'full' (default) β€” cal in [0.5, 1.5] (legacy: can promote AND demote).
659
+ # 'boost' β€” cal in [1.0, 1.5]: outcome-credit can PROMOTE a proven memory but never DEMOTE one below
660
+ # its relevance, so a wrong/random credit cannot suppress a correct memory (kills backfire).
661
+ # 'gated' β€” disable cal (->1.0) for this recall when the pooled signal looks weaker than 1/(1+D).
662
+ mode = getattr(self, "cal_mode", "full")
663
+ gate_off = False
664
+ if mode == "gated" and cands:
665
+ top = max(c[0] for c in cands)
666
+ near = [c for c in cands if c[0] >= top * 0.95] # candidates relevance can't separate
667
+ D = len(near)
668
+ if D >= 2:
669
+ g = sum(float(c[3].get("good", 0) or 0) for c in near)
670
+ b = sum(float(c[3].get("bad", 0) or 0) for c in near)
671
+ if (g + 1.0) / (g + b + 2.0) <= 1.0 / (1.0 + D):
672
+ gate_off = True
673
+ scored = []
674
+ for sim, prov, evalue, r in cands:
675
+ if gate_off:
676
+ cal = 1.0
677
+ else:
678
+ cal = 0.5 + self._reliability(r)
679
+ if mode == "boost" and cal < 1.0:
680
+ cal = 1.0
681
+ score = sim * (1.0 + math.log1p(max(0.0, evalue))) * prov * cal
682
+ scored.append((score, sim, r))
683
+ scored.sort(key=lambda x: -x[0])
684
+ out = []
685
+ _top_sim = scored[0][1] if scored else 1.0 # normalize reinforcement by this query's best match
686
+ for score, sim, r in scored[:k]:
687
+ # Relevance-weighted reinforcement: a strong, on-target hit reinforces value MORE than a
688
+ # marginal one that merely squeaked into the top-k. A flat +bump lets a memory that is a
689
+ # weak false-positive for many queries become 'immortal' β€” the popular-but-irrelevant
690
+ # failure mode. Weighting by this recall's relevance (normalized to the query's best hit)
691
+ # ties reinforcement to how well the memory actually answered. (Independently converged on
692
+ # in production by the Dakera and mem0 teams: weight access events by recall score, not raw
693
+ # count.)
694
+ rel = (sim / _top_sim) if _top_sim > 0 else 1.0
695
+ r["value"] += 0.25 * rel
696
+ r["last_access"] = _now # ...and resets the per-type decay clock
697
+ # Type GRADUATION: an episodic memory recalled into high accrued value has proven durable,
698
+ # so promote it to semantic β€” it stops fading on the fast 7-day episodic clock and decays
699
+ # on the slow semantic one instead. (Dakera's access-driven episodic->semantic promotion,
700
+ # gated on accrued VALUE rather than raw access count, so a popular-but-trivial memory
701
+ # doesn't graduate.)
702
+ # POISON guard (HARDENED 2026-06-25): durability must be EARNED by INDEPENDENT corroboration, not
703
+ # mere recall-frequency. The value bump above is correctness-blind, so a confabulation recalled
704
+ # enough would otherwise graduate to the durable (slow-decay) tier and entrench itself. A
705
+ # self-assertable `source` string or a SINGLE `links` edge is attacker-settable (AgentPoison /
706
+ # MINJA / OWASP-ASI06), so neither alone may confer durability. Require either an EARNED net-positive
707
+ # outcome (good>0 and good>=bad β€” set only by credit() resolving real work, not self-assertable), OR
708
+ # >=2 DISTINCT corroborating links (no single self-created edge suffices). An uncorroborated popular
709
+ # memory stays episodic and fades on the fast clock unless earned.
710
+ # SYBIL HARDENING (entity resolution): count DISTINCT CANONICAL sources among the corroborating
711
+ # links, not the raw link count. A naive "β‰₯2 links" lets an attacker mint independence by naming
712
+ # one origin many ways ("Wikipedia" / "wikipedia.org" / a full URL β†’ 3 links, 1 real source).
713
+ # Canonicalizing source identifiers before counting collapses those to one; a link whose record
714
+ # has no source counts as its own id, so genuinely source-less corroboration is unchanged.
715
+ _good = float(r.get("good", 0) or 0); _bad = float(r.get("bad", 0) or 0)
716
+ corroborated = (_good > 0 and _good >= _bad) or self._distinct_sources(r.get("links"), _by_id) >= 2
717
+ if r.get("mtype") == "episodic" and r["value"] >= _GRADUATE_VALUE and corroborated:
718
+ r["mtype"] = "semantic"
719
+ r.setdefault("meta", {})["graduated_from_episodic"] = True
720
+ out.append({"id": r["id"], "text": r["text"], "tags": r["tags"], "iso": r["iso"],
721
+ "value": round(r["value"], 2), "relevance": round(sim, 3),
722
+ "score": round(score, 3), "links": r["links"],
723
+ "reliability": round(self._reliability(r), 3),
724
+ "source": r.get("source"), # re-checkable origin (provenance), surfaced so a recalled fact can be traced back
725
+ "stale_derived": bool(r.get("_stale_derived"))})
726
+ # NOTE: recall is a READ. It nudges in-memory access value / graduation, but must NOT persist the
727
+ # whole store here β€” serializing (json.dumps) on every recall, across many agents' stores,
728
+ # saturated the thread pool and FROZE the world. The in-memory nudges are persisted on the next
729
+ # remember()/consolidate()/flush(); losing recent access metadata on a hard crash is harmless.
730
+ if out:
731
+ self._dirty = True # mark for the next throttled/forced save; do NOT serialize on the read path
732
+ return out
733
+
734
+ @staticmethod
735
+ def _canon_source(doc) -> str:
736
+ """Entity-resolution canonicalization of a source identifier, so sybil variants of one origin
737
+ ('Wikipedia', 'wikipedia.org', 'https://www.wikipedia.org/wiki/X') collapse to a single key."""
738
+ s = str(doc or "").strip().lower()
739
+ s = re.sub(r"^[a-z]+://", "", s) # strip scheme
740
+ s = re.sub(r"^www\.", "", s) # strip www.
741
+ s = s.split("/")[0].split("?")[0] # host / first path segment only
742
+ s = re.sub(r"\.(org|com|net|io|gov|edu|co|ai|dev|info|news)$", "", s) # strip a common TLD
743
+ s = re.sub(r"[^a-z0-9]+", "", s) # collapse remaining punctuation
744
+ return s
745
+
746
+ @staticmethod
747
+ def _distinct_sources(links, by_id) -> int:
748
+ """Count DISTINCT canonical sources among corroborating links β€” entity resolution BEFORE counting,
749
+ so 'three names for one source' sybil variants count as one. A link whose record carries no source
750
+ counts as its own id, so genuinely source-less corroboration is not penalised (no regression)."""
751
+ keys = set()
752
+ for lid in (links or []):
753
+ lr = by_id.get(lid)
754
+ if lr is None:
755
+ continue
756
+ src = lr.get("source")
757
+ doc = src.get("doc") if isinstance(src, dict) else (src if isinstance(src, str) else None)
758
+ keys.add(Mnemo._canon_source(doc) if doc else "id:" + lid)
759
+ return len(keys)
760
+
761
+ @staticmethod
762
+ def _is_corroborated(rec: dict, by_id: dict) -> bool:
763
+ """The corroboration bar shared by episodic->semantic graduation and the recall influence gate:
764
+ an EARNED net-positive outcome (good>0 and good>=bad β€” set by credit() on real work, not
765
+ self-assertable), OR an already-graduated 'semantic' memory, OR >=2 DISTINCT-canonical-source
766
+ corroborating links (sybil variants of one source collapse to one). A single fresh self-asserted
767
+ memory (the AgentPoison single-instance poison) meets none of these."""
768
+ good = float(rec.get("good", 0) or 0)
769
+ bad = float(rec.get("bad", 0) or 0)
770
+ if good > 0 and good >= bad:
771
+ return True
772
+ if rec.get("mtype") == "semantic":
773
+ return True
774
+ return Mnemo._distinct_sources(rec.get("links"), by_id) >= 2
775
+
776
+ @staticmethod
777
+ def _reliability(r: dict) -> float:
778
+ """Per-memory track record as a Beta(1+good, 1+bad) posterior MEAN: 0.5 with no outcomes yet,
779
+ ->1 if recalls into it kept resolving WELL, ->0 if they kept resolving badly. Counts only grow."""
780
+ g = float(r.get("good", 0) or 0)
781
+ b = float(r.get("bad", 0) or 0)
782
+ return (g + 1.0) / (g + b + 2.0)
783
+
784
+ def credit(self, ids, outcome, weight: float = 1.0) -> dict:
785
+ """Close the accuracy loop onto the substrate. When the work a set of memories was recalled into
786
+ gets a real verdict (a forecast resolves, a replication is ruled REPRODUCED/FAILED, a hypothesis is
787
+ severe-tested), call credit(recalled_ids, outcome): each memory's Beta(good,bad) track record is
788
+ nudged so future recall ranks by WAS-IT-RIGHT, not merely was-it-recalled. Append-only to the
789
+ counts; never edits raw text. `outcome` may be a bool, a sign (>0 good), or a verdict string
790
+ (good/right/correct/reproduced/hit vs bad/wrong/failed/miss)."""
791
+ if isinstance(outcome, bool):
792
+ good = outcome
793
+ elif isinstance(outcome, (int, float)):
794
+ good = outcome > 0
795
+ else:
796
+ s = str(outcome).strip().lower()
797
+ good = s in ("good", "right", "correct", "reproduced", "hit", "true", "win", "+")
798
+ by_id = {x["id"]: x for x in self.items}
799
+ key, updated = ("good" if good else "bad"), []
800
+ for i in (ids or []):
801
+ rec = by_id.get(i)
802
+ if rec is None:
803
+ continue
804
+ rec[key] = float(rec.get(key, 0) or 0) + float(weight)
805
+ updated.append(i)
806
+ if updated:
807
+ self._save()
808
+ return {"updated": updated, "outcome": key, "weight": weight}
809
+
810
+ def _effective_value(self, r: dict, now: float) -> float:
811
+ """Recall weight = stored value decayed by time since last access, at the memory's TYPE
812
+ half-life (episodic fades fast, semantic slow, procedural barely). Access resets the clock,
813
+ so memories that keep being useful stay alive while stored-but-never-recalled ones fade.
814
+ Reversible: raw value/text are untouched; only the effective ranking weight decays."""
815
+ hl = _HALFLIFE_S.get(r.get("mtype", "episodic"), _HALFLIFE_S["episodic"])
816
+ age = max(0.0, now - r.get("last_access", r.get("ts", now)))
817
+ return r["value"] * (0.5 ** (age / hl))
818
+
819
+ # ── consolidation (the "dream" pass) ──────────────────────────────────────
820
+ def _common_vocab(self, active: list[dict], min_df_frac: float = 0.002):
821
+ """Token sets per memory + the corpus's COMMON vocabulary (tokens shared by enough
822
+ memories to be real content, not one-off noise). Cheap, O(total tokens)."""
823
+ from collections import Counter
824
+ df: Counter = Counter()
825
+ toks = []
826
+ for r in active:
827
+ tk = _tokens(r["text"]); toks.append(tk); df.update(tk)
828
+ min_df = max(3, int(min_df_frac * len(active)))
829
+ common = {w for w, c in df.items() if c >= min_df}
830
+ return toks, common
831
+
832
+ def recall_iterative(self, query: str, ask_followup, k: int = 6, rounds: int = 1,
833
+ **recall_kw) -> list[dict]:
834
+ """Multi-hop recall. One-shot top-k misses evidence reachable only via a BRIDGE entity (a fact whose
835
+ detail lives in a memory NOT similar to the query). This does: retrieve -> let a capable model read the
836
+ results and name what's missing, emitting follow-up queries -> retrieve again -> merge (dedup by id).
837
+ `ask_followup(query, current_results) -> list[str]` is caller-supplied, so mnemo stays model-agnostic
838
+ (inject any model/LLM). MEASURED ~3.3x multi-hop full-evidence recall vs one-shot top-k on LoCoMo
839
+ (0.057 -> 0.186, n=70 across 3 conversations) β€” the one mechanism that moved the multi-hop bottleneck
840
+ where static retrieval tricks (dense-neighbor, lexical bridges) did not. More expensive (a model call
841
+ in the loop), so it's an explicit mode, not the default."""
842
+ seen: dict = {}
843
+ for r in self.recall(query, k=k, **recall_kw):
844
+ seen[r["id"]] = r
845
+ for _ in range(max(0, int(rounds))):
846
+ try:
847
+ followups = ask_followup(query, list(seen.values())) or []
848
+ except Exception:
849
+ followups = []
850
+ for fq in followups:
851
+ if not isinstance(fq, str) or not fq.strip():
852
+ continue
853
+ for r in self.recall(fq, k=k, **recall_kw):
854
+ seen.setdefault(r["id"], r)
855
+ return list(seen.values())
856
+
857
+ def consolidate(self, keep: int | None = None, dup_threshold: float = 0.82,
858
+ hub_coverage: float = 0.12, link_duplicates: bool = True) -> dict:
859
+ """The dream pass. ADDS a derived layer (status + links); never edits raw text. Three steps:
860
+
861
+ 1. HUB PASS β€” flag indiscriminate "universal-matcher" memories. Under lexical recall the
862
+ similarity is the overlap coefficient |q∩t|/min(|q|,|t|), so a memory whose token set
863
+ covers a large fraction of the corpus's common vocabulary scores ~1.0 against ALMOST ANY
864
+ query and drowns the specific memory the user actually wanted (measured on a 6k-note
865
+ vault: such hubs sat in the top-10 for ~47% of queries). We mark them `status:'hub'`
866
+ (reversible; recall skips them unless include_hubs) β€” measured to lift recall@5 ~+22%.
867
+ 2. near-duplicate LINKING (dedup without delete) β€” EXCEPT a polarity clash, which is a
868
+ STATE TOGGLE (preference flip): supersede the OLDER, since a contradiction is not a dup.
869
+ 3. keep-budget: mark the lowest-value surplus `superseded`.
870
+
871
+ hub_coverage: a memory covering β‰₯ this fraction of the common vocabulary is a hub (0 disables).
872
+ link_duplicates: the dup pass is O(nΒ²); pass False to skip it on large stores."""
873
+ active = [r for r in self.items if r["status"] == "active"]
874
+ hubs = 0
875
+ if hub_coverage and len(active) >= 50:
876
+ toks, common = self._common_vocab(active)
877
+ nv = len(common) or 1
878
+ for r, tk in zip(active, toks):
879
+ shared = len(tk & common)
880
+ cov = shared / nv
881
+ # A genuine 'universal matcher' overlaps MANY of the corpus's common words. Requiring an absolute
882
+ # floor (>= 3 shared common words) on top of the coverage fraction prevents the low-diversity /
883
+ # templated-store failure: when the common vocabulary is tiny (e.g. a handful of repeated attribute
884
+ # words), a legitimate memory trivially covers >= hub_coverage of it with just ONE common word, which
885
+ # would wrongly flag every memory a hub and SILENTLY EMPTY recall. (Measured: 3-5 shared attrs -> 100%
886
+ # hub-flagged, 0% recall, before this floor.)
887
+ if shared >= 3 and cov >= hub_coverage:
888
+ r["status"] = "hub"
889
+ r.setdefault("meta", {})["hub"] = True
890
+ r["meta"]["hub_coverage"] = round(cov, 3)
891
+ r["superseded_ts"] = time.time()
892
+ hubs += 1
893
+ active = [r for r in active if r["status"] == "active"]
894
+ active.sort(key=lambda r: -r["value"])
895
+ linked = toggled = 0
896
+ if link_duplicates:
897
+ # Pairwise near-duplicate pass. A high-similarity pair is normally LINKED (dedup without
898
+ # delete) β€” UNLESS it's a polarity clash (one negates the other), which is a STATE TOGGLE
899
+ # (a preference flip / contradiction), not a duplicate. Then we supersede the OLDER memory
900
+ # so recall returns the NEW state, instead of letting high vector similarity silently
901
+ # merge a contradiction into one blob. (state-toggle guard.)
902
+ for i, a in enumerate(active):
903
+ if a["status"] != "active": # superseded by an earlier toggle this pass
904
+ continue
905
+ avec = self._qvec(a["text"]) # embed each anchor once, not once per partner
906
+ for b in active[i + 1:]:
907
+ if b["status"] != "active" or b["id"] in a["links"]:
908
+ continue
909
+ if self._similarity(a["text"], b, avec) >= dup_threshold:
910
+ if _negation_clash(a["text"], b["text"]) or _value_clash(a["text"], b["text"]):
911
+ # Resolve by VALIDITY time (valid_from = when the fact is TRUE), not ingest order
912
+ # (ts = when it was stored). A fact learned LATE about an EARLIER state (e.g. a
913
+ # back-filled record) must NOT overwrite the genuinely-current one just because it
914
+ # arrived later. valid_from defaults to ts, so ingest-ordered streams are unchanged;
915
+ # only out-of-order arrivals (the bi-temporal case) flip vs the old ts rule.
916
+ _vf = lambda r: r.get("valid_from", r["ts"])
917
+ older, newer = (a, b) if _vf(a) <= _vf(b) else (b, a)
918
+ # Fast-novelty guard (opt-in): supersede only on a CORROBORATED contradiction
919
+ # (earned credit, or >=2 links β€” same bar as graduation). An uncorroborated
920
+ # single contradiction is recorded as a link but does NOT override a standing
921
+ # fact (resists single-shot poison flips). Default OFF -> legacy fast behavior.
922
+ if self.supersede_requires_corroboration:
923
+ _ng = float(newer.get("good", 0) or 0); _nb = float(newer.get("bad", 0) or 0)
924
+ if not ((_ng > 0 and _ng >= _nb) or len(newer.get("links") or []) >= 2):
925
+ a["links"].append(b["id"]); linked += 1
926
+ continue
927
+ # Persistence (CUSUM) guard: supersede only once the NEW state is asserted by
928
+ # >= supersede_persistence independent records (the change has persisted). Count
929
+ # active records that (i) match newer's value/polarity and (ii) contradict older β€”
930
+ # an isolated poison flip stays below the threshold and is merely linked.
931
+ if self.supersede_persistence > 1:
932
+ nvec = self._qvec(newer["text"])
933
+ support = sum(
934
+ 1 for r in active if r["status"] == "active"
935
+ and self._similarity(newer["text"], r, nvec) >= dup_threshold
936
+ and not _value_clash(newer["text"], r["text"])
937
+ and not _negation_clash(newer["text"], r["text"])
938
+ and (_value_clash(older["text"], r["text"]) or _negation_clash(older["text"], r["text"])))
939
+ if support < self.supersede_persistence:
940
+ a["links"].append(b["id"]); linked += 1
941
+ continue
942
+ older["status"] = "superseded"
943
+ older["superseded_ts"] = time.time()
944
+ older["invalidated_at"] = _vf(newer) # bi-temporal: when this record stopped being current
945
+ older.setdefault("meta", {})["superseded_by_toggle"] = newer["id"]
946
+ # Accuracy loop, live consumer: being OVERTURNED by a later contradiction is
947
+ # a was-wrong signal β€” debit the superseded claim, credit the one that
948
+ # corrected the record. So the consolidation pass continuously feeds each
949
+ # memory's reliability from real outcomes, not just external scoring.
950
+ older["bad"] = float(older.get("bad", 0) or 0) + 1.0
951
+ newer["good"] = float(newer.get("good", 0) or 0) + 1.0
952
+ toggled += 1
953
+ if older is a:
954
+ break # this anchor is gone; advance to the next
955
+ else:
956
+ a["links"].append(b["id"]); linked += 1
957
+ staled = 0
958
+ if keep is not None and len(active) > keep:
959
+ # active is sorted by -raw value (above). Legacy = keep the top-`keep` by raw value. Two-tier =
960
+ # protect the top kprot by raw value (recency-immune), then fill the remaining budget from the
961
+ # REST by EFFECTIVE (decay-weighted) value, so a stale high-raw-value memory can't crowd out a
962
+ # freshly-useful one. (kprot=0 for tiny budgets -> pure recency-aware fill.)
963
+ if self.two_tier_keep:
964
+ now = time.time()
965
+ kprot = int(self.protect_frac * keep)
966
+ protected, rest = active[:kprot], active[kprot:]
967
+ rest_keep = set(id(r) for r in
968
+ sorted(rest, key=lambda r: -self._effective_value(r, now))[:keep - kprot])
969
+ drop = [r for r in rest if id(r) not in rest_keep]
970
+ else:
971
+ drop = active[keep:]
972
+ for r in drop:
973
+ r["status"] = "superseded"; r["superseded_ts"] = time.time(); staled += 1
974
+ self._save()
975
+ return {"active": len([r for r in self.items if r["status"] == "active"]),
976
+ "hubs_flagged": hubs, "linked_pairs": linked, "toggled": toggled,
977
+ "staled": staled, "kept": keep, "total": len(self.items)}
978
+
979
+ # ── cluster-triggered consolidation ───────────────────────────────────────
980
+ def _cluster_active(self, sim_threshold: float = 0.5) -> list[list[dict]]:
981
+ """Cheap greedy single-pass clustering of ACTIVE memories by similarity (O(nΒ·#clusters)).
982
+ Highest-value member is the cluster representative; each memory joins the most-similar
983
+ cluster above the threshold, else starts its own. Lexical or semantic per the store's mode."""
984
+ active = sorted([r for r in self.items if r["status"] == "active"], key=lambda r: -r["value"])
985
+ cents: list[dict] = []
986
+ for r in active:
987
+ rvec = self._qvec(r["text"])
988
+ best = None
989
+ for c in cents:
990
+ s = self._similarity(c["rec"]["text"], r, c["vec"])
991
+ if s >= sim_threshold and (best is None or s > best[1]):
992
+ best = (c, s)
993
+ if best:
994
+ best[0]["members"].append(r)
995
+ else:
996
+ cents.append({"rec": r, "vec": rvec, "members": [r]})
997
+ return [c["members"] for c in cents]
998
+
999
+ def consolidate_clusters(self, threshold: int = 15, cluster_sim: float = 0.5,
1000
+ dup_threshold: float = 0.82, keep_per_cluster: int | None = None) -> dict:
1001
+ """Cluster-TRIGGERED consolidation: consolidate a semantic cluster only once it has grown past
1002
+ `threshold` members β€” not a global nightly blanket. Avoids (1) prematurely consolidating sparse
1003
+ topics, where the raw episodes are still the best representation, and (2) unbounded growth in
1004
+ dense ones. Cheap to call often (no-op until a cluster is ripe). Runs dedup + the state-toggle
1005
+ guard (+ optional keep-budget) WITHIN each ripe cluster only."""
1006
+ clusters = self._cluster_active(cluster_sim)
1007
+ fired = linked = toggled = staled = 0
1008
+ for members in clusters:
1009
+ if len(members) < threshold:
1010
+ continue # sparse β€” leave the raw episodes alone
1011
+ fired += 1
1012
+ members.sort(key=lambda r: -r["value"])
1013
+ for i, a in enumerate(members):
1014
+ if a["status"] != "active":
1015
+ continue
1016
+ avec = self._qvec(a["text"])
1017
+ for b in members[i + 1:]:
1018
+ if b["status"] != "active" or b["id"] in a["links"]:
1019
+ continue
1020
+ if self._similarity(a["text"], b, avec) >= dup_threshold:
1021
+ if _negation_clash(a["text"], b["text"]) or _value_clash(a["text"], b["text"]):
1022
+ older, newer = (a, b) if a["ts"] <= b["ts"] else (b, a)
1023
+ older["status"] = "superseded"; older["superseded_ts"] = time.time()
1024
+ older.setdefault("meta", {})["superseded_by_toggle"] = newer["id"]
1025
+ toggled += 1
1026
+ if older is a:
1027
+ break
1028
+ else:
1029
+ a["links"].append(b["id"]); linked += 1
1030
+ if keep_per_cluster is not None:
1031
+ act = sorted([r for r in members if r["status"] == "active"], key=lambda r: -r["value"])
1032
+ for r in act[keep_per_cluster:]:
1033
+ r["status"] = "superseded"; r["superseded_ts"] = time.time(); staled += 1
1034
+ self._save()
1035
+ return {"clusters_total": len(clusters), "clusters_fired": fired, "threshold": threshold,
1036
+ "linked_pairs": linked, "toggled": toggled, "staled": staled}
1037
+
1038
+ # ── contradiction surfacing (flag, never auto-delete) ─────────────────────
1039
+ def contradictions(self, sim_threshold: float = 0.5, incompatible=None) -> list[dict]:
1040
+ """Flag mutually-incompatible memories among RELATED ones (similarity-gated) for human review.
1041
+ `incompatible(a_text, b_text)->bool` defaults to a negation/polarity heuristic."""
1042
+ inc = incompatible or _negation_clash
1043
+ active = [r for r in self.items if r["status"] == "active"]
1044
+ flags = []
1045
+ for i, a in enumerate(active):
1046
+ avec = self._qvec(a["text"]) # embed each anchor once, not once per partner
1047
+ for b in active[i + 1:]:
1048
+ if self._similarity(a["text"], b, avec) >= sim_threshold and inc(a["text"], b["text"]):
1049
+ flags.append({"a": a["id"], "b": b["id"],
1050
+ "a_text": a["text"][:120], "b_text": b["text"][:120]})
1051
+ return flags
1052
+
1053
+ # ── value, reported at the COHORT level ───────────────────────────────────
1054
+ def value_by_cohort(self) -> dict:
1055
+ """Per-TAG value rollup. Deliberately not per-memory: at n-of-1, per-item value is noise;
1056
+ the cohort (tag / time-block) is where the signal is real."""
1057
+ out: dict[str, dict] = {}
1058
+ for r in self.items:
1059
+ if r["status"] != "active":
1060
+ continue
1061
+ for tag in (r["tags"] or ["(untagged)"]):
1062
+ c = out.setdefault(tag, {"count": 0, "value": 0.0})
1063
+ c["count"] += 1; c["value"] += r["value"]
1064
+ return {k: {"count": v["count"], "value": round(v["value"], 2),
1065
+ "avg": round(v["value"] / v["count"], 2)} for k, v in out.items()}
1066
+
1067
+ def _save(self, force: bool = False):
1068
+ if not self.path:
1069
+ return
1070
+ # Throttle: coalesce frequent writes (e.g. one per recall) so a large store isn't re-serialized
1071
+ # on the hot path. force=True (or flush()) bypasses it for shutdown/critical persistence.
1072
+ now = time.time()
1073
+ if not force and (now - self._last_save) < self._save_min_s:
1074
+ self._dirty = True
1075
+ return
1076
+ try:
1077
+ # Persist text/metadata only; the `vec` embedding arrays are a re-derivable in-memory CACHE
1078
+ # and are STRIPPED here. json.dumps of N x 768-dim float vectors is huge, slow, and holds the
1079
+ # GIL for many seconds - which froze the whole event loop even from a worker thread (the
1080
+ # frozen-world bug, 2026-06-20). Vectors stay in self.items (RAM) so recall is unaffected this
1081
+ # session; on reload they are re-embedded lazily. Keeps the store file small + the save fast.
1082
+ slim = [{k: v for k, v in r.items() if k != "vec"} for r in self.items]
1083
+ # Atomic write: a partial/interleaved write can't corrupt the store (crash- and
1084
+ # concurrent-writer-safe β€” last writer wins, never a torn JSON file).
1085
+ data = json.dumps(slim, ensure_ascii=False, indent=1)
1086
+ tmp = self.path.with_name(self.path.name + ".tmp")
1087
+ tmp.write_text(data, encoding="utf-8")
1088
+ os.replace(tmp, self.path)
1089
+ self._last_save = now
1090
+ self._dirty = False
1091
+ except Exception:
1092
+ pass
1093
+
1094
+ def flush(self):
1095
+ """Force-persist any pending throttled changes (call on clean shutdown)."""
1096
+ if self._dirty:
1097
+ self._save(force=True)
1098
+
1099
+
1100
+ # ── per-type decay priors (the half-life a memory's ranking value decays at, by kind) ──────────
1101
+ # episodic = events (fade fast); semantic = durable facts (fade slow); procedural = rules/prefs
1102
+ # (barely fade). Access resets the decay clock (see Mnemo._effective_value). Tunable.
1103
+ _HALFLIFE_S = {"episodic": 7 * 86400, "semantic": 180 * 86400, "procedural": 3650 * 86400}
1104
+ # accrued value at which a repeatedly-recalled EPISODIC memory graduates to semantic (β‰ˆ16 strong
1105
+ # recalls from the 1.0 floor); proven-durable, so it should decay on the slow clock, not the fast one.
1106
+ _GRADUATE_VALUE = 5.0
1107
+ _PROCEDURAL_RE = re.compile(r"\b(always|never|prefers?|rule|workflow|convention|policy|habit|"
1108
+ r"setting|must|should|avoid|don't|do not)\b", re.I)
1109
+ _SEMANTIC_RE = re.compile(r"\b(means|defined|definition|theorem|law of|equals|consists? of|"
1110
+ r"is a |is an |is the |refers to)\b", re.I)
1111
+
1112
+
1113
+ def _infer_type(text: str) -> str:
1114
+ """Conservative type inference: default EPISODIC (fast decay) and only promote on clear markers.
1115
+ Callers that know the kind should pass mtype explicitly."""
1116
+ t = text or ""
1117
+ if _PROCEDURAL_RE.search(t):
1118
+ return "procedural"
1119
+ if _SEMANTIC_RE.search(t):
1120
+ return "semantic"
1121
+ return "episodic"
1122
+
1123
+
1124
+ def _negation_clash(a: str, b: str) -> bool:
1125
+ """Cheap default: two highly-related statements where exactly one negates. Replace with an
1126
+ LLM judge for production β€” but gate it behind similarity first to keep it O(neighbourhood)."""
1127
+ neg = re.compile(r"\b(not|no|never|cannot|can't|doesn't|isn't|won't|fails?|false)\b", re.I)
1128
+ return bool(neg.search(a)) != bool(neg.search(b))
1129
+
1130
+
1131
+ _NUM = re.compile(r"-?\d+(?:\.\d+)?")
1132
+
1133
+
1134
+ def _value_clash(a: str, b: str) -> bool:
1135
+ """A VALUE UPDATE: two already-near-duplicate statements that are identical EXCEPT for a differing
1136
+ numeric value ('retry limit is 5' -> '... is 12'). This is a state toggle (the fact's value changed),
1137
+ NOT a duplicate β€” so the older should be superseded, not merged. Gated behind the caller's similarity
1138
+ check; the tight 'non-numeric remainder is identical' condition keeps genuinely-distinct facts safe."""
1139
+ # A value UPDATE keeps the same numbers in the same ORDER except ONE position whose value changed
1140
+ # ('timeout is 5' -> 'is 12'; '5 of 10' -> '7 of 10'). Compare numbers POSITIONALLY, not as sets: a
1141
+ # set view is ambiguous for ENUMERATED facts because an index can equal another row's value
1142
+ # ('step 1 takes 5 min' vs 'step 5 takes 13 min' share the literal 5), which set-math reads as a
1143
+ # single change and would silently supersede a coexisting record. (Measured: a 6-item enumerated store
1144
+ # lost 5/6 facts under the set rule; 0/6 under this positional rule.)
1145
+ na, nb = _NUM.findall(a), _NUM.findall(b) # ORDERED, not sets
1146
+ if not na or len(na) != len(nb):
1147
+ return False # no numbers, or different count -> not a single update
1148
+ if sum(1 for x, y in zip(na, nb) if x != y) != 1:
1149
+ return False # exactly one positional value changed
1150
+ # Compare the word-skeleton with ALL numbers stripped: _tokens keeps 3+ digit numbers as tokens
1151
+ # (_WORD requires length >= 3), so a multi-digit value ('...is 123') would otherwise spuriously make
1152
+ # the skeletons differ and miss the update. Strip numbers first, exactly as before this guard existed.
1153
+ return _tokens(_NUM.sub("", a)) == _tokens(_NUM.sub("", b)) # identical apart from the one value
1154
+
1155
+
1156
+ if __name__ == "__main__":
1157
+ m = Mnemo() # no path, no embedder β€” pure in-memory + lexical
1158
+ m.remember("SGD converges slowly due to gradient variance.", tags=["optimization"], value=3)
1159
+ m.remember("SGD does not converge slowly.", tags=["optimization"], value=1)
1160
+ m.remember("Pre-trend tests catch only 31% of fatal DiD bias.", tags=["causal"], value=2)
1161
+ print("recall 'SGD variance':", [r["text"][:46] for r in m.recall("SGD variance", k=3)])
1162
+ print("consolidate:", m.consolidate(keep=10))
1163
+ print("contradictions:", m.contradictions()) # flags the SGD pair (related + one negates)
1164
+ print("value_by_cohort:", m.value_by_cohort())
1165
+ print("(For semantic recall, pass embed=your_model to Mnemo(); lexical is the zero-dep fallback.)")