feat(mem0-agent): capture gate, eval harness, and mixed-window fix

- triggers.py: hard-drop rules now filter PER TURN, so a durable fact sitting
  between two progress lines survives instead of being dropped with them. The
  eval harness caught this: 3 mixed-window fixtures were being destroyed
  client-side before the extractor ever saw them.
- widened flag rules for 13 plainly durable windows the gate was missing
  (standing preferences, X-over-Y decisions, stated rules, diagnoses, verified
  procedures); verified procedures moved from aggressive to balanced
- eval/: 56 labeled fixtures from the audited v1 corpus plus offline and live
  runners with a --check regression gate
- un-ignore integrations/mem0-agent/eval (root .gitignore excluded it silently)

Offline scorecard: hard_drop_recall 1.000, extract_recall 1.000,
flag_precision 1.000, 0 misclassified.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Deshraj Yadav
2026-07-28 02:02:37 -07:00
parent 3ea00cf88c
commit 4d5f2e653c
8 changed files with 1735 additions and 11 deletions
+3
View File
@@ -186,6 +186,9 @@ notebooks/*.yaml
# local directories for testing
eval/
# ...but the plugin's evaluation harness is source, not local scratch
!integrations/mem0-agent/eval/
integrations/mem0-agent/eval/last_report.json
qdrant_storage/
.crossnote
testing.ipynb
+193
View File
@@ -0,0 +1,193 @@
# Write-gate evaluation
The v1 plugin had no way to measure extraction quality, so it degraded for three months
without anyone noticing. When the corpus was finally audited, 20.5% of it was
near-duplicate heartbeats — the largest single cluster was 119 near-identical
training-progress frames — and organic searches had fallen from 257/month to 75/month.
The memory got noisy, people stopped trusting it, and nothing in the system said so.
This directory is the instrument that would have said so. It scores the write gate
against a labeled fixture set and fails a build when the score regresses.
The gate has two halves, and each mode measures one of them:
| Half | Where it lives | Measured by |
| --- | --- | --- |
| Local trigger rules — mechanical noise never leaves the machine | `src/mem0_agent/triggers.py` | `--offline` |
| Custom instructions — the platform stores nothing from a narration/activity/repo-file window | `src/mem0_agent/config/project_config.py` (`INSTRUCTIONS`) | `--live` |
## Running it
### Offline (default, no network, no credentials)
```bash
cd integrations/mem0-agent
PYTHONPATH=src python3 eval/run.py --offline
PYTHONPATH=src python3 eval/run.py --offline --level aggressive # sweep a capture level
PYTHONPATH=src python3 eval/run.py --offline --check # exit non-zero on regression
```
Runs all fixtures through `mem0_agent.triggers.classify` and scores the decisions
against the labels. It costs nothing, so it belongs in CI on every commit that touches
`triggers.py`.
`triggers.py` is imported lazily. If it is missing or its `classify()` signature is
unreadable, the run reports the reason and scores nothing rather than crashing — but
`--check` then exits non-zero, because a gate that cannot measure must not pass.
### Live (writes to a scratch project)
```bash
export MEM0_API_KEY=...
PYTHONPATH=src python3 eval/run.py --live \
--project-id proj_SCRATCH --org-id org_YOURS --cleanup --check
```
Replays every `exclude` and `extract` fixture through `mem0_agent.api.Api` with
`infer=True`, waits for extraction, reads back, and scores what the platform actually
stored. This is the only way to test the custom instructions; they are a prompt, and
prompts are not unit-testable.
Safety and correctness properties, each of which exists because of something that went
wrong before:
- **`--project-id` and `--org-id` are required and have no defaults.** A live run writes
real memories. v1's benchmark data ended up in the production project because a
harness defaulted to whatever the API key resolved to. Point these at a scratch
project.
- **Every fixture gets its own user id**, `eval-<runid>-<fixture id>`, under app id
`eval-<runid>`. Results are attributable to a fixture and a run, and anything left
behind is trivially findable.
- **It polls; it does not sleep.** Extraction landed anywhere between 20s and 5min
during validation, so any fixed wait is either wrong or wasteful. `--timeout`
(default 360s) bounds the wait; `--min-settle` (default 60s) is the minimum time
before a zero read is allowed to count as suppression, since "nothing stored" and
"not stored yet" look identical.
- **It scores `metadata.type`, never `categories`.** Platform categorization lags ~3.9h
at the median and is 0% for memories under an hour old, so a category-based score
would read zero on a fresh run and tell you nothing.
- **`--cleanup` deletes everything the run wrote.** Without it the report prints the
`app_id` and user-id prefix needed to clean up by hand.
One honesty caveat on the live `type_match` metric: `metadata.type` is stamped by the
*client* at write time (from `triggers.classify` when available, otherwise from the
fixture's `expect_type`). So `type_match` measures that the type survives the round
trip, and — when the classifier is present — that the classifier chose the right one.
It is not the platform's independent opinion. The report records `stamped_by` per
fixture so a reader can tell which case they are looking at.
Both modes write `eval/last_report.json`: every score, every per-fixture row, the
thresholds in force, and the pass/fail verdict.
## The fixture set
`fixtures.py` holds the labeled windows. Each entry is exactly:
```python
{"id": str, "window": [{"role", "content"}, ...], "label": "drop"|"exclude"|"extract",
"expect_type": str | None, "note": str}
```
| Label | Contract | Enforced by |
| --- | --- | --- |
| `drop` | Must never be sent to the platform at all | client trigger rules |
| `exclude` | May be sent, but the platform must store nothing from it | custom instructions |
| `extract` | Must produce ≥1 memory, ideally of `expect_type` | both halves |
The `drop` and `exclude` windows use the wording of the audited v1 corpus: the
epoch/loss/ETA training heartbeat, the `N of M chunks processed (X%)` frame, the
markdown-ingest `pid`/`elapsed` line, monitoring status lines, tool-only turns,
assistant narration, file-modification lists, `CLAUDE.md`/README excerpts, one-off
directives (`you do it`), and session-only status (`at 12:45:24 the job was wrapping
up`). Several noise classes appear two or three times with only the numbers changed —
that is what a near-duplicate cluster actually looks like, and a gate that catches the
first frame but not the third has learned the numbers rather than the shape.
The `extract` block covers all five durable types plus four **mixed** windows, where a
durable fact sits between two progress lines. Mixed windows are the most informative
fixtures in the set: they fail in both directions. A gate too eager to drop heartbeats
destroys the fact along with them; a gate too eager to store keeps the heartbeats.
Print coverage at any time:
```bash
python3 eval/fixtures.py
```
## Thresholds
`--check` exits non-zero when a gated metric falls below its floor:
| Metric | Floor | Why |
| --- | --- | --- |
| `hard_drop_recall` | 0.95 | Fraction of `drop` fixtures the client never sends. This is the number that was silently 0 in v1. |
| `extract_recall` | 0.80 | Fraction of `extract` fixtures that survive. Offline: flagged for capture. Live: ≥1 memory stored. A gate that stores nothing scores perfectly on noise. |
Both floors must hold; they measure opposite failure directions and either one alone is
trivially gamed.
Other metrics are reported but not gated, because they diagnose rather than decide:
- `hard_drop_explicit` — how much of the noise containment comes from a real drop rule
rather than from no flag rule happening to match. `capture.py` only forwards windows
whose action is `flag`, so a `skip` does contain the noise — but only until someone
adds a flag rule that matches it. A gap between `hard_drop_recall` and
`hard_drop_explicit` is a list of heartbeats held back by luck.
- `hard_drop_precision` — of everything hard-dropped, how much was safe to drop. This is
where over-broad drop rules show up, and it is the metric the mixed fixtures exist to
move. A hard drop on an `extract` window is the one unrecoverable error: the memory is
gone and nothing logs a miss. Hard-dropping an `exclude` window is *not* counted
against this score — those are meant to be discarded, and discarding them locally is
simply cheaper than having the platform do it.
- `extract_skipped` vs `extract_hard_dropped` — same lost memory, different repair. A
skip means a flag rule is missing; a hard drop means a drop rule is too greedy.
- `flag_precision`, `type_accuracy`, `type_coverage`, `exclude_suppression`,
`noise_leak_count` — the rest of the picture.
`noise_leak_count` deserves a note: it regex-matches stored memory text for heartbeat
markers (`epoch`, `ETA`, `N% complete`, `chunks processed`, `pid NNNN`, …). Any hit
means the class that ate 20.5% of the v1 corpus has found a new way through, even if
every count-based score looks fine.
Raising a floor is cheap and should be done once a score has held above the new bar for
a while. Lowering one is a decision that belongs in a PR description, next to the reason.
## Baseline (design validation)
Established against the live platform during the v2 design phase, using the same
fixture wording:
- **Exclude classes: 5 of 6 suppressed by the custom instructions alone.** Training
heartbeats, chunk-progress frames, file-modification lists, session-only status, and
assistant narration all stored nothing.
- **The repo-file class is client-side only.** A pasted `CLAUDE.md` excerpt produced
three confident "preference" memories. This is not fixable by prompting: a pasted
convention and a stated convention are textually identical, so the extractor is right
to store it and the client must never send it. `triggers.py` owns this rule, and it is
mandatory rather than tunable.
- **Extract classes: 5 of 5 captured** — preference, decision, convention, insight and
runbook — plus the mixed window, where the bastion-host fact was stored and neither
surrounding progress line was.
- **Categories were empty on every memory** at read time, which is what fixed the read
path on `metadata.type` and this harness with it.
Write latency was 0.38–0.51s for the fire-and-forget `add`; extraction landed 20s–5min
later. That gap is the whole reason the live mode polls.
## Adding a fixture
1. Append it to `_DROP`, `_EXCLUDE` or `_EXTRACT` in `fixtures.py` with a fresh `id`.
2. Use real wording. A fixture invented to be easy to classify measures nothing —
prefer a window copied from an actual session or an audit.
3. Write a `note` saying which failure the fixture guards. Six months from now that
sentence is the only thing standing between a red score and someone "fixing" it by
deleting the fixture.
4. `extract` fixtures need an `expect_type` from
`mem0_agent.config.project_config.TYPES`; `drop` and `exclude` fixtures must have
`expect_type = None`.
5. Run `PYTHONPATH=src python3 -m pytest tests/test_fixtures.py -q` — it enforces the
schema, unique ids, valid labels and types, and that every durable type is covered.
6. Re-run `--offline` and, when the change touches the instructions, `--live`.
Adding a fixture usually lowers a score. That is the point: the score was previously
measuring a smaller world.
+702
View File
@@ -0,0 +1,702 @@
"""Labeled windows for the write-gate evaluation.
Every fixture is one turn-window exactly as the client would hand it to the capture
path. The wording of the `drop` and `exclude` entries is taken from the audited v1
corpus -- these are the classes that made 20.5% of that corpus near-duplicate
heartbeats and drove organic searches from 257/month down to 75/month. The `extract`
entries are the knowledge the gate must never throw away.
Labels
------
drop The client's local trigger rules must hard-drop the window. It never reaches
the platform, so it costs nothing and can never be stored. Failing to drop
these is what produced the v1 heartbeat corpus.
exclude The window may legitimately be sent (a local rule cannot cheaply tell it
apart from real content), but the platform's custom instructions must store
NOTHING from it. This is the layer that catches narration, activity logs,
repo-file contents and one-off directives.
extract The window must produce at least one memory, ideally of `expect_type`.
Losing these is the expensive failure: the gate gets quiet and useless.
Schema (every entry, exactly these keys)
----------------------------------------
id stable identifier, also the per-fixture user_id suffix in live runs
window list of {"role", "content"} messages
label "drop" | "exclude" | "extract"
expect_type one of mem0_agent.config.project_config.TYPES, or None
note why this fixture exists / what regression it guards
Adding a fixture: append it to the right block, give it a fresh id, and say in `note`
which real failure it represents. tests/test_fixtures.py enforces the schema.
"""
from __future__ import annotations
LABELS: tuple[str, ...] = ("drop", "exclude", "extract")
# ---------------------------------------------------------------------------
# DROP -- mechanical noise. The client must never send these.
# ---------------------------------------------------------------------------
_DROP: list[dict] = [
{
"id": "d01_train_epoch_eta",
"window": [
{"role": "assistant", "content": "Task notification (task-id bukn4vw5n): v4 train metrics at epoch "
"0.7381/2 (37% complete) with loss 0.4727, gradient norm 0.4716, "
"ETA 124 minutes."},
{"role": "assistant", "content": "Still training. Next update in about 10 minutes."},
],
"label": "drop",
"expect_type": None,
"note": "The single most duplicated shape in the audited v1 corpus: epoch/loss/ETA training heartbeat.",
},
{
"id": "d02_train_epoch_eta_later",
"window": [
{"role": "assistant", "content": "Task notification (task-id bukn4vw5n): v4 train metrics at epoch "
"0.9124/2 (46% complete) with loss 0.4412, gradient norm 0.5031, "
"ETA 101 minutes."},
],
"label": "drop",
"expect_type": None,
"note": "Same shape as d01 twenty minutes later. Repeated same-shape turns are the near-duplicate engine.",
},
{
"id": "d03_train_step_metrics",
"window": [
{"role": "assistant", "content": "step 4200/12000 | loss 0.3318 | lr 1.2e-05 | grad_norm 0.61 | "
"throughput 812 tok/s | eta 02:41:15"},
],
"label": "drop",
"expect_type": None,
"note": "Bare training metric line with no prose at all -- pure telemetry.",
},
{
"id": "d04_chunks_progress",
"window": [
{"role": "assistant", "content": "Progress for task bnzbd1uay: 218 of 928 chunks processed "
"(23% complete), approximately 5,141 synthetic memories generated, "
"11 chunk failures, ETA about 55 minutes."},
],
"label": "drop",
"expect_type": None,
"note": "Verbatim v1 corpus wording: the 'N of M chunks processed (X%)' heartbeat.",
},
{
"id": "d05_chunks_progress_mid",
"window": [
{"role": "assistant", "content": "Progress for task bnzbd1uay: 466 of 928 chunks processed "
"(50% complete), approximately 10,884 synthetic memories generated, "
"19 chunk failures, ETA about 31 minutes."},
],
"label": "drop",
"expect_type": None,
"note": "Second emission of d04. Differs only in the numbers -- the classic near-duplicate pair.",
},
{
"id": "d06_chunks_progress_late",
"window": [
{"role": "assistant", "content": "Progress for task bnzbd1uay: 902 of 928 chunks processed "
"(97% complete), approximately 21,330 synthetic memories generated, "
"24 chunk failures, ETA about 2 minutes."},
],
"label": "drop",
"expect_type": None,
"note": "Third emission. A gate that drops d04 but keeps this one has learned the numbers, not the shape.",
},
{
"id": "d07_markdown_ingest_pid",
"window": [
{"role": "assistant", "content": "markdown ingest still running (pid 48213, elapsed 00:14:52); "
"3,204 files indexed so far, 0 errors."},
],
"label": "drop",
"expect_type": None,
"note": "Markdown-ingest heartbeat with pid/elapsed -- second-largest duplicate cluster in the audit.",
},
{
"id": "d08_markdown_ingest_pid_repeat",
"window": [
{"role": "assistant", "content": "markdown ingest still running (pid 48213, elapsed 00:29:18); "
"6,771 files indexed so far, 2 errors."},
],
"label": "drop",
"expect_type": None,
"note": "Repeat of d07. Same pid, later elapsed.",
},
{
"id": "d09_ingest_complete_stats",
"window": [
{"role": "assistant", "content": "markdown ingest finished (pid 48213, elapsed 00:41:07): "
"9,118 files indexed, 2 errors, 0 skipped."},
],
"label": "drop",
"expect_type": None,
"note": "Terminal heartbeat. Completion counts are still telemetry, not knowledge.",
},
{
"id": "d10_monitor_status_line",
"window": [
{"role": "assistant", "content": "monitor: api p50 118ms p99 640ms | queue depth 3 | workers 8/8 "
"healthy | last deploy 41m ago"},
],
"label": "drop",
"expect_type": None,
"note": "Monitoring status line. True for one instant, useless in any future session.",
},
{
"id": "d11_monitor_all_green",
"window": [
{"role": "assistant", "content": "Health check at 09:14:03 - all green. Nothing to do."},
],
"label": "drop",
"expect_type": None,
"note": "Monitoring no-op turn.",
},
{
"id": "d12_tool_only_turn",
"window": [
{"role": "assistant", "content": "<tool_use name=\"Bash\">git status --porcelain</tool_use>"},
{"role": "user", "content": "<tool_result> M src/mem0_agent/api.py\n M tests/test_api.py</tool_result>"},
],
"label": "drop",
"expect_type": None,
"note": "Tool-only turn: no natural-language content from either party.",
},
{
"id": "d13_tool_only_test_output",
"window": [
{"role": "user", "content": "<tool_result>============ 412 passed, 3 skipped in 38.21s "
"============</tool_result>"},
],
"label": "drop",
"expect_type": None,
"note": "Raw tool output with no interpretation. The lesson, if any, comes in a later turn.",
},
{
"id": "d14_tool_only_ls",
"window": [
{"role": "assistant", "content": "<tool_use name=\"Bash\">ls -la integrations/mem0-agent</tool_use>"},
{"role": "user", "content": "<tool_result>total 0\ndrwxr-xr-x docs\ndrwxr-xr-x eval\ndrwxr-xr-x "
"hooks\ndrwxr-xr-x src\ndrwxr-xr-x tests</tool_result>"},
],
"label": "drop",
"expect_type": None,
"note": "Directory listing round-trip. Derivable from the repo at any time.",
},
{
"id": "d15_eta_only",
"window": [
{"role": "assistant", "content": "Next update in about 10 minutes."},
],
"label": "drop",
"expect_type": None,
"note": "The shortest heartbeat there is. v1 stored dozens of these.",
},
{
"id": "d16_percent_bar",
"window": [
{"role": "assistant", "content": "[####################........] 71% | 1,412/1,988 rows migrated | "
"eta 6m"},
],
"label": "drop",
"expect_type": None,
"note": "Progress bar render. Mechanical by construction.",
},
{
"id": "d17_still_running_ack",
"window": [
{"role": "assistant", "content": "Still running. 41% now."},
{"role": "user", "content": "ok"},
],
"label": "drop",
"expect_type": None,
"note": "Heartbeat plus bare acknowledgement. Nothing durable can be extracted from 'ok'.",
},
{
"id": "d18_backfill_job_notification",
"window": [
{"role": "assistant", "content": "Task notification (task-id q7z1m4d0c): backfill job 'memories_v3' "
"is 62% complete, 3.1M of 5.0M rows, ETA 22 minutes."},
],
"label": "drop",
"expect_type": None,
"note": "Job-notification wrapper, a different job than d01/d04 but the same shape.",
},
]
# ---------------------------------------------------------------------------
# EXCLUDE -- may be sent; the platform instructions must store nothing.
# ---------------------------------------------------------------------------
_EXCLUDE: list[dict] = [
{
"id": "x01_narration_browser_test",
"window": [
{"role": "assistant", "content": "Should I drive a browser test (log in, open the agent, send a "
"message), or would you rather click through the UI yourself?"},
{"role": "user", "content": "you do it"},
],
"label": "exclude",
"expect_type": None,
"note": "The v1 corpus turned this into 'user prefers the assistant to drive browser tests'. It is a "
"one-off directive plus assistant narration, not a preference.",
},
{
"id": "x02_narration_plan",
"window": [
{"role": "assistant", "content": "I'm going to start by reading api.py and settings.py, then sketch "
"the change, then run the tests before I touch anything else."},
],
"label": "exclude",
"expect_type": None,
"note": "Assistant stating its own plan mid-task. Attribution trap: this is not the user's preference.",
},
{
"id": "x03_narration_asked_whether",
"window": [
{"role": "assistant", "content": "I asked whether to keep the old endpoint around during the "
"migration and you said you would think about it."},
],
"label": "exclude",
"expect_type": None,
"note": "Narration of an unresolved exchange. No decision was reached, so nothing durable exists yet.",
},
{
"id": "x04_file_modification_list",
"window": [
{"role": "assistant", "content": "I modified VERSION, chat.py, agent.py, types.py, chunking.py, the "
"slack adapter, the router, the tests, and several web components "
"in this session."},
],
"label": "exclude",
"expect_type": None,
"note": "Verbatim v1 wording. Activity that git already records, in higher fidelity, forever.",
},
{
"id": "x05_commit_list",
"window": [
{"role": "assistant", "content": "Committed 4 changes: e102b21 foundation, b357a5a release bump, "
"d653b63 milvus guard, cc46715 cassandra filters."},
],
"label": "exclude",
"expect_type": None,
"note": "Commit log restated in prose.",
},
{
"id": "x06_pr_activity",
"window": [
{"role": "assistant", "content": "Opened PR #6589 against main and requested review from two "
"teammates; CI is running now."},
],
"label": "exclude",
"expect_type": None,
"note": "PR bookkeeping, derivable from the forge.",
},
{
"id": "x07_claude_md_excerpt",
"window": [
{"role": "user", "content": "## Coding Standards\n\n- Python source files: snake_case.py\n"
"- Test files: test_<module>.py\n- Ruff line length 120\n"
"(this is the contents of our CLAUDE.md)"},
],
"label": "exclude",
"expect_type": None,
"note": "Repo-file contents. This is the one exclude class the platform instructions alone do NOT "
"suppress -- it reads as genuine convention. The client must drop repo-file pastes locally.",
},
{
"id": "x08_readme_excerpt",
"window": [
{"role": "user", "content": "From the README:\n\n## Installation\n\n```bash\npip install mem0ai\n```\n"
"\n## Quickstart\n\n```python\nfrom mem0 import Memory\nm = Memory()\n```"},
],
"label": "exclude",
"expect_type": None,
"note": "README excerpt. Already in the repo; storing it duplicates a file that will drift.",
},
{
"id": "x09_config_file_excerpt",
"window": [
{"role": "user", "content": "Here is our pyproject:\n\n[tool.ruff]\nline-length = 120\n"
"target-version = \"py310\"\n\n[tool.pytest.ini_options]\n"
"testpaths = [\"tests\"]"},
],
"label": "exclude",
"expect_type": None,
"note": "Config file paste. Same class as x07/x08 -- the file is the source of truth, not memory.",
},
{
"id": "x10_one_off_you_do_it",
"window": [
{"role": "assistant", "content": "Do you want to run the migration, or should I?"},
{"role": "user", "content": "you do it"},
],
"label": "exclude",
"expect_type": None,
"note": "One-off task directive. v1 generalized these into standing preferences.",
},
{
"id": "x11_one_off_skip_tests",
"window": [
{"role": "user", "content": "skip tests for now, I just want to see if it compiles"},
],
"label": "exclude",
"expect_type": None,
"note": "Scoped to this moment. Storing it as a preference would suppress tests forever.",
},
{
"id": "x12_one_off_run_yourself",
"window": [
{"role": "user", "content": "run it yourself this time, I'm on a call"},
],
"label": "exclude",
"expect_type": None,
"note": "Explicitly a one-time instruction ('this time'), the exact wording the instructions call out.",
},
{
"id": "x13_session_status_time",
"window": [
{"role": "assistant", "content": "At 12:45:24 the job was wrapping up; I'll give you a final summary "
"when it completes."},
],
"label": "exclude",
"expect_type": None,
"note": "Verbatim v1 wording. Session-only status with a wall-clock timestamp.",
},
{
"id": "x14_session_only_step",
"window": [
{"role": "assistant", "content": "We're on step 3 of 5 of the migration right now; the last two "
"steps are the index rebuild and the cutover."},
],
"label": "exclude",
"expect_type": None,
"note": "Where we are in this session. Session state has its own record and TTL; it is not knowledge.",
},
{
"id": "x15_assistant_self_attribution",
"window": [
{"role": "assistant", "content": "I prefer to read the whole file before editing so I don't miss "
"context, so that's what I did here."},
],
"label": "exclude",
"expect_type": None,
"note": "Attribution trap: the assistant's own habit is not the developer's preference.",
},
{
"id": "x16_credentials",
"window": [
{"role": "user", "content": "here's the staging key so you can test: "
"MEM0_API_KEY=m0-abc123FAKEnotreal456 and the DSN is "
"postgres://app:hunter2@staging-db:5432/app"},
],
"label": "exclude",
"expect_type": None,
"note": "Secrets must never be stored, however useful they look. Explicit exclusion in the instructions.",
},
]
# ---------------------------------------------------------------------------
# EXTRACT -- durable knowledge. At least one memory, ideally of expect_type.
# ---------------------------------------------------------------------------
_EXTRACT: list[dict] = [
{
"id": "e01_pref_test_output_first",
"window": [
{"role": "user", "content": "Stop dumping the whole diff at me every time. Show me the failing test "
"output first, then the fix. That's how I want it from now on."},
{"role": "assistant", "content": "Understood - failing test output first, then the fix."},
],
"label": "extract",
"expect_type": "preference",
"note": "Explicit standing preference ('from now on'). Baseline: captured.",
},
{
"id": "e02_pref_no_summary_tables",
"window": [
{"role": "user", "content": "In general, don't end your answers with a summary table. Just tell me "
"what changed in two sentences. Applies to every task, not just this one."},
],
"label": "extract",
"expect_type": "preference",
"note": "Communication preference, stated as a general rule -- the discriminator against x11/x12.",
},
{
"id": "e03_pref_ask_before_force_push",
"window": [
{"role": "user", "content": "Rule for me, always: never force-push a shared branch without asking "
"first. I've been burned by that twice."},
],
"label": "extract",
"expect_type": "preference",
"note": "Workflow guardrail with stated motivation.",
},
{
"id": "e04_pref_pnpm_only",
"window": [
{"role": "assistant", "content": "Should I use npm install here?"},
{"role": "user", "content": "No - I always use pnpm for anything TypeScript, in every repo. npm and "
"yarn produce lockfile churn I then have to clean up."},
],
"label": "extract",
"expect_type": "preference",
"note": "Tooling preference scoped to the developer, not the repo. Should land at user scope.",
},
{
"id": "e05_decision_pgvector",
"window": [
{"role": "user", "content": "Let's go with pgvector instead of Pinecone. I don't want a second "
"vendor to manage, and the latency is fine at our scale."},
{"role": "assistant", "content": "Going with pgvector then, for vendor consolidation."},
],
"label": "extract",
"expect_type": "decision",
"note": "Resolved choice plus reasoning. Baseline: captured (as two memories).",
},
{
"id": "e06_decision_arq_over_celery",
"window": [
{"role": "user", "content": "We're dropping Celery and moving the workers to arq. Celery's redis "
"broker config kept drifting between environments and arq is asyncio "
"native, which matches the rest of the service."},
],
"label": "extract",
"expect_type": "decision",
"note": "Migration decision with two stated reasons.",
},
{
"id": "e07_decision_metadata_over_categories",
"window": [
{"role": "assistant", "content": "We could filter reads on categories or on metadata.type."},
{"role": "user", "content": "Filter on metadata.type. Categories are assigned by a background job "
"hours later, so a memory written this session would be invisible to a "
"category filter."},
],
"label": "extract",
"expect_type": "decision",
"note": "The read-path decision this harness itself depends on.",
},
{
"id": "e08_decision_keep_worktrees",
"window": [
{"role": "user", "content": "We'll keep using git worktrees for parallel agent work rather than "
"branch switching - switching branches invalidates the build cache and "
"costs us four minutes every time."},
],
"label": "extract",
"expect_type": "decision",
"note": "Process decision with a measured justification.",
},
{
"id": "e09_convention_branch_naming",
"window": [
{"role": "assistant", "content": "The push was rejected: this repo's hook requires branch names in "
"the form user/<id>/<name>."},
{"role": "user", "content": "Right, that's the rule here - always name branches that way."},
],
"label": "extract",
"expect_type": "convention",
"note": "Team rule confirmed by the developer, enforced by tooling but not written down.",
},
{
"id": "e10_convention_conventional_commits",
"window": [
{"role": "user", "content": "Every commit message in this repo has to be a conventional commit with "
"the package scope, like fix(mem0-agent): .... The release router parses "
"the scope, and it isn't documented anywhere."},
],
"label": "extract",
"expect_type": "convention",
"note": "Undocumented team rule with the consequence of breaking it.",
},
{
"id": "e11_convention_no_core_deps",
"window": [
{"role": "user", "content": "Never add anything to the core dependencies list - new deps go in an "
"optional group. That's a hard rule on this team; core has to stay "
"installable with no extras."},
],
"label": "extract",
"expect_type": "convention",
"note": "Project rule stated as a hard constraint.",
},
{
"id": "e12_insight_pytest_needs_compose",
"window": [
{"role": "assistant", "content": "Root cause found: pytest in server/ fails with a misleading "
"postgres connection error unless `docker compose up` is running "
"first. The tests need the compose stack."},
{"role": "user", "content": "good catch"},
],
"label": "extract",
"expect_type": "insight",
"note": "Root-caused gotcha with a misleading symptom. Baseline: captured.",
},
{
"id": "e13_insight_project_id_body",
"window": [
{"role": "assistant", "content": "Found it. project_id and org_id have to go in the request body - "
"as query params the API silently ignores them and the write lands "
"in whatever project the API key defaults to. No error, no warning."},
],
"label": "extract",
"expect_type": "insight",
"note": "Silent-failure constraint. Exactly the class of thing that costs an hour when forgotten.",
},
{
"id": "e14_insight_not_takes_list",
"window": [
{"role": "assistant", "content": "The 400 was the filter shape: NOT takes a list of clauses, not a "
"single object. `{\"NOT\": {...}}` is rejected, `{\"NOT\": [{...}]}` "
"works."},
],
"label": "extract",
"expect_type": "insight",
"note": "Non-obvious API shape, generalized past the incident.",
},
{
"id": "e15_insight_latest_only",
"window": [
{"role": "user", "content": "So that's why we saw duplicates - reads default to returning superseded "
"memories next to the ones that replaced them. latest_only has to be set "
"on every read or the context pack is full of stale pairs."},
],
"label": "extract",
"expect_type": "insight",
"note": "Constraint discovered from a symptom, stated as the general lesson.",
},
{
"id": "e16_runbook_release",
"window": [
{"role": "assistant", "content": "Verified the release procedure end to end: bump VERSION, run "
"`make build`, tag with `cli-v<version>`, push the tag, then the "
"release router dispatches the package workflow. Confirmed working "
"on the last release."},
],
"label": "extract",
"expect_type": "runbook",
"note": "Multi-step procedure explicitly verified end to end. Baseline: captured.",
},
{
"id": "e17_runbook_local_stack",
"window": [
{"role": "assistant", "content": "Local stack bring-up works, confirmed twice: `docker compose up` in "
"server/, wait for neo4j to report healthy on 8474, then "
"`uvicorn main:app --reload` from openmemory/api, then seed with "
"`python scripts/seed.py --demo`. Starting uvicorn before neo4j is "
"healthy fails the first request every time."},
],
"label": "extract",
"expect_type": "runbook",
"note": "Ordered procedure with the failure mode of doing it out of order.",
},
{
"id": "e18_runbook_republish",
"window": [
{"role": "user", "content": "For a re-publish, don't delete and recreate the GitHub release - "
"dispatch the package workflow by hand instead: "
"`gh workflow run <package>-cd.yml --ref refs/tags/<tag> -f tag=<tag>`. "
"We did that last week and it worked."},
],
"label": "extract",
"expect_type": "runbook",
"note": "Verified recovery procedure including the thing not to do.",
},
{
"id": "e19_mixed_insight_bastion",
"window": [
{"role": "assistant", "content": "Training at epoch 1.2/3 (40%), ETA 38 minutes."},
{"role": "user", "content": "while that runs - remember that our staging DB only accepts connections "
"through the bastion host, direct psql always times out."},
{"role": "assistant", "content": "Noted. Still training, 41% now."},
],
"label": "extract",
"expect_type": "insight",
"note": "MIXED: durable fact buried between two heartbeats. The fact must survive, the progress must "
"not. Baseline: captured, with no heartbeat leakage.",
},
{
"id": "e20_mixed_decision_queue",
"window": [
{"role": "assistant", "content": "Progress for task bnzbd1uay: 301 of 928 chunks processed (32% "
"complete), ETA about 44 minutes."},
{"role": "user", "content": "One thing while we wait: we've decided the ingest queue stays at "
"concurrency 4. Anything higher and the embedding provider starts "
"rate-limiting us, which costs more time than it saves."},
{"role": "assistant", "content": "Progress for task bnzbd1uay: 318 of 928 chunks processed (34% "
"complete), ETA about 41 minutes."},
],
"label": "extract",
"expect_type": "decision",
"note": "MIXED: decision plus reasoning sandwiched between two identical-shape progress lines.",
},
{
"id": "e21_mixed_preference_monitoring",
"window": [
{"role": "assistant", "content": "monitor: api p50 121ms p99 655ms | queue depth 2 | workers 8/8 "
"healthy"},
{"role": "user", "content": "Going forward, when something breaks, show me the smallest repro before "
"you propose a fix. I don't want the fix until I've seen the repro."},
{"role": "assistant", "content": "monitor: api p50 119ms p99 640ms | queue depth 3 | workers 8/8 "
"healthy"},
],
"label": "extract",
"expect_type": "preference",
"note": "MIXED: standing preference between two monitoring lines.",
},
{
"id": "e22_mixed_convention_ingest",
"window": [
{"role": "assistant", "content": "markdown ingest still running (pid 48213, elapsed 00:22:40); "
"5,010 files indexed so far."},
{"role": "user", "content": "Also, house rule you should know: every new vector-store provider needs "
"a test directory under tests/vector_stores/ before the PR can merge. "
"Reviewers reject without it and it's not in the contributing guide."},
{"role": "assistant", "content": "markdown ingest still running (pid 48213, elapsed 00:24:11); "
"5,402 files indexed so far."},
],
"label": "extract",
"expect_type": "convention",
"note": "MIXED: undocumented team rule between two ingest heartbeats.",
},
]
FIXTURES: list[dict] = [*_DROP, *_EXCLUDE, *_EXTRACT]
def by_label(label: str) -> list[dict]:
"""All fixtures carrying `label`."""
return [f for f in FIXTURES if f["label"] == label]
def counts() -> dict[str, int]:
"""Fixture count per label, so the harness can report its own coverage."""
return {label: len(by_label(label)) for label in LABELS}
def counts_by_type() -> dict[str, int]:
"""Extract-fixture count per expected memory type."""
out: dict[str, int] = {}
for f in by_label("extract"):
t = f["expect_type"] or "unspecified"
out[t] = out.get(t, 0) + 1
return out
def get(fixture_id: str) -> dict | None:
return next((f for f in FIXTURES if f["id"] == fixture_id), None)
def coverage_line() -> str:
c = counts()
types = ", ".join(f"{k}={v}" for k, v in sorted(counts_by_type().items()))
return (f"{len(FIXTURES)} fixtures: drop={c['drop']} exclude={c['exclude']} extract={c['extract']} "
f"({types})")
if __name__ == "__main__": # `python3 eval/fixtures.py` prints coverage
print(coverage_line())
+580
View File
@@ -0,0 +1,580 @@
#!/usr/bin/env python3
"""Write-gate evaluation harness.
v1 shipped with no way to measure extraction quality, so it degraded unnoticed for
three months: 20.5% of the corpus turned into near-duplicate heartbeats and organic
searches fell from 257/month to 75/month. Nobody noticed because nobody could look.
This harness is the thing that looks.
Two modes:
--offline (default, no network)
Runs every fixture through mem0_agent.triggers.classify and scores the client's
local rules: hard-drop recall, hard-drop precision, flag precision/recall and
per-type accuracy. Costs nothing, so it can run in CI on every commit.
--live --project-id P --org-id O
Replays the `exclude` and `extract` fixtures against a SCRATCH project through
mem0_agent.api.Api with infer=True, polls until extraction lands, reads back and
scores what the platform actually stored. This measures the half of the gate that
lives in the custom instructions and cannot be unit-tested.
Both modes print a scorecard and write eval/last_report.json. `--check` turns the run
into a gate: it exits non-zero if hard-drop recall drops below 0.95 or extract recall
below 0.80, so no edit to the trigger rules or the custom instructions ships unmeasured.
Classifier contract (offline mode)
----------------------------------
`mem0_agent.triggers.classify(window)` is expected to take the message-window list and
return a decision. The adapter below accepts every reasonable shape -- a bool, a string,
a (decision, type) pair, a dict, or an object with attributes -- because the module is
built in parallel with this one. If the module is missing, offline scoring is skipped
with a clear message rather than crashing; if it exists but returns something
unreadable, that is reported as a contract mismatch, not as a score.
"""
from __future__ import annotations
import argparse
import json
import os
import re
import sys
import time
import uuid
from pathlib import Path
from typing import Any, Callable
EVAL_DIR = Path(__file__).resolve().parent
sys.path.insert(0, str(EVAL_DIR))
sys.path.insert(0, str(EVAL_DIR.parent / "src"))
import fixtures as fx # noqa: E402
REPORT_PATH = EVAL_DIR / "last_report.json"
# --check gates. Raise them only when the measured score has been above the new bar for
# a while; lowering one is a decision that belongs in a PR description.
THRESHOLDS: dict[str, float] = {
"hard_drop_recall": 0.95,
"extract_recall": 0.80,
}
# Text that must never appear inside a stored memory. If one of these survives the gate
# the heartbeat class is leaking again, which is precisely how v1 rotted.
NOISE_PATTERNS = [
re.compile(r"\beta\b[^.]{0,20}\b\d", re.I),
re.compile(r"\bepoch\b", re.I),
re.compile(r"\d+\s*%\s*(complete|done)", re.I),
re.compile(r"chunks?\s+processed", re.I),
re.compile(r"\bpid\s*\d+", re.I),
re.compile(r"\bgradient\s+norm\b", re.I),
re.compile(r"\bqueue\s+depth\b", re.I),
re.compile(r"\belapsed\s+\d\d:\d\d", re.I),
]
def noise_leaks(text: str) -> list[str]:
return [p.pattern for p in NOISE_PATTERNS if p.search(text or "")]
# ---------------------------------------------------------------------------
# classifier adapter
# ---------------------------------------------------------------------------
# Three outcomes, not two. "drop" is a hard drop -- mechanical noise the client refuses
# to send at any level. "skip" is "nothing worth storing right now", which is a soft miss:
# harmless on an exclude window, a lost memory on an extract one. Collapsing the two
# would hide exactly the regression this harness exists to catch.
_DROP_WORDS = {"drop", "noise", "ignore", "suppress", "reject", "block", "hard_drop"}
_SKIP_WORDS = {"skip", "none", "no", "defer", "wait", "noop", "no_trigger"}
_SEND_WORDS = {"send", "capture", "flag", "flagged", "store", "keep", "extract", "accept", "allow",
"pass", "yes"}
DECISIONS = ("drop", "skip", "send", "unknown")
class ContractMismatch(RuntimeError):
"""classify() exists but neither its signature nor its return value is readable."""
def load_classifier() -> tuple[Callable | None, str]:
"""Import lazily. triggers.py is written by a parallel workstream and may not exist."""
try:
from mem0_agent import triggers # type: ignore
except ImportError as e:
return None, f"mem0_agent.triggers is not available yet ({e})"
except Exception as e: # a broken module is a different problem than a missing one
return None, f"mem0_agent.triggers failed to import: {type(e).__name__}: {e}"
fn = getattr(triggers, "classify", None)
if not callable(fn):
return None, "mem0_agent.triggers exists but has no callable classify()"
return fn, ""
def _invoke(fn: Callable, window: list[dict], level: str | None = None) -> Any:
"""Try the plausible call shapes, most likely first."""
attempts: tuple[Callable[[], Any], ...] = ()
if level:
attempts += (lambda: fn(window, level), lambda: fn(window, level=level))
attempts += (
lambda: fn(window),
lambda: fn(messages=window),
lambda: fn(window=window),
lambda: fn("\n".join(m.get("content", "") for m in window)),
)
last: Exception | None = None
for call in attempts:
try:
return call()
except TypeError as e:
last = e
raise ContractMismatch(f"classify() rejected every call shape: {last}")
def _type_of(value: Any) -> str | None:
if isinstance(value, str) and value.lower() in _known_types():
return value.lower()
return None
def _known_types() -> set[str]:
try:
from mem0_agent.config.project_config import TYPES
return set(TYPES)
except Exception:
return set(fx.counts_by_type()) | {"session_state"}
def _look(result: Any, names: tuple[str, ...]) -> Any:
for n in names:
if isinstance(result, dict):
if n in result:
return result[n]
elif hasattr(result, n):
return getattr(result, n)
return None
def normalize(result: Any) -> tuple[str, str | None]:
"""Map whatever classify() returned onto (decision, type|None). See DECISIONS."""
if result is None:
return "skip", None
if isinstance(result, bool):
return ("send" if result else "skip"), None
if isinstance(result, str):
low = result.strip().lower()
if low in _known_types():
return "send", low
if low in _DROP_WORDS:
return "drop", None
if low in _SKIP_WORDS:
return "skip", None
if low in _SEND_WORDS:
return "send", None
return "unknown", None
if isinstance(result, (tuple, list)):
if not result:
return "unknown", None
decision, _ = normalize(result[0])
mtype = _type_of(result[1]) if len(result) > 1 else None
if decision == "unknown" and mtype:
decision = "send"
return decision, mtype
mtype = _type_of(_look(result, ("type", "mtype", "memory_type", "kind")))
# A string verdict is the most expressive shape, so it wins over the booleans:
# `flagged=False` cannot tell a hard drop apart from a soft skip.
verdict = _look(result, ("action", "decision", "verdict", "outcome", "status", "result"))
if isinstance(verdict, str):
decision, vtype = normalize(verdict)
if decision != "unknown":
return decision, (mtype or vtype)
drop = _look(result, ("drop", "dropped", "is_drop", "should_drop"))
if isinstance(drop, bool) and drop:
return "drop", mtype
send = _look(result, ("send", "capture", "should_capture", "flag", "flagged", "should_send", "keep"))
if isinstance(send, bool):
return ("send" if send else ("skip" if drop is False else "drop")), mtype
if isinstance(drop, bool):
return "send", mtype
if mtype:
return "send", mtype
return "unknown", None
def reason_of(result: Any) -> str:
"""Whatever the classifier called this rule -- the most useful column in the report."""
raw = _look(result, ("reason", "rule", "why", "trigger", "explanation"))
return str(raw) if isinstance(raw, (str, int)) else ""
# ---------------------------------------------------------------------------
# scoring
# ---------------------------------------------------------------------------
def _ratio(num: int, den: int) -> float | None:
return round(num / den, 4) if den else None
def score_offline(rows: list[dict]) -> dict[str, Any]:
"""rows: {"id","label","expect_type","decision","got_type"}."""
drops = [r for r in rows if r["label"] == "drop"]
extracts = [r for r in rows if r["label"] == "extract"]
excludes = [r for r in rows if r["label"] == "exclude"]
dropped = [r for r in rows if r["decision"] == "drop"]
flagged = [r for r in rows if r["decision"] == "send"]
per_type: dict[str, dict[str, int]] = {}
for r in extracts:
want = r["expect_type"] or "unspecified"
bucket = per_type.setdefault(want, {"n": 0, "correct": 0, "typed": 0})
bucket["n"] += 1
if r["got_type"]:
bucket["typed"] += 1
if r["got_type"] == want:
bucket["correct"] += 1
typed = sum(b["typed"] for b in per_type.values())
correct = sum(b["correct"] for b in per_type.values())
return {
# The metric that keeps the corpus clean: mechanical noise must never be sent.
# capture.py only forwards windows whose action is "flag", so both "drop" and
# "skip" satisfy the contract -- the label says "never sent", not "matched a
# rule named drop".
"hard_drop_recall": _ratio(sum(1 for r in drops if r["decision"] in ("drop", "skip")), len(drops)),
# Advisory, not gated: how much of that containment comes from an explicit drop
# rule rather than from no flag rule happening to match. Noise held back only by
# the absence of a flag rule leaks the day someone adds one.
"hard_drop_explicit": _ratio(sum(1 for r in drops if r["decision"] == "drop"), len(drops)),
# Of everything hard-dropped, how much was safe to drop. A hard drop on an
# `extract` window is the one unrecoverable error -- the memory is gone and
# nothing logs a miss. Hard-dropping an `exclude` window is not counted against
# this: those are meant to be discarded, and doing it locally is simply cheaper.
"hard_drop_precision": _ratio(sum(1 for r in dropped if r["label"] != "extract"), len(dropped)),
# The metric that keeps the gate useful: durable knowledge must survive locally.
"extract_recall": _ratio(sum(1 for r in extracts if r["decision"] == "send"), len(extracts)),
# Of everything forwarded to the platform, how much was worth forwarding.
"flag_precision": _ratio(sum(1 for r in flagged if r["label"] == "extract"), len(flagged)),
# Extract windows lost to a soft "skip" rather than a hard drop. Same lost
# memory, different fix: a missing flag rule, not an over-broad drop rule.
"extract_skipped": sum(1 for r in extracts if r["decision"] == "skip"),
"extract_hard_dropped": sum(1 for r in extracts if r["decision"] == "drop"),
# Exclude windows the client suppressed locally -- free wins, not required, since
# the platform instructions are the designated owner of that class.
"exclude_suppressed_early": _ratio(
sum(1 for r in excludes if r["decision"] in ("drop", "skip")), len(excludes)),
"type_accuracy": _ratio(correct, typed),
"type_coverage": _ratio(typed, len(extracts)),
"per_type": per_type,
"unreadable": sum(1 for r in rows if r["decision"] == "unknown"),
}
def run_offline(level: str | None = None) -> dict[str, Any]:
fn, why = load_classifier()
if fn is None:
return {"mode": "offline", "skipped": True, "reason": why, "scores": {}, "rows": [],
"level": level}
rows: list[dict] = []
try:
for f in fx.FIXTURES:
try:
raw = _invoke(fn, f["window"], level)
decision, got_type = normalize(raw)
reason, err = reason_of(raw), ""
except ContractMismatch:
raise
except Exception as e: # a fixture that blows up the classifier is a real finding
decision, got_type, reason, err = "unknown", None, "", f"{type(e).__name__}: {e}"
rows.append({
"id": f["id"], "label": f["label"], "expect_type": f["expect_type"],
"decision": decision, "got_type": got_type, "reason": reason, "error": err,
})
except ContractMismatch as e:
return {"mode": "offline", "skipped": True, "reason": str(e), "scores": {}, "rows": [],
"level": level}
return {"mode": "offline", "skipped": False, "reason": "", "level": level,
"scores": score_offline(rows), "rows": rows}
# ---------------------------------------------------------------------------
# live mode
# ---------------------------------------------------------------------------
def _memory_type(mem: dict) -> str | None:
"""Read the type from metadata. NEVER from categories: categorization lags ~3.9h
(median) and is 0% for memories under an hour old, so a fresh read sees nothing."""
meta = mem.get("metadata") or {}
return meta.get("type") if isinstance(meta, dict) else None
def score_live(rows: list[dict]) -> dict[str, Any]:
extracts = [r for r in rows if r["label"] == "extract"]
excludes = [r for r in rows if r["label"] == "exclude"]
per_type: dict[str, dict[str, int]] = {}
for r in extracts:
want = r["expect_type"] or "unspecified"
bucket = per_type.setdefault(want, {"n": 0, "stored": 0, "typed_correct": 0})
bucket["n"] += 1
if r["stored"] >= 1:
bucket["stored"] += 1
if want in (r["got_types"] or []):
bucket["typed_correct"] += 1
return {
"extract_recall": _ratio(sum(1 for r in extracts if r["stored"] >= 1), len(extracts)),
"exclude_suppression": _ratio(sum(1 for r in excludes if r["stored"] == 0), len(excludes)),
"type_match": _ratio(sum(b["typed_correct"] for b in per_type.values()),
sum(b["n"] for b in per_type.values())),
"noise_leak_count": sum(len(r["leaks"]) for r in rows),
"memories_written": sum(r["stored"] for r in rows),
"per_type": per_type,
}
def run_live(args: argparse.Namespace) -> dict[str, Any]:
from mem0_agent.api import Api, results_of # noqa: PLC0415
from mem0_agent.config.filters import all_in_scope # noqa: PLC0415
from mem0_agent.config.project_config import POLICY_VERSION # noqa: PLC0415
key = args.api_key or os.environ.get("MEM0_API_KEY")
if not key:
raise SystemExit("live mode needs an API key: --api-key or MEM0_API_KEY")
runid = args.run_id or uuid.uuid4().hex[:8]
app_id = f"eval-{runid}"
api = Api(key, org_id=args.org_id, project_id=args.project_id, strict=True)
classifier, why = load_classifier()
if classifier is None:
print(f"note: {why}; metadata.type will be stamped from expect_type (self-stamped, "
f"so type_match measures the round trip only)")
targets = [f for f in fx.FIXTURES if f["label"] in ("exclude", "extract")]
print(f"live run {runid}: project={args.project_id} app_id={app_id} fixtures={len(targets)}")
sent: list[dict] = []
for f in targets:
uid = f"eval-{runid}-{f['id']}"
mtype = f["expect_type"]
stamped_by = "expect_type"
if classifier is not None:
try:
decision, got = normalize(_invoke(classifier, f["window"], args.level))
if got:
mtype, stamped_by = got, "classifier"
except Exception:
pass
meta = {"session_id": f"eval-{runid}", "editor": "eval-harness",
"policy": POLICY_VERSION, "fixture": f["id"]}
if mtype:
meta["type"] = mtype
t0 = time.time()
status, body = api.add(f["window"], user_id=uid, app_id=app_id, infer=True, metadata=meta)
ok = status in (200, 201, 202)
print(f" sent {f['id']:34s} {status} {round(time.time() - t0, 2)}s"
+ ("" if ok else f" <- {body}"))
sent.append({"id": f["id"], "label": f["label"], "expect_type": f["expect_type"],
"user_id": uid, "add_status": status, "stamped_type": mtype,
"stamped_by": stamped_by, "accepted": ok})
def read(user_id: str) -> list[dict]:
status, body = api.get_all(all_in_scope(user_id, app_id), page_size=100)
return results_of(body) if status == 200 else []
# Poll rather than sleep: extraction landed anywhere from 20s to 5min in validation,
# so any fixed wait is either wrong or wasteful. Only the extract fixtures have a
# target to poll for; a zero read on an exclude fixture is only meaningful once the
# settle window has passed, so those are read once at the end.
started = time.time()
deadline = started + args.timeout
pending = {r["id"]: r["user_id"] for r in sent if r["label"] == "extract" and r["accepted"]}
landed: dict[str, list[dict]] = {}
interval = 5.0
print(f"polling for extraction (timeout {args.timeout:.0f}s, settle {args.min_settle:.0f}s)...")
while pending and time.time() < deadline:
time.sleep(min(interval, max(1.0, deadline - time.time())))
interval = min(interval * 1.4, 20.0)
for fid, uid in list(pending.items()):
mems = read(uid)
if mems:
landed[fid] = mems
pending.pop(fid, None)
print(f" t+{int(time.time() - started)}s extract landed "
f"{len(landed)}/{len(landed) + len(pending)}")
remaining_settle = args.min_settle - (time.time() - started)
if remaining_settle > 0:
print(f" settling {int(remaining_settle)}s before the suppression read...")
time.sleep(remaining_settle)
# Final read for everything: extract fixtures may have gained a second memory since
# they first landed, and the exclude fixtures are read here for the first time.
found = {r["id"]: read(r["user_id"]) for r in sent}
rows: list[dict] = []
for r in sent:
mems = found.get(r["id"], [])
texts = [m.get("memory") or "" for m in mems]
rows.append({**r,
"stored": len(mems),
"got_types": sorted({t for t in (_memory_type(m) for m in mems) if t}),
"memories": [{"id": m.get("id"), "text": t, "type": _memory_type(m),
"categories": m.get("categories")} for m, t in zip(mems, texts)],
"leaks": [lk for t in texts for lk in noise_leaks(t)]})
report = {"mode": "live", "skipped": False, "reason": "", "run_id": runid,
"project_id": args.project_id, "app_id": app_id,
"scores": score_live(rows), "rows": rows}
if args.cleanup:
deleted = 0
for r in rows:
status, _ = api.delete_all(user_id=r["user_id"], app_id=app_id)
deleted += r["stored"] if status in (200, 202, 204) else 0
report["cleanup"] = {"deleted_scopes": len(rows), "deleted_memories": deleted}
print(f"cleanup: removed {deleted} memories across {len(rows)} scopes")
else:
report["cleanup"] = {"skipped": True,
"hint": f"delete with app_id={app_id} / user_id prefix eval-{runid}-"}
return report
# ---------------------------------------------------------------------------
# reporting
# ---------------------------------------------------------------------------
def print_scorecard(report: dict[str, Any]) -> None:
counts = fx.counts()
print()
print("=" * 72)
title = f" mem0-agent write gate -- {report['mode']} scorecard"
if report.get("level"):
title += f" (level={report['level']})"
print(title)
print("=" * 72)
print(f" coverage: {fx.coverage_line()}")
if report.get("skipped"):
print(f" SKIPPED: {report['reason']}")
print("=" * 72)
return
scores = report["scores"]
for name, value in scores.items():
if name == "per_type" or isinstance(value, dict):
continue
gate = THRESHOLDS.get(name)
shown = "n/a" if value is None else f"{value:.3f}" if isinstance(value, float) else str(value)
mark = ""
if gate is not None:
mark = " FAIL" if (value is None or value < gate) else " ok"
shown += f" (min {gate:.2f}){mark}"
print(f" {name:24s} {shown}")
per_type = scores.get("per_type") or {}
if per_type:
print(" per expected type:")
for t, b in sorted(per_type.items()):
body = " ".join(f"{k}={v}" for k, v in b.items())
print(f" {t:14s} {body}")
if report["mode"] == "offline":
soft = [r for r in report["rows"] if r["label"] == "drop" and r["decision"] == "skip"]
if soft:
print(f" noise contained only by the absence of a flag rule ({len(soft)}, advisory):")
for r in soft:
print(f" {r['id']:38s} {r.get('reason', '')}")
bad = [r for r in report["rows"]
if (r["label"] == "drop" and r["decision"] not in ("drop", "skip"))
or (r["label"] == "extract" and r["decision"] != "send")
or r["decision"] == "unknown"]
if bad:
print(f" misclassified ({len(bad)}):")
for r in bad:
want = "drop" if r["label"] == "drop" else "send"
detail = r.get("error") or r.get("reason") or ""
print(f" {r['id']:38s} want={want:5s} got={r['decision']:8s} {detail}")
else:
for r in report["rows"]:
want = "0" if r["label"] == "exclude" else ">=1"
ok = (r["stored"] == 0) if r["label"] == "exclude" else (r["stored"] >= 1)
print(f" [{'PASS' if ok else 'FAIL'}] {r['id']:34s} stored={r['stored']} want={want} "
f"types={r['got_types']}")
for m in r["memories"]:
print(f" -> {(m['text'] or '')[:96]}")
if scores.get("noise_leak_count"):
print(f" NOISE LEAK: {scores['noise_leak_count']} stored memories match heartbeat patterns")
print(f" labels: {counts}")
print("=" * 72)
def check(report: dict[str, Any]) -> tuple[bool, list[str]]:
"""A gate that cannot measure must not pass."""
if report.get("skipped"):
return False, [f"nothing was measured: {report['reason']}"]
failures = []
scores = report["scores"]
applicable = {k: v for k, v in THRESHOLDS.items() if k in scores}
if not applicable:
return False, [f"no gated metric present in {report['mode']} scores"]
for name, floor in applicable.items():
value = scores.get(name)
if value is None:
failures.append(f"{name} not measured (threshold {floor})")
elif value < floor:
failures.append(f"{name}={value:.3f} below threshold {floor}")
return (not failures), failures
def main(argv: list[str] | None = None) -> int:
p = argparse.ArgumentParser(description="mem0-agent write-gate evaluation")
mode = p.add_mutually_exclusive_group()
mode.add_argument("--offline", action="store_true", help="score the local trigger rules (default)")
mode.add_argument("--live", action="store_true", help="replay fixtures against a scratch project")
p.add_argument("--project-id", help="SCRATCH project id (required for --live)")
p.add_argument("--org-id", help="org id (required for --live)")
p.add_argument("--api-key", help="defaults to $MEM0_API_KEY")
p.add_argument("--run-id", help="override the generated run id")
p.add_argument("--timeout", type=float, default=360.0, help="live: max seconds to poll (default 360)")
p.add_argument("--min-settle", type=float, default=60.0,
help="live: minimum seconds before a zero read counts as suppression (default 60)")
p.add_argument("--level", help="capture aggressiveness passed to classify() "
"(conservative|balanced|aggressive); default is the module's own")
p.add_argument("--cleanup", action="store_true", help="live: delete everything this run wrote")
p.add_argument("--check", action="store_true", help="exit non-zero if scores regress below thresholds")
p.add_argument("--report", default=str(REPORT_PATH), help=f"report path (default {REPORT_PATH})")
args = p.parse_args(argv)
if args.live:
# A live run writes real memories. It must never be able to land in production
# by omission, so both ids are required and neither has a default.
missing = [n for n, v in (("--project-id", args.project_id), ("--org-id", args.org_id)) if not v]
if missing:
p.error("live mode requires " + " and ".join(missing)
+ " -- point them at a scratch project, never a production one")
report = run_live(args)
else:
report = run_offline(args.level)
report["thresholds"] = THRESHOLDS
report["fixture_counts"] = fx.counts()
report["fixture_counts_by_type"] = fx.counts_by_type()
report["generated_at"] = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())
passed, failures = check(report)
report["check"] = {"passed": passed, "failures": failures}
print_scorecard(report)
Path(args.report).write_text(json.dumps(report, indent=1))
print(f"report written to {args.report}")
if args.check:
if failures:
print("CHECK FAILED:")
for f in failures:
print(f" - {f}")
return 1
print("CHECK PASSED")
return 0
if __name__ == "__main__":
sys.exit(main())
@@ -108,7 +108,9 @@ def observe(ctx, window: list[dict], level: str | None = None) -> TriggerResult:
try:
buf.note_shape(window)
if result.action == "flag" and result.mtype:
buf.append(window, result.mtype, result.reason)
# Buffer the FILTERED window: classify strips noise turns, so a durable fact
# that sat between two progress lines is sent without them.
buf.append(result.payload(window), result.mtype, result.reason)
ctx.log("capture_observe", action=result.action, mtype=result.mtype, reason=result.reason)
except Exception:
pass
@@ -24,7 +24,7 @@ from __future__ import annotations
import hashlib
import re
from dataclasses import dataclass
from dataclasses import dataclass, field
from typing import Any, Callable, Sequence
# --------------------------------------------------------------------------
@@ -177,11 +177,20 @@ class TriggerResult:
action: str # "drop" | "skip" | "flag"
mtype: str | None
reason: str
window: tuple[Any, ...] = field(default=(), compare=False)
"""The turns worth sending. Noise turns are filtered out, so a durable fact sitting
between two progress lines survives instead of being dropped with them.
Excluded from equality: it is the payload, not the verdict.
"""
@property
def flagged(self) -> bool:
return self.action == "flag"
def payload(self, original: Sequence[Any]) -> list[Any]:
return list(self.window) if self.window else list(original or [])
# --------------------------------------------------------------------------
# HARD DROP rules -- applied at every level
@@ -382,6 +391,7 @@ REMEMBER_RULE = Rule(
_rx(
r"\bremember (this|that|to|:)",
r"\b(please )?remember\b[^.\n]{0,40}\bfor (next time|the future|future sessions)\b",
r"\bremember\s+(that|this|the|to|about)\b",
r"\bdon'?t forget\b",
r"\bdo not forget\b",
r"\bnote that\b",
@@ -508,7 +518,10 @@ RUNBOOK_RULE = Rule(
),
predicate=_verified_procedure,
mtype="runbook",
min_level="aggressive",
# A procedure the user states they VERIFIED is strong, specific and among the most
# useful things to recall, so it lands at balanced. `aggressive` remains for completed
# goals and procedures the assistant merely proposes (COMPLETED_GOAL_RULE).
min_level="balanced",
)
COMPLETED_GOAL_RULE = Rule(
@@ -519,14 +532,98 @@ COMPLETED_GOAL_RULE = Rule(
)
# Order matters: the first match wins, so the most specific intent leads.
# --------------------------------------------------------------------------
# Widened rules, added after the eval harness showed 13 plainly durable windows
# falling through as `no_trigger`. Each pattern below is traceable to a fixture in
# eval/fixtures.py; re-run `eval/run.py --offline` after touching any of them.
# --------------------------------------------------------------------------
STATED_RULE_RULE = Rule(
"stated_rule",
_rx(
r"\brule for me\b",
r"\b(that'?s|this is) a hard rule\b",
r"\bhouse rule\b",
r"\bhard rule (here|for)\b",
r"\bnever (add|put|place|introduce|merge|force-?push)\b",
r"\bevery new \w+[\w\s-]{0,30} (needs|requires|must)\b",
r"\bin general,? (don'?t|do not|never|always)\b",
r"\balways:? never\b",
),
mtype="convention",
min_level="balanced",
)
HABITUAL_PREFERENCE_RULE = Rule(
"habitual_preference",
_rx(
r"\bi always (use|run|want|prefer|do)\b",
r"\bi never (use|run|want|do)\b",
r"\bwe only use\b",
r"\bjust tell me\b",
r"\bdon'?t end your (answers?|responses?)\b",
r"\bapply (that|this) (everywhere|to every|going forward)\b",
),
mtype="preference",
min_level="conservative",
scope="user",
)
CHOICE_RULE = Rule(
"explicit_choice",
_rx(
r"\bwe'?(ve)? decided\b",
r"\bwe'?(re| are) (dropping|moving|switching|migrating)\b",
r"\bwe'?ll keep (using|the)\b",
r"\bwe'?(re| are) keeping\b",
r"\brather than\b[^.\n]{0,120}\b(because|since)\b",
r"\b(use|filter on|go with) \w[\w.\-]* (instead of|over) \w[\w.\-]*",
r"\bstays at\b[^.\n]{0,60}\banything higher\b",
),
mtype="decision",
min_level="balanced",
)
DIAGNOSIS_RULE = Rule(
"diagnosis",
_rx(
r"\bthe \d{3} was\b",
r"\bis rejected\b",
r"\bwas the (filter|config|shape|schema|encoding|ordering)\b",
r"\btakes a list\b",
r"\bare assigned by a background job\b",
r"\bwon'?t (match|return|work) (unless|until|without)\b",
),
mtype="insight",
min_level="balanced",
)
PROCEDURE_RULE = Rule(
"verified_procedure_phrasing",
_rx(
r"\bverified .{0,40}\bend to end\b",
r"\bconfirmed (twice|three times|repeatedly)\b",
r"\bworks,? confirmed\b",
r"\bfor a (re-?publish|re-?deploy|rollback|re-?run)\b[^.\n]{0,80}\b(dispatch|run|use)\b",
r"\bbring-?up works\b",
r"\bthe steps? (are|were)\b[^.\n]{0,40}:",
),
mtype="runbook",
min_level="balanced",
)
FLAG_RULES: list[Rule] = [
REMEMBER_RULE,
CORRECTION_RULE,
STANDING_PREFERENCE_RULE,
HABITUAL_PREFERENCE_RULE,
DECISION_RULE,
CHOICE_RULE,
INSIGHT_RULE,
DIAGNOSIS_RULE,
CONVENTION_RULE,
STATED_RULE_RULE,
RUNBOOK_RULE,
PROCEDURE_RULE,
COMPLETED_GOAL_RULE,
]
@@ -589,25 +686,57 @@ def classify(
value = level_rank(level)
# Rules that only mean something across a whole window run first, on the original turns.
for rule in HARD_DROP_RULES:
if rule.matches(turns, _RANK["aggressive"]): # hard drops ignore the level
if rule.name in _WINDOW_LEVEL_DROPS and rule.matches(turns, _RANK["aggressive"]):
return TriggerResult("drop", None, rule.name)
# Noise is removed turn by turn, not window by window. Dropping a whole window because
# one line in it was a progress update is how a durable fact gets lost: during live
# validation a window of [progress, "the staging DB only accepts the bastion host",
# progress] correctly yielded the bastion fact, and the client must not pre-empt that.
kept: list[Any] = []
dropped_reasons: list[str] = []
for turn in turns:
reason = _turn_drop_reason([turn])
if reason:
dropped_reasons.append(reason)
else:
kept.append(turn)
if not kept:
return TriggerResult("drop", None, dropped_reasons[0] if dropped_reasons else "noise")
# Window-level drops that only make sense across turns.
if recent_shapes and shape_signature(turns) in set(recent_shapes):
return TriggerResult("drop", None, "repeated_shape")
reason = repo_content_reason(turns)
reason = repo_content_reason(kept)
if reason:
return TriggerResult("drop", None, reason)
if natural_words(window_text(turns)) < 4:
if natural_words(window_text(kept)) < 4:
return TriggerResult("skip", None, "no_prose")
for rule in FLAG_RULES:
if rule.matches(turns, value):
if rule.matches(kept, value):
mtype = rule.mtype
if rule.name == "remember_intent":
mtype = _refine_remember_type(turns)
return TriggerResult("flag", mtype, rule.name)
mtype = _refine_remember_type(kept)
return TriggerResult("flag", mtype, rule.name, tuple(kept))
return TriggerResult("skip", None, "no_trigger")
def _turn_drop_reason(one_turn: Sequence[Any]) -> str | None:
"""Name of the hard-drop rule this single turn matches, if any."""
for rule in HARD_DROP_RULES:
if rule.name in _WINDOW_LEVEL_DROPS:
continue
if rule.matches(one_turn, _RANK["aggressive"]): # hard drops ignore the level
return rule.name
return None
# Rules whose meaning depends on seeing the whole window, so they are not applied per turn.
_WINDOW_LEVEL_DROPS = {"repeated_shape_in_window"}
@@ -0,0 +1,113 @@
"""Mixed windows: a durable fact next to mechanical noise.
The realistic case, and the one the first implementation got wrong. Live validation showed
a window of [progress, "the staging DB only accepts the bastion host", progress] correctly
yields the bastion fact -- so the client must strip the noise turns, not discard the window.
"""
import pytest
from mem0_agent import capture
from mem0_agent.settings import SessionState, Settings
from mem0_agent.triggers import classify
PROGRESS = "Progress for task bnzbd1uay: 301 of 928 chunks processed (32% complete), ETA about 44 minutes."
HEARTBEAT = "markdown ingest still running (pid 48213, elapsed 00:22:40); 5,010 files indexed so far."
DURABLE = "One thing while we wait: we've decided the ingest queue stays at concurrency 4. Anything higher and the DB starts timing out."
BASTION = "Also remember the staging database only accepts connections through the bastion host; direct psql always times out."
class FakeApi:
def __init__(self):
self.added = []
def add(self, messages, **kw):
self.added.append({"messages": messages, **kw})
return (200, {"event_id": "e", "status": "PENDING"})
def get_all(self, filters, **kw):
return (200, {"results": []})
def update(self, mid, **kw):
return (200, {})
class FakeCtx:
def __init__(self, tmp_path):
self.api = FakeApi()
self.settings = Settings(data={"capture": "balanced"}, path=tmp_path / "s.json")
self.state = SessionState("sess-mixed", root=tmp_path / "sessions")
self.user_id, self.app_id = "dev", "acme-repo"
self.session_id, self.branch, self.ready, self.reason = "sess-mixed", "main", True, ""
def provenance(self, mtype):
return {"type": mtype, "session_id": self.session_id}
def log(self, *a, **k):
pass
@pytest.fixture
def ctx(tmp_path):
return FakeCtx(tmp_path)
def u(text):
return {"role": "user", "content": text}
def a(text):
return {"role": "assistant", "content": text}
MIXED = [a(PROGRESS), u(DURABLE), a(HEARTBEAT)]
def test_mixed_window_is_not_dropped():
result = classify(MIXED, "balanced")
assert result.action == "flag", "the durable fact must survive its noisy neighbours"
assert result.mtype == "decision"
def test_noise_turns_are_stripped_from_the_payload():
result = classify(MIXED, "balanced")
sent = [t["content"] for t in result.payload(MIXED)]
assert DURABLE in sent
assert PROGRESS not in sent and HEARTBEAT not in sent
def test_pure_noise_window_still_drops():
assert classify([a(PROGRESS), a(HEARTBEAT)], "balanced").action == "drop"
def test_only_the_durable_turn_reaches_the_api(ctx):
capture.observe(ctx, MIXED, "balanced")
capture.flush(ctx)
assert len(ctx.api.added) == 1
body = " ".join(m["content"] for m in ctx.api.added[0]["messages"])
assert "concurrency 4" in body
for noise in ("928 chunks", "pid 48213", "ETA about"):
assert noise not in body
def test_remember_intent_survives_noise(ctx):
window = [a(HEARTBEAT), u(BASTION), a(PROGRESS)]
result = capture.observe(ctx, window, "balanced")
assert result.action == "flag"
capture.flush(ctx)
body = " ".join(m["content"] for m in ctx.api.added[0]["messages"])
assert "bastion" in body.lower()
assert "pid 48213" not in body
def test_a_window_of_noise_plus_prose_without_a_trigger_is_skipped():
"""Stripping noise must not turn an ordinary exchange into a memory."""
window = [a(PROGRESS), u("what does that number mean?"), a("It is the chunk count.")]
assert classify(window, "balanced").action == "skip"
def test_window_level_repeat_detection_still_applies():
"""Three identically-shaped turns in one window is noise regardless of filtering."""
repeated = [a(f"Progress for task abc: {i} of 928 chunks processed ({i}% complete), ETA {i} minutes.")
for i in (11, 12, 13)]
assert classify(repeated, "aggressive").action == "drop"
@@ -353,9 +353,11 @@ def test_decision_is_gated_to_balanced_and_up():
assert classify(DECISION_WINDOW, "aggressive").mtype == "decision"
def test_runbook_is_gated_to_aggressive_only():
def test_verified_runbook_is_captured_from_balanced_up():
"""A procedure the user says they verified is durable knowledge, not a stretch goal.
`aggressive` is for completed goals and procedures the assistant merely proposes."""
assert classify(RUNBOOK_WINDOW, "conservative").action == "skip"
assert classify(RUNBOOK_WINDOW, "balanced").action == "skip"
assert classify(RUNBOOK_WINDOW, "balanced").mtype == "runbook"
assert classify(RUNBOOK_WINDOW, "aggressive").mtype == "runbook"