Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
The table of contents is too big for display.
Diff view
Diff view
  •  
  •  
  •  
51 changes: 51 additions & 0 deletions .github/workflows/e2e-harness.yml
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,57 @@ jobs:
with:
version: ${{ env.UV_VERSION }}

- name: Verify OceanBase Atomic Memory and HTTP regressions
if: matrix.database == 'oceanbase'
env:
COMPOSE_PROJECT_NAME: powercontext-e2e-oceanbase
POWERCONTEXT_E2E_OUTPUT: ${{ github.workspace }}/.powercontext-e2e/oceanbase/acceptance
run: |
mkdir -p "$POWERCONTEXT_E2E_OUTPUT"
compose=(docker compose -f e2e/bub/compose.yaml -f e2e/bub/compose.oceanbase.yaml)
"${compose[@]}" up --detach --wait --wait-timeout 660 oceanbase
container_id=$("${compose[@]}" ps -q oceanbase)
oceanbase_ip=$(docker inspect --format '{{range .NetworkSettings.Networks}}{{.IPAddress}}{{end}}' "$container_id")
test -n "$oceanbase_ip"
export POWERCONTEXT_TEST_OCEANBASE_URL="mysql+aoceanbase://root%40test:powercontext-e2e@${oceanbase_ip}:2881/powercontext?charset=utf8mb4"
uv run --locked --extra server python - <<'PY'
import asyncio
import os

from pydantic import SecretStr

from powercontext.builtin.persistence.oceanbase import OceanBaseConfig, OceanBaseProfile

async def version():
config = OceanBaseConfig(url=SecretStr(os.environ["POWERCONTEXT_TEST_OCEANBASE_URL"]))
async with OceanBaseProfile.open(config, tables=()) as profile:
async with profile.database.transaction() as connection:
result = await connection.exec_driver_sql("SELECT VERSION()")
print("OceanBase engine:", result.scalar_one())

asyncio.run(version())
PY
uv run --locked --extra server python -m pytest -vv -ra \
tests/e2e/test_atomic_memory_search_authorization_snapshot.py \
tests/e2e/test_atomic_memory_read_snapshot.py \
tests/e2e/test_runtime_server.py -k oceanbase \
--junitxml="$RUNNER_TEMP/atomic-memory-oceanbase-snapshots.xml"
uv run --locked --extra server python - <<'PY'
import json
import os
import xml.etree.ElementTree as ET
from pathlib import Path

suites = ET.parse(Path(os.environ["RUNNER_TEMP"]) / "atomic-memory-oceanbase-snapshots.xml").iter("testsuite")
counts = {key: 0 for key in ("tests", "failures", "errors", "skipped")}
for suite in suites:
for key in counts:
counts[key] += int(suite.get(key, "0"))
print("OceanBase Atomic Memory and HTTP results:", json.dumps(counts))
if not counts["tests"] or any(counts[key] for key in ("failures", "errors", "skipped")):
raise SystemExit("OceanBase Atomic Memory and HTTP regressions must execute without failures, errors, or skips")
PY

- name: Run deterministic scenarios
id: acceptance_scenarios
env:
Expand Down
181 changes: 131 additions & 50 deletions benchmark/locomo/runner.py

Large diffs are not rendered by default.

112 changes: 86 additions & 26 deletions benchmark/locomo_plus/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,8 +34,17 @@

from benchmark.locomo.dataset import LoCoMoSession
from benchmark.locomo.metrics import retrieval_metrics
from benchmark.locomo.runner import load_settings, normalize_run_id, public_configuration
from powercontext.builtin.artifacts.memory.prompts import memory_extraction_instructions_version
from benchmark.locomo.runner import (
_all_atomic_records,
_atomic_snapshot,
_lineage_source_ids,
load_settings,
normalize_run_id,
public_configuration,
)
from powercontext.builtin.artifacts.atomic_memory.extraction import atomic_memory_extraction_instructions
from powercontext.builtin.artifacts.atomic_memory.models import AtomicMemoryStateValue
from powercontext.builtin.artifacts.atomic_memory.reconciliation import ATOMIC_MEMORY_RECONCILIATION_INSTRUCTIONS
from powercontext.builtin.inference import InvalidInferenceOutputError, character_token_estimator
from powercontext.builtin.persistence.sqlite import SQLiteConfig
from powercontext.builtin.runtime import (
Expand All @@ -61,7 +70,7 @@
)

ARMS = ("memory", "memory-source", "query-only", "oracle-cue", "full-context")
HARNESS_VERSION = "powercontext.locomo-plus.v1"
HARNESS_VERSION = "powercontext.locomo-plus.v2"

# Only application-authored details are safe to copy verbatim. Provider exception messages can contain credentials.
_KNOWN_DETAILS = frozenset({
Expand Down Expand Up @@ -92,7 +101,7 @@ def describe_error(error: BaseException) -> dict[str, Any]:
seen.add(id(current))
item: dict[str, Any] = {"type": type(current).__name__}
if isinstance(current, InvalidInferenceOutputError):
if current.operation in {"generate", "embed", "memory-extract"}:
if current.operation in {"generate", "embed", "atomic-memory-extract", "atomic-memory-reconcile"}:
item["operation"] = current.operation
if current.detail in _KNOWN_DETAILS:
item["detail"] = current.detail
Expand Down Expand Up @@ -286,7 +295,12 @@ def _configuration(settings, judge_model, max_tokens) -> dict[str, Any]:
"embedding": _digest(str(inference.embedding_base_url)),
},
"memory_extraction_profile": "conversation",
"memory_extraction_instructions": memory_extraction_instructions_version(MemoryExtractionProfile.CONVERSATION),
"memory_extraction_instructions_sha256": hashlib.sha256(
atomic_memory_extraction_instructions(MemoryExtractionProfile.CONVERSATION).encode("utf-8")
).hexdigest(),
"memory_reconciliation_instructions_sha256": hashlib.sha256(
ATOMIC_MEMORY_RECONCILIATION_INSTRUCTIONS.encode("utf-8")
).hexdigest(),
}


Expand Down Expand Up @@ -346,16 +360,18 @@ async def _ingest(runtime, case, sessions, scope, output_directory, records, pri
cursor = await memory_app.cursor()
record["processed_session_count"] = cursor.sequence
_write_json(output_directory / "ingestion.json", records)
page = await memory_app.list()
memories = await _all_atomic_records(memory_app, include_inactive=True)
record.update({
"status": "ok",
"session_count": len(sessions),
"memory_count": len(page.entries),
"memories": [entry.model_dump(mode="json") for entry in page.entries],
"schema": "powercontext.benchmark.locomo-plus.ingestion.v2",
"atomic_memory_count": sum(memory.state.state is AtomicMemoryStateValue.ACTIVE for memory in memories),
"atomic_memory_snapshot": [read.model_dump(mode="json") for read in _atomic_snapshot(memories)],
"memories": [memory.model_dump(mode="json", by_alias=True) for memory in memories],
})
record.pop("error_type", None)
record.pop("error", None)
return page # noqa: TRY300
return memories # noqa: TRY300
except Exception as error:
record.update({"status": "error", "error_type": type(error).__name__, "error": describe_error(error)})
if flush_inflight:
Expand Down Expand Up @@ -420,23 +436,37 @@ def _finish_usage(observation, stage, usage, model, prices):

async def _recall_usage(runtime, scope):
statistics = await runtime.statistics.for_scope(scope).overview()
return _sum_usage([
row.embedding.model_dump() for row in statistics.usage.by_purpose if row.purpose.value == "memory_recall"
])
rows = [row for row in statistics.usage.by_purpose if row.purpose.value == "memory_recall"]
return {
"embedding": _sum_usage([{**row.embedding.model_dump(), "output_tokens": 0} for row in rows]),
"generation": _sum_usage([row.generation.model_dump() for row in rows]),
}


async def _retrieve(runtime, case, page, scope, sessions, top_k, source_expansion):
async def _retrieve(runtime, case, scope, sessions, top_k, source_expansion):
result = await runtime.memory.for_scope(scope).search(
SearchMemoryRequest(query=case.question, limit=top_k, mode="hybrid")
)
sources = {(record.entry.entry_id, record.entry.entry_version_id): record.entry.sources for record in page.entries}
records = runtime.records.for_scope(scope)
cache: dict[tuple[str, str, int], tuple[str, ...]] = {}
rendered: list[str] = []
hits: list[dict[str, Any]] = []
session_map = {session.session_id: session for session in sessions}
selected_ids: list[str] = []
for hit in result.hits:
ids = tuple(ref.source_id for ref in sources.get((hit.entry_id, hit.entry_version_id), ()))
hits.append({**hit.model_dump(mode="json"), "source_ids": list(ids)})
for rank, wrapper in enumerate(result.hits, 1):
hit = wrapper.hit
ids = await _lineage_source_ids(records, hit.artifact_ref, cache)
hits.append({
"rank": rank,
"artifact_ref": hit.artifact_ref.model_dump(mode="json"),
"state_version": hit.state_version,
"kind": hit.kind,
"text": hit.text,
"score": hit.score,
"distance": hit.distance,
"matched_by": list(wrapper.matched_by),
"source_ids": list(ids),
})
rendered.append(f"Memory: {hit.text}\nSources: {', '.join(ids)}")
for source_id in ids:
local_id = source_id.rsplit(":", maxsplit=1)[-1]
Expand All @@ -449,7 +479,29 @@ async def _retrieve(runtime, case, page, scope, sessions, top_k, source_expansio
metrics = retrieval_metrics(
evidence_sessions=evidence_sessions, hit_source_ids=tuple(tuple(hit["source_ids"]) for hit in hits)
)
return "\n\n".join(rendered), hits, selected_ids, metrics
return (
"\n\n".join(rendered),
hits,
selected_ids,
{
**metrics,
"mode": result.mode,
"score_kind": "rrf-ranking-score",
"embedding_calls": result.embedding_calls,
"generation_calls": result.generation_calls,
"rerank": None
if result.rerank is None
else {
"policy_id": result.rerank.policy_id,
"candidate_count": len(result.rerank.candidate_hits),
"selected_ranks": list(result.rerank.selected_ranks),
"discarded_rank_count": result.rerank.discarded_rank_count,
"used_fallback": result.rerank.used_fallback,
"latency_ms": result.rerank.latency_ms,
"usage": result.rerank.usage.model_dump(mode="json"),
},
},
)


async def _evaluate(
Expand Down Expand Up @@ -477,6 +529,7 @@ async def _evaluate(
if previous
else {
"case_id": case.case_id,
"schema": "powercontext.benchmark.locomo-plus.observation.v2",
"category": case.category,
"constraint_type": case.relation_type,
"question": case.question,
Expand Down Expand Up @@ -505,22 +558,29 @@ async def _evaluate(
)
scope = registered.scope_id
observation["scope_id"] = scope
page = await _ingest(runtime, case, sessions, scope, output_directory, ingestion, prices, settings)
await _ingest(runtime, case, sessions, scope, output_directory, ingestion, prices, settings)
phase = "retrieval"
queried = perf_counter()
before = await _recall_usage(runtime, scope)
_start_usage(observation, "retrieval", settings.inference.embedding_model, prices)
context, hits, selected_ids, retrieval = await _retrieve(
runtime, case, page, scope, sessions, top_k, arm == "memory-source"
runtime, case, scope, sessions, top_k, arm == "memory-source"
)
observation["latency_ms"]["query"] = (perf_counter() - queried) * 1_000
after = await _recall_usage(runtime, scope)
usage = {
key: None if after[key] is None or before[key] is None else after[key] - before[key]
for key in ("requests", "input_tokens")
}
usage["output_tokens"] = 0
_finish_usage(observation, "retrieval", usage, settings.inference.embedding_model, prices)
for channel, stage, model in (
("embedding", "retrieval", settings.inference.embedding_model),
("generation", "retrieval_generation", settings.inference.generation_model),
):
usage = {
key: None
if after[channel][key] is None or before[channel][key] is None
else after[channel][key] - before[channel][key]
for key in ("requests", "input_tokens", "output_tokens")
}
if stage == "retrieval_generation":
_start_usage(observation, stage, model, prices)
_finish_usage(observation, stage, usage, model, prices)
elif arm == "query-only":
context = ""
else:
Expand Down
22 changes: 19 additions & 3 deletions benchmark/memory_capacity/__main__.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@
# See the License for the specific language governing permissions and
# limitations under the License.

"""Legacy Memory manifest capacity benchmark; this does not measure Atomic Memory."""

from __future__ import annotations

import argparse
Expand All @@ -27,8 +29,9 @@
from pydantic import SecretStr
from sqlalchemy import event, text

from powercontext.builtin.artifacts.memory import MemoryEntryInput
from powercontext.builtin.artifacts.memory import MemoryCompactionPolicy, MemoryEntryInput, MemoryService
from powercontext.builtin.artifacts.memory.canonical import memory_content_bytes
from powercontext.builtin.persistence.memory import RelationalMemoryBackend
from powercontext.builtin.persistence.oceanbase import OceanBaseConfig
from powercontext.builtin.persistence.sqlite import SQLiteConfig
from powercontext.builtin.runtime import BuiltinConfig, open_builtin_contexts
Expand All @@ -52,6 +55,7 @@ async def run(args): # noqa: C901
config = BuiltinConfig(database=database, runtime=RuntimeConfig(memory_compaction_enabled=True))
output = {
"backend": args.backend,
"memory_model": "legacy-memory-manifest",
"status": "running",
"python": platform.python_version(),
"platform": platform.platform(),
Expand All @@ -64,7 +68,17 @@ async def run(args): # noqa: C901
}
async with open_builtin_contexts(config) as contexts:
scope_id = "capacity-benchmark-" + uuid4().hex
service = (await contexts.get(scope_id)).artifacts.memory
service = MemoryService(
backend=RelationalMemoryBackend(
database=contexts.database,
scope_id=scope_id,
artifacts=contexts.repositories.artifacts,
index=contexts.index,
),
compaction=MemoryCompactionPolicy(
enabled=True, min_tombstone_revisions=config.runtime.memory_compaction_min_tombstone_revisions
),
)
projection_rows = 0

def record(_connection, cursor, statement, _parameters, _context, _executemany):
Expand Down Expand Up @@ -192,7 +206,9 @@ async def snapshot(memory):


def main():
parser = argparse.ArgumentParser(description="Measure the Memory capacity envelope without inference.")
parser = argparse.ArgumentParser(
description="Measure legacy Memory manifest capacity without inference; excludes Atomic Memory."
)
parser.add_argument("--backend", choices=("sqlite", "oceanbase"), default="sqlite")
parser.add_argument("--counts", nargs="+", type=int, default=[200, 1000, 5000])
parser.add_argument("--final-window", type=int, default=100)
Expand Down
6 changes: 5 additions & 1 deletion docs/en/docs/develop/api-quickstart.md
Original file line number Diff line number Diff line change
Expand Up @@ -111,7 +111,7 @@ source_exchange = post(
source_ref = source_exchange["source"]

# Explicit long-term writes require application or user authorization.
post(
saved = post(
"/v1/memory/remember",
{
"scope_id": SCOPE_ID,
Expand All @@ -121,6 +121,10 @@ post(
},
)

# Preserve the real independent Artifact identities from the successful write.
memories = saved["records"]
print(json.dumps(memories, ensure_ascii=False, indent=2))

# Prepare bounded historical context for one model request.
question = "How should the assistant handle a refund request for an expired order?"
prepared = post(
Expand Down
Loading
Loading