Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
41 changes: 41 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,47 @@ python -m benchmarks.beam.run \
--chat-sizes 100K --conversations 0-9
```

### Option C: GoodMemory (local, no Docker)

Runs against [GoodMemory](https://github.com/hjqcan/GoodMemory)'s packaged
HTTP bridge (local-first; requires [Bun](https://bun.sh)). Ingestion is
deterministic (every turn is written verbatim, no LLM extractor needed), and
session timestamps are preserved as UTC observation prefixes. Start the bridge
with the provider-free `recommended` preset for multi-granular BM25, entity,
and reciprocal-rank fusion. Embeddings are optional and add a dense channel;
they are not required for representative provider-free recall.

```bash
npm install -g goodmemory@0.7.5

# Use an ephemeral store for an isolated benchmark run. The official
# goodmemory-client dependency from requirements.txt owns the wire contract.
GOODMEMORY_HTTP_BRIDGE_TOKEN=your-token \
GOODMEMORY_STORAGE_PROVIDER=memory \
goodmemory-http-bridge --recommended # serves http://localhost:8739

# Run a benchmark against it
GOODMEMORY_HTTP_BRIDGE_TOKEN=your-token python -m benchmarks.locomo.run \
--project-name my-goodmemory-test \
--backend goodmemory \
--top-k 10 \
--top-k-cutoffs 10
```

The adapter requests `auto` recall by default and logs the bridge's requested
and resolved routing whenever it falls back; override with
`GOODMEMORY_RECALL_STRATEGY`.
GoodMemory 0.7.5's `recall-context` contract returns at most 12 selected items
and does not expose a caller-controlled item limit. Use cutoff 10 for a
comparable configured cutoff; 20/50/200 cannot retrieve additional items in
this release. The adapter preserves the bridge's ranking but records score
`0.0` because the bridge does not publish a numeric relevance score.
Point at a non-default bridge with `--goodmemory-host` or `GOODMEMORY_HOST`.
Each run should use a fresh `--project-name`: GoodMemory isolates state by
scope (the user id embeds the project name) and the bridge intentionally has
no bulk delete endpoint. In-memory storage is fresh per bridge process, which
suits a benchmark run.

### Option B: Mem0 OSS (Self-Hosted)

Requires Docker and Docker Compose. This starts a local Mem0 server backed by Qdrant.
Expand Down
25 changes: 17 additions & 8 deletions benchmarks/beam/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@

from benchmarks.common.llm_client import LLMClient
from benchmarks.common.mem0_client import Mem0Client, format_search_results
from benchmarks.common.goodmemory_client import GoodMemoryClient
from benchmarks.common.metrics import compute_kendall_tau_b, compute_overall_metrics
from benchmarks.common.schema import (
CutoffResult,
Expand Down Expand Up @@ -955,10 +956,12 @@ def parse_args() -> argparse.Namespace:
help="Comma-separated question types to evaluate (default: all)",
)
parser.add_argument("--rpm", type=int, default=200, help="Requests per minute for LLM")
parser.add_argument("--backend", default="oss", choices=["oss", "cloud"],
help="Mem0 backend: 'oss' for self-hosted server (default), 'cloud' for api.mem0.ai")
parser.add_argument("--backend", default="oss", choices=["oss", "cloud", "goodmemory"],
help="Mem0 backend: 'oss' for self-hosted server (default), 'cloud' for api.mem0.ai, 'goodmemory' for a GoodMemory HTTP bridge")
parser.add_argument("--mem0-host", default=None,
help="Mem0 server URL")
parser.add_argument("--goodmemory-host", default=None,
help="GoodMemory HTTP bridge URL (default: GOODMEMORY_HOST env or http://localhost:8739)")
parser.add_argument("--mem0-api-key", default=None,
help="Mem0 API key (cloud mode only)")
return parser.parse_args()
Expand Down Expand Up @@ -1021,12 +1024,18 @@ async def async_main() -> None:

# Init clients
backend = os.getenv("MEM0_BACKEND", args.backend)
mem0 = Mem0Client(
mode=backend,
host=args.mem0_host,
api_key=args.mem0_api_key if backend == "cloud" else None,
rpm=args.rpm,
)
if backend == "goodmemory":
mem0 = GoodMemoryClient(
host=args.goodmemory_host,
rpm=args.rpm,
)
else:
mem0 = Mem0Client(
mode=backend,
host=args.mem0_host,
api_key=args.mem0_api_key if backend == "cloud" else None,
rpm=args.rpm,
)
answerer = LLMClient(
model=args.answerer_model, provider=args.provider, rpm=args.rpm
)
Expand Down
218 changes: 218 additions & 0 deletions benchmarks/common/goodmemory_client.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,218 @@
"""
Async benchmark adapter for the official GoodMemory Python bridge client.

Setup (no Docker required)::

npm install -g goodmemory@0.7.5
pip install goodmemory-client==0.1.0
GOODMEMORY_HTTP_BRIDGE_TOKEN=replace-me \
goodmemory-http-bridge --recommended

The official client owns the HTTP wire contract, caller/scope identity,
authorization headers, retries, and recall-routing metadata. This module only
adapts that synchronous client to the async ``Mem0Client`` interface used by
the benchmark runners.
"""

from __future__ import annotations

import asyncio
import logging
import os
from datetime import datetime, timezone
from typing import Any

from aiolimiter import AsyncLimiter
from goodmemory_client import (
GoodMemoryClient as BridgeClient,
GoodMemoryClientError,
Scope,
)

logger = logging.getLogger(__name__)

PUBLISHED_RECALL_ITEM_LIMIT = 12


def format_observed_content(
content: str,
observation_date: str | None = None,
timestamp: int | None = None,
) -> str:
"""Preserve benchmark event time in the text GoodMemory indexes."""
if timestamp is not None:
observed_at = datetime.fromtimestamp(timestamp, tz=timezone.utc).isoformat()
observed_at = observed_at.replace("+00:00", "Z")
elif observation_date:
observed_at = observation_date
else:
return content
return f"[Observed at {observed_at}] {content}"


class GoodMemoryClient:
"""Expose the benchmark suite's async memory-client interface."""

def __init__(
self,
host: str | None = None,
token: str | None = None,
max_retries: int = 5,
retry_delay: float = 5.0,
rpm: int = 120,
timeout: float = 300.0,
recall_strategy: str | None = None,
) -> None:
self.host = (host or os.getenv("GOODMEMORY_HOST", "http://localhost:8739")).rstrip("/")
self.token = token or os.getenv("GOODMEMORY_HTTP_BRIDGE_TOKEN") or None
self.max_retries = max_retries
self.retry_delay = retry_delay
self.timeout = timeout
self.recall_strategy = (
recall_strategy or os.getenv("GOODMEMORY_RECALL_STRATEGY") or "auto"
)
self.limiter = AsyncLimiter(rpm, 60)

def _bridge_client(self, user_id: str) -> BridgeClient:
return BridgeClient(
self.host,
scope=Scope(user_id=user_id),
token=self.token,
timeout_seconds=self.timeout,
max_attempts=self.max_retries,
retry_delay_seconds=self.retry_delay,
)

async def close(self) -> None:
"""The stdlib bridge client has no persistent session to close."""

async def __aenter__(self) -> "GoodMemoryClient":
return self

async def __aexit__(self, *exc: Any) -> None:
await self.close()

async def add(
self,
messages: list[dict[str, str]],
user_id: str,
observation_date: str | None = None,
timestamp: int | None = None,
custom_instructions: str | None = None,
metadata: dict | None = None,
) -> dict | None:
"""Import each turn as a verified fact, preserving its source role."""
source_messages = [message for message in messages if message.get("content")]
kept = [
{
"role": "user",
"content": format_observed_content(
(
message["content"]
if (message.get("role") or "user") == "user"
else f"[role={message.get('role') or 'user'}] {message['content']}"
),
observation_date=observation_date,
timestamp=timestamp,
),
}
for message in source_messages
]
if not kept:
return {"skipped": True}

annotations = [
{
"remember": "always",
"confirmed": True,
"verified": True,
"kindHint": "fact",
"messageIndex": index,
"metadataPatch": {
"attributes": {
"sourceRole": source_messages[index].get("role") or "user",
}
},
}
for index in range(len(kept))
]

try:
async with self.limiter:
return await asyncio.to_thread(
self._bridge_client(user_id).remember,
kept,
mode="sync",
extraction_strategy="rules-only",
annotations=annotations,
)
except (GoodMemoryClientError, OSError, ValueError) as error:
logger.error("ADD failed for user=%s: %s", user_id, str(error)[:200])
return None

async def search(
self,
query: str,
user_id: str,
top_k: int = 20,
rerank: bool = False,
score_debug: bool = False,
) -> list[dict]:
"""Recall memories and normalize them to the Mem0 result shape."""
effective_top_k = min(top_k, PUBLISHED_RECALL_ITEM_LIMIT)
if top_k > PUBLISHED_RECALL_ITEM_LIMIT:
logger.warning(
"GoodMemory 0.7.5 recall-context requested top_k=%d but returns "
"at most %d selected items; use cutoff 10 for comparable runs.",
top_k,
PUBLISHED_RECALL_ITEM_LIMIT,
)

try:
async with self.limiter:
result = await asyncio.to_thread(
self._bridge_client(user_id).recall_context,
query,
strategy=self.recall_strategy,
)
except (GoodMemoryClientError, OSError, ValueError) as error:
logger.error("SEARCH failed for user=%s: %s", user_id, str(error)[:200])
return []

routing = result.routing
if (
routing.fallback_reason
or (
routing.requested_strategy not in ("", "auto")
and routing.resolved_strategy != routing.requested_strategy
)
):
logger.warning(
"GoodMemory recall routing: requested=%s resolved=%s fallback=%s",
routing.requested_strategy,
routing.resolved_strategy,
routing.fallback_reason,
)

normalised: list[dict[str, Any]] = []
for item in result.items[:effective_top_k]:
content = item.get("content", "")
if not content:
continue
normalised.append(
{
"memory": content,
"score": 0.0,
"id": item.get("memoryId", ""),
}
)
return normalised

async def delete_user(self, user_id: str) -> bool:
"""Runs are isolated by the project-derived GoodMemory user id."""
logger.info(
"GoodMemory delete_user(%s): no bulk delete endpoint; "
"use a fresh --project-name per run.",
user_id,
)
return True
25 changes: 17 additions & 8 deletions benchmarks/locomo/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@

from benchmarks.common.llm_client import LLMClient
from benchmarks.common.mem0_client import Mem0Client, format_search_results
from benchmarks.common.goodmemory_client import GoodMemoryClient
from benchmarks.common.metrics import compute_overall_metrics
from benchmarks.common.schema import (
CutoffResult,
Expand Down Expand Up @@ -709,10 +710,12 @@ def parse_args() -> argparse.Namespace:
parser.add_argument("--user-profile", action="store_true", help="Fetch user profiles")
parser.add_argument("--max-questions", type=int, default=None, help="Max questions to process (for quick testing)")
parser.add_argument("--rpm", type=int, default=200, help="Requests per minute for LLM")
parser.add_argument("--backend", default="oss", choices=["oss", "cloud"],
help="Mem0 backend: 'oss' for self-hosted server (default), 'cloud' for api.mem0.ai")
parser.add_argument("--backend", default="oss", choices=["oss", "cloud", "goodmemory"],
help="Mem0 backend: 'oss' for self-hosted server (default), 'cloud' for api.mem0.ai, 'goodmemory' for a GoodMemory HTTP bridge")
parser.add_argument("--mem0-host", default=None,
help="Mem0 server URL (default: http://localhost:8888 for oss, https://api.mem0.ai for cloud)")
parser.add_argument("--goodmemory-host", default=None,
help="GoodMemory HTTP bridge URL (default: GOODMEMORY_HOST env or http://localhost:8739)")
parser.add_argument("--mem0-api-key", default=None,
help="Mem0 API key (cloud mode only)")
return parser.parse_args()
Expand Down Expand Up @@ -829,12 +832,18 @@ async def judge_one(qid: str, conv_idx: int, qi: int, qa: dict) -> None:

# Init Mem0 (not used for --evaluate-only)
backend = os.getenv("MEM0_BACKEND", args.backend)
mem0 = Mem0Client(
mode=backend,
host=args.mem0_host,
api_key=args.mem0_api_key if backend == "cloud" else None,
rpm=args.rpm,
)
if backend == "goodmemory":
mem0 = GoodMemoryClient(
host=args.goodmemory_host,
rpm=args.rpm,
)
else:
mem0 = Mem0Client(
mode=backend,
host=args.mem0_host,
api_key=args.mem0_api_key if backend == "cloud" else None,
rpm=args.rpm,
)
shutdown = GracefulShutdown()
checkpoint = Checkpoint(output_dir)

Expand Down
Loading