diff --git a/.gitignore b/.gitignore index 568da3757..a524a87c2 100644 --- a/.gitignore +++ b/.gitignore @@ -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 diff --git a/integrations/mem0-agent/eval/README.md b/integrations/mem0-agent/eval/README.md new file mode 100644 index 000000000..d0c944b1c --- /dev/null +++ b/integrations/mem0-agent/eval/README.md @@ -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--`, under app id + `eval-`. 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. diff --git a/integrations/mem0-agent/eval/fixtures.py b/integrations/mem0-agent/eval/fixtures.py new file mode 100644 index 000000000..0d1d8499d --- /dev/null +++ b/integrations/mem0-agent/eval/fixtures.py @@ -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": "git status --porcelain"}, + {"role": "user", "content": " M src/mem0_agent/api.py\n M tests/test_api.py"}, + ], + "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": "============ 412 passed, 3 skipped in 38.21s " + "============"}, + ], + "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": "ls -la integrations/mem0-agent"}, + {"role": "user", "content": "total 0\ndrwxr-xr-x docs\ndrwxr-xr-x eval\ndrwxr-xr-x " + "hooks\ndrwxr-xr-x src\ndrwxr-xr-x tests"}, + ], + "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_.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//."}, + {"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`, 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 -cd.yml --ref refs/tags/ -f 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()) diff --git a/integrations/mem0-agent/eval/run.py b/integrations/mem0-agent/eval/run.py new file mode 100644 index 000000000..b278d4d1f --- /dev/null +++ b/integrations/mem0-agent/eval/run.py @@ -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()) diff --git a/integrations/mem0-agent/src/mem0_agent/capture.py b/integrations/mem0-agent/src/mem0_agent/capture.py index 8e871e4ce..a557162ff 100644 --- a/integrations/mem0-agent/src/mem0_agent/capture.py +++ b/integrations/mem0-agent/src/mem0_agent/capture.py @@ -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 diff --git a/integrations/mem0-agent/src/mem0_agent/triggers.py b/integrations/mem0-agent/src/mem0_agent/triggers.py index c56e83316..14ed399fa 100644 --- a/integrations/mem0-agent/src/mem0_agent/triggers.py +++ b/integrations/mem0-agent/src/mem0_agent/triggers.py @@ -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"} diff --git a/integrations/mem0-agent/tests/test_mixed_windows.py b/integrations/mem0-agent/tests/test_mixed_windows.py new file mode 100644 index 000000000..54e01739f --- /dev/null +++ b/integrations/mem0-agent/tests/test_mixed_windows.py @@ -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" diff --git a/integrations/mem0-agent/tests/test_triggers.py b/integrations/mem0-agent/tests/test_triggers.py index 13cf1e4eb..144c93579 100644 --- a/integrations/mem0-agent/tests/test_triggers.py +++ b/integrations/mem0-agent/tests/test_triggers.py @@ -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"