Skip to content
Open
Show file tree
Hide file tree
Changes from 10 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
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
42 changes: 15 additions & 27 deletions docs/en/docs/develop/http-api.md
Original file line number Diff line number Diff line change
Expand Up @@ -74,8 +74,9 @@ curl --fail \
"$POWERCONTEXT_URL/v1/memory/remember"
```

The response contains an exact citation. Keep that citation when a later request must revise, retire, or read that
specific immutable revision.
The response's `records` contain independent memories with an `artifact` reference, content, state, and `state_version`.
Read history by exact ArtifactRef, edit through the generic Artifact API with its content ETag, and use Atomic Memory
operations to forget or restore. See [Atomic Memory](../workflows/atomic-memory.md) for complete examples.

Search active entries in the same scope:

Expand Down Expand Up @@ -162,8 +163,9 @@ available to Server administrators through `/v1/access/audit/list`. When authent
each audit event keeps the effective `principal` and the trusted `actor` as separate opaque identities.

The Access wire contract has only three Resource Kinds: `server`, `scope`, and `artifact`. An Artifact Resource uses
the logical identity `{family, artifact_id}` and deliberately contains no Revision. Memory can narrow a grant with a
`memory_entry` selector containing only `entry_id`. Unknown Families, `prompt` when no Prompt lifecycle is implemented,
the logical identity `{family, artifact_id}` and deliberately contains no Revision. Each Atomic Memory has its own
Artifact authority without an entry selector. Legacy `memory_entry` selectors identify old entries; offline migration
retargets their valid grants to the corresponding Atomic Memory. Unknown Families, `prompt` when no Prompt lifecycle is implemented,
and mismatched selectors or roles never create a Binding. `/v1/access/me` reports the current mode, Provider
capabilities, and each Artifact Family's enabled state.

Expand All @@ -173,7 +175,7 @@ target Scope. Consequently, one logical sharing grant covers earlier and later s
publication still records the exact copied Revision and its provenance. Host-local projection remains an
operational surface protected by the corresponding Scope and Artifact checks.

Prompt publication returns `422 / artifact_publication_unsupported` without creating a target Artifact. To configure
Atomic Memory does not support cross-Scope publication. Prompt publication returns `422 / artifact_publication_unsupported` without creating a target Artifact. To configure
a Prompt in another Scope, use `POST /v1/scopes/{scope_id}/artifacts` with `family=prompt` and a registered `prompt_key`,
or update it through `PUT /v1/scopes/{scope_id}/artifacts/prompt/{prompt_key}` with `If-Match`. These operations preserve
the fixed Prompt identity and validate its content.
Expand All @@ -200,7 +202,8 @@ use the same policy enforcement point; MCP tool visibility is not permission.
| Source and context | `/v1/sources/content`, `/v1/context/prepare` | Capture evidence and prepare bounded context |
| Work continuity | `/v1/work/*` | Create work contracts, prepare or acknowledge Handoffs, and record outcomes |
| Low-level Handoff | `/v1/handoff/*` | Activate, prepare, finalize, commit, or continue a Handoff |
| Memory | `/v1/memory/*` | Flush, remember, search, list, get, revise, retire, and inspect changes |
| Atomic Memory | `/v1/atomic-memory/*`, generic Artifact routes | Search, list, inspect state, merge, forget, and restore; use generic Artifact routes for content reads and writes |
| Memory compatibility | `/v1/memory/*` | Flush, remember, search, list, and legacy identity reads; legacy collection mutations return an explicit unsupported error |
| Experience and Skill | `/v1/experience/*`, `/v1/skill/*`, `/v1/skills/*` | Propose, review, package, govern, distribute, and read managed Skill revisions |
| Review | `/v1/artifact-candidates/*` | List, inspect, revise, approve, or reject pending Candidates |
| External Skills | `/v1/external-skills/*` | Scan configured targets and resolve or import packages |
Expand All @@ -224,25 +227,11 @@ Errors use one JSON envelope:
}
```

For `/v1/memory/remember` and `/v1/memory/entries/revise`, entry text is limited to 8192 UTF-8 bytes
after Unicode NFC normalization and trimming leading and trailing whitespace. This is a byte limit, not a
character limit. An oversized entry returns HTTP `422` with the existing top-level `invalid_request` code:

```json
{
"error": {
"code": "invalid_request",
"message": "The request is invalid.",
"details": {
"code": "text-too-long",
"message": "memory entry text must not exceed 8192 UTF-8 bytes"
}
}
}
```

Clients can use `error.details.code` to identify the canonical validation failure. Details can still be `null`
for other failures; internal exception text is not returned for unstructured Memory errors.
Atomic Memory text is limited to 8192 UTF-8 bytes; oversized content returns HTTP `422`.
Legacy citation mutations and collection revision preconditions return `legacy_memory_operation_unsupported`.
Do not retry by silently discarding a precondition. Content edits return `428` for a missing `If-Match` or `412`
for a stale ETag; merge and lifecycle state conflicts return `409`. See [Atomic Memory](../workflows/atomic-memory.md)
for error handling and replacement operations.

Common statuses are:

Expand All @@ -257,6 +246,5 @@ Common statuses are:
| `503` | A required Runtime binding or dependency is unavailable |
| `500` | The Server failed without exposing internal details |

Every response includes `X-PowerContext-Request-ID`; record it when diagnosing a failed call. Preserve exact citations
for Memory revision and retirement. Candidate review writes require the current `expected_version`; after a `409`, read
Every response includes `X-PowerContext-Request-ID`; record it when diagnosing a failed call. Keep the content ETag for Memory edits, and the exact ArtifactRef and state_version for forgetting. Candidate review writes require the current `expected_version`; after a `409`, read
the Candidate again before deciding whether to retry.
Loading
Loading