Activate and verify a workspace ontology
|
Available on NAMS: Yes. This procedure targets a NAMS workspace. Since 0.7.0 a Bolt client has the same |
To apply a domain schema to future NAMS writes, clone a template, choose the validation mode and activate its version ID.
Prerequisites
-
Python 3.10+ and a POSIX shell.
-
A dedicated, nonconcurrent sandbox workspace with a key in
MEMORY_API_KEYand permission to manage its ontology. -
A selected template from the current catalog. The program below uses
healthcareonly after checking it is present. -
A known active version to restore if validation fails. For a temporary exercise that restores automatically, use the ontology tutorial.
Prepare the local helpers and install the published client:
Create a local folder and virtual environment for the examples on this page. The commands reuse an existing environment without changing its files:
mkdir -p ~/agent-memory-tutorials
cd ~/agent-memory-tutorials
if [ -e .venv ]; then
printf '%s\n' 'Using the existing virtual environment.'
else
python3 -m venv .venv
fi
source .venv/bin/activate
Expected: ~/agent-memory-tutorials is your working directory and its virtual environment is active. Install the published SDK with the command below.
Each complete code block labelled Save as names a file to create in this folder using your editor. Copy the entire block, including imports and the entry point. Expand each helper disclosure and use Copy code to copy its full source. Keep all files together so their imports resolve.
When continuing from another tutorial or guide, retain the existing environment, configuration, session files, and .tutorial-state/. Reuse unchanged helper files; compare an existing file before replacing it, and finish any pending cleanup or recovery before changing the code that owns its state.
python -m pip install 'neo4j-agent-memory[nams]==0.7.0'
|
In 0.7.0, |
1. Save a complete activation program
The task is four calls: clone a template, update it into a strict revision, activate that version, then verify it took. Save this as activate_ontology.py in agent-memory-tutorials/. It lists the catalog, creates the strict revision, verifies its active ID with read_active_binding and leaves that version active on success, printing both IDs needed for an intentional rollback. The two helper modules it imports, hosted_tutorial_helpers and hosted_tutorial_state, are in the appendix at the end of this page — save them alongside this file before running it.
activate_ontology.pyimport asyncio
from neo4j_agent_memory import NamsSettings, connect
from hosted_tutorial_helpers import read_active_binding
from hosted_tutorial_state import http_client
async def main():
settings = NamsSettings()
client = await connect(settings)
try:
async with http_client(settings) as http:
catalog = await client.ontology.list()
names = [item.name for item in catalog]
print("Available ontologies:", names)
if "healthcare" not in names:
raise RuntimeError("Expected healthcare template is unavailable")
previous = await read_active_binding(http)
print("Previous version:", previous["version_id"])
clone = await client.ontology.clone("healthcare")
print("Workspace ontology:", clone.ontology_id)
try:
version = await client.ontology.update(
clone.ontology_id, clone.document, validation_mode="strict"
)
await client.ontology.activate(version.id)
active = await read_active_binding(http)
if active["version_id"] != version.id or active["validation_mode"] != "strict":
raise RuntimeError("Active version readback did not match")
except Exception:
await client.ontology.activate(previous["version_id"])
restored = await read_active_binding(http)
if (restored["version_id"], restored["validation_mode"]) != (
previous["version_id"], previous["validation_mode"]
):
raise RuntimeError(f"Restore not verified; inspect {clone.ontology_id}")
await client.ontology.delete(clone.ontology_id)
raise
print("Verified active strict version:", active["version_id"])
print("Entity types:", [item["label"] for item in active["document"]["entity_types"]])
finally:
await client.close()
asyncio.run(main())
2. Run the program
python activate_ontology.py
Use the printed entity labels for writes governed by this schema. Edit clone.document before update when customizing a domain; the full document model and validation modes are in ontology reference. update and delete take an ontology ID; activate takes a version ID. Do not interchange them.
3. Verify and retain rollback information
Success prints Verified active strict version with the same version returned by update. A failing activation or readback attempts to restore and verify the prior version before deleting the clone. If restoration fails, retain the clone and printed identifiers for inspection.
Readback verifies the selected active version. For a temporary activation that restores the previous binding and removes its clone, run the activation and restoration lesson. That lesson does not test rejection of an unknown entity label.
Verify extraction separately
Use a bounded wait tied to its returned conversation ID or exact expected names. Ontology activation, direct entity writes and asynchronous message extraction have different success criteria.
On Bolt, extraction completes within the write, so there is nothing to wait for; long_term.wait_for_extraction() returns True immediately there.
Limits and gotchas
-
Activating an ontology changes the workspace’s active ontology; it does not rewrite existing entities.
-
On Bolt,
get_active()reports the version activated in the database, whileclient.ontology_documentreports the ontology the client resolved when it connected, which is the one driving its extraction and validation. They can differ: a file passed throughschema_config.ontology_pathorMemoryClient(ontology=…)outranks the stored version, and a version activated after the client connected applies only from its next connection. -
Offline tests do not certify the deployed service. Python and TypeScript currently expose different transport exception mappings; inspect the REST reference before handling errors.
-
A completed message write is not proof of extraction readiness, and a generic nonempty nearest-neighbor search can match older workspace entities rather than the ones just written.
Appendix: helper files
activate_ontology.py above imports these two modules. hosted_tutorial_state.py provides the shared HTTP client and the on-disk run ledger; hosted_tutorial_helpers.py provides read_active_binding, which reads and validates the authoritative GET /ontologies/active binding directly and fails closed when it is incomplete (see the get_active() note under Prerequisites). Save both files in agent-memory-tutorials/ before running the activation program. For a guided walkthrough of what this code does and why, see the ontology tutorial.
Open the helper below and save its complete source as hosted_tutorial_state.py in ~/agent-memory-tutorials.
Show hosted_tutorial_state.py
hosted_tutorial_state.py"""Private, restartable state for the hosted REST tutorials; no I/O at import."""
from __future__ import annotations
import hashlib
import json
import os
import re
import tempfile
from datetime import datetime, timezone
from pathlib import Path
from urllib.parse import urlsplit, urlunsplit
from uuid import uuid4
def read_state(path):
"""Read a local ledger without credentials, SDK imports, or service access."""
path = Path(path)
if path.is_symlink():
raise RuntimeError("Refusing a symlink for tutorial state")
data = json.loads(path.read_text())
if not isinstance(data, dict) or data.get("schema_version") != 1:
raise RuntimeError("Unsupported or incomplete tutorial state")
if not isinstance(data.get("identity"), dict):
raise RuntimeError("Missing tutorial identity")
if (
not isinstance(data.get("resources"), dict)
or not data.get("run_token")
or "pending" not in data
):
raise RuntimeError("Incomplete tutorial state")
if any(not isinstance(entries, dict) for entries in data["resources"].values()):
raise RuntimeError("Invalid resource ledger")
if any(
not isinstance(entry, dict)
for entries in data["resources"].values()
for entry in entries.values()
):
raise RuntimeError("Invalid resource disposition")
datetime.fromisoformat(data["started_at"])
return data
def new_run_status(data):
"""Classify recorded evidence only; never infer absence from an empty ledger."""
if data.get("pending"):
return "needs_reconciliation"
resources = data["resources"]
entries = [entry for group in resources.values() for entry in group.values()]
if (
data.get("remote_write_started") is False
and not entries
and not any(data.get(key) for key in ("conversation_id", "run_id", "skill_id"))
and not data.get("ontology", {}).get("clone_id")
):
return "no_remote_write_started"
disposed = data.get("cleanup_complete") is True or (
data.get("ontology", {}).get("status") == "complete"
and data["ontology"].get("restored") is True
and data["ontology"].get("deleted") is True
)
if (
disposed
and entries
and all(entry.get("status") in {"deleted", "deleted_with_ontology"} for entry in entries)
):
return "recorded_disposition_complete"
return "owner_disposition_required"
# Restoration target and exercise outcome shown by inspect-file.
ONTOLOGY_SUMMARY_KEYS = (
"previous",
"strict",
"clone_id",
"template",
"status",
"restored",
"deleted",
"exercise_error",
)
def inspect_file(path):
"""Return a redacted local summary; this never authenticates or enables writes."""
data = read_state(path)
endpoint = urlsplit(str(data["identity"].get("endpoint", "")))
# Do not expose credentials/query parameters even from a manually altered file.
endpoint = urlunsplit(
(endpoint.scheme, endpoint.netloc.rsplit("@", 1)[-1], endpoint.path, "", "")
)
summary = {
"local_only": True,
"service_checked": False,
"lesson": data.get("lesson"),
"run_token": data["run_token"],
"started_at": data["started_at"],
"identity": {
"endpoint": endpoint,
"workspace_id": data["identity"].get("workspace_id"),
},
"operator_context": data.get("operator_context", {}),
"pending": data.get("pending"),
"new_run_status": new_run_status(data),
"resources": {
kind: {
resource_id: {
key: entry[key] for key in ("status", "owned", "disposition") if key in entry
}
for resource_id, entry in entries.items()
}
for kind, entries in data["resources"].items()
},
}
ontology = data.get("ontology")
if isinstance(ontology, dict):
# The owner needs the saved prior binding to restore the workspace.
# Schema documents are omitted; the IDs identify them exactly.
binding_keys = ("version_id", "ontology_id", "revision", "validation_mode")
summary["ontology"] = {
key: (
{k: value[k] for k in binding_keys if k in value}
if isinstance(value, dict)
else value
)
for key, value in ontology.items()
if key in ONTOLOGY_SUMMARY_KEYS
}
return summary
def check_new_run(path, next_state):
"""Permit a separate path only after no-write proof or recorded full disposal."""
status = new_run_status(read_state(path))
if status not in {"no_remote_write_started", "recorded_disposition_complete"}:
raise RuntimeError(
"No automatic new run: uncertain, retained, or legacy state needs owner disposition"
)
next_state = Path(next_state)
if next_state.exists() or next_state.is_symlink():
raise RuntimeError("Choose an unused state path; the previous run must remain intact")
if next_state.parent.resolve() == Path(path).parent.resolve():
raise RuntimeError("Use a separate run directory so reports and archives cannot collide")
if next_state.parent.exists() and any(next_state.parent.iterdir()):
raise RuntimeError("Choose a new or empty run directory; preserve its existing reports")
return status
def digest(text: str) -> str:
return hashlib.sha256(text.encode()).hexdigest()
def connection(settings):
"""Use the same resolved NamsSettings for SDK and tutorial HTTP requests."""
config = settings.nams
parts = urlsplit(str(config.endpoint).rstrip("/"))
if parts.scheme not in {"http", "https"} or not parts.netloc:
raise ValueError("A hosted HTTP(S) endpoint is required")
if parts.username or parts.password or parts.query or parts.fragment:
raise ValueError("Endpoint must not contain credentials, a query, or a fragment")
if config.transport_mode == "bridge" or (
config.transport_mode == "auto" and not re.search(r"/v\d+(?:/|$)", parts.path)
):
raise ValueError("These tutorials require the hosted REST API")
endpoint = urlunsplit((parts.scheme.lower(), parts.netloc.lower(), parts.path, "", ""))
key = config.api_key.get_secret_value() if config.api_key else ""
if not key:
raise ValueError("Set MEMORY_API_KEY before starting the tutorial")
# Tutorial identity must not differ between the SDK and direct HTTP calls.
headers = dict(config.headers)
if any(k.lower() in {"authorization", "x-workspace-id"} for k in headers):
raise ValueError("Use api_key/workspace_id settings, not overriding auth/workspace headers")
headers["Authorization"] = f"Bearer {key}"
if config.workspace_id:
headers["X-Workspace-Id"] = config.workspace_id
resolved = {"endpoint": endpoint, "workspace_id": config.workspace_id}
return resolved, headers, config.timeout, key
def credential_digest(key: str, salt: str) -> str:
"""Bind state to the credential that seeded it, without keeping a directly
comparable digest of that credential on disk. A per-run salt and a
deliberately slow KDF mean the stored value answers one question -- "is this
the same key as last time?" -- and is not worth attacking offline. Message
content keeps ``digest``: it is a change detector, not a secret."""
return hashlib.scrypt(
key.encode(), salt=bytes.fromhex(salt), n=16384, r=8, p=1, maxmem=64 * 1024 * 1024, dklen=32
).hex()
def identity(settings, salt=None):
resolved, _headers, _timeout, key = connection(settings)
salt = salt or os.urandom(16).hex()
return {**resolved, "credential_salt": salt, "credential_digest": credential_digest(key, salt)}
def http_client(settings, **kwargs):
import httpx
resolved, headers, timeout, _key = connection(settings)
return httpx.AsyncClient(
base_url=resolved["endpoint"] + "/", headers=headers, timeout=timeout, **kwargs
)
def private_json(path: Path, data, *, exclusive=False):
"""Write without exposing partial JSON or replacing a previous exercise."""
path = Path(path)
path.parent.mkdir(parents=True, exist_ok=True, mode=0o700)
if path.is_symlink():
raise RuntimeError("Refusing a symlink for tutorial state")
content = json.dumps(data, indent=2, allow_nan=False) + "\n"
if exclusive:
try:
fd = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600)
except FileExistsError as exc:
raise RuntimeError("State already exists; inspect or recover the prior run") from exc
with os.fdopen(fd, "w") as output:
output.write(content)
output.flush()
os.fsync(output.fileno())
return
fd, temporary = tempfile.mkstemp(prefix=path.name + ".", dir=path.parent)
try:
with os.fdopen(fd, "w") as output:
output.write(content)
output.flush()
os.fsync(output.fileno())
os.replace(temporary, path)
finally:
if os.path.exists(temporary):
os.unlink(temporary)
class TutorialState:
def __init__(self, path, data):
self.path = Path(path)
self.data = data
@classmethod
def create(cls, path, settings, lesson, *, workspace_label=None, workspace_owner=None):
data = {
"schema_version": 1,
"lesson": lesson,
"identity": identity(settings),
"run_token": uuid4().hex,
"started_at": datetime.now(timezone.utc).isoformat(),
"resources": {},
"pending": None,
# Written durably by begin() before the first remote mutation.
# Missing in an older ledger means unknown, never proof of no writes.
"remote_write_started": False,
"operator_context": {
"workspace_label": workspace_label,
"workspace_owner": workspace_owner,
},
}
private_json(Path(path), data, exclusive=True)
return cls(path, data)
@classmethod
def load(cls, path, settings, lesson=None):
path = Path(path)
data = read_state(path)
stored = data["identity"]
salt = stored.get("credential_salt")
if not isinstance(salt, str) or not re.fullmatch(r"[0-9a-f]{32}", salt):
raise RuntimeError("Tutorial identity is missing its credential salt; do not reseed")
if stored != identity(settings, salt):
raise RuntimeError(
"Endpoint, workspace, or credential changed; verify ownership before recovery"
)
if lesson is not None and data.get("lesson") != lesson:
raise RuntimeError("State belongs to another tutorial")
return cls(path, data)
def save(self):
private_json(self.path, self.data)
def record(self, kind, resource_id, **metadata):
if not isinstance(resource_id, str) or not resource_id.strip():
raise RuntimeError(f"Missing returned {kind} ID; retained pending operation")
resources = self.data["resources"].setdefault(kind, {})
resources.setdefault(resource_id, {"status": "retained"}).update(metadata)
self.save()
def begin(self, operation):
if self.data.get("pending"):
raise RuntimeError("Unresolved write; inspect the retained state before retrying")
self.data["remote_write_started"] = True
self.data["pending"] = operation
self.save()
def finish(self):
self.data["pending"] = None
self.save()
def inspect(self):
data = json.loads(json.dumps(self.data))
data["identity"].pop("credential_salt", None)
data["identity"].pop("credential_digest", None)
return data
def remember_message(state, message):
role = getattr(message.role, "value", message.role)
state.record("message", str(message.id), content_sha256=digest(message.content), role=role)
async def verify_messages(client, state):
"""Read-only exact-ID/content verification; safe to run in a fresh process."""
expected = state.data["resources"].get("message", {})
if not expected or not state.data.get("conversation_id"):
raise RuntimeError("No complete message seed to verify")
history = await client.short_term.get_conversation(state.data["conversation_id"])
actual = {str(m.id): m for m in history.messages}
if set(actual) != set(expected):
raise RuntimeError("Conversation has missing or unrecorded messages; retained state")
for message_id, entry in expected.items():
message = actual[message_id]
role = getattr(message.role, "value", message.role)
if digest(message.content) != entry["content_sha256"] or role != entry["role"]:
raise RuntimeError(f"Message readback differs: {message_id}")
print(f"Verified: exact IDs, roles and content for {len(expected)} stored messages")
def main(argv=None):
"""Local diagnostics only; importing this module still performs no I/O."""
import argparse
parser = argparse.ArgumentParser(description="Inspect hosted lesson state without credentials")
commands = parser.add_subparsers(dest="command", required=True)
inspect = commands.add_parser("inspect-file", help="Print a redacted local summary; no network")
inspect.add_argument("state", type=Path)
retry = commands.add_parser("check-new-run", help="Check a separate path; preserve both files")
retry.add_argument("state", type=Path)
retry.add_argument("--next-state", type=Path, required=True)
args = parser.parse_args(argv)
if args.command == "inspect-file":
print(json.dumps(inspect_file(args.state), indent=2))
else:
status = check_new_run(args.state, args.next_state)
print(f"Recorded new-run check: {status}; original state preserved")
print(f"Use --state {args.next_state} only after rechecking the lesson prerequisites")
print("No service state was checked and no new file was created")
if __name__ == "__main__":
main()
Open the helper below and save its complete source as hosted_tutorial_helpers.py in ~/agent-memory-tutorials.
Show hosted_tutorial_helpers.py
hosted_tutorial_helpers.py"""Bounded control flow shared by the hosted documentation exercises.
These helpers know only the documented run outcomes and SDK ontology methods.
They do not contact a service at import time.
"""
from __future__ import annotations
import asyncio
import json
import time
from collections.abc import Awaitable, Callable
from contextlib import asynccontextmanager
from typing import Any
async def poll_skill_run(
fetch: Callable[[], Awaitable[dict[str, Any]]],
*,
timeout: float = 120.0,
interval: float = 2.0,
clock: Callable[[], float] = time.monotonic,
sleep: Callable[[float], Awaitable[None]] = asyncio.sleep,
) -> dict[str, Any]:
"""Wait through queued/running; return a documented outcome or fail closed."""
if timeout <= 0 or interval <= 0:
raise ValueError("timeout and interval must be positive")
deadline = clock() + timeout
while True:
remaining = deadline - clock()
if remaining <= 0:
raise TimeoutError("Skill run did not settle before the tutorial deadline")
run = await asyncio.wait_for(fetch(), timeout=remaining)
outcome = run.get("outcome")
if outcome in {"Created", "Withheld", "Failed"}:
if outcome == "Created" and not run.get("skillId"):
raise RuntimeError("Created run is missing its returned skillId")
return run
if outcome is not None or run.get("status") not in {"queued", "running"}:
raise RuntimeError(f"Unrecognized skill run state: {run!r}")
await sleep(min(interval, max(0.0, deadline - clock())))
def ontology_binding(active):
"""Capture authoritative IDs, mode and schema; never infer a latest revision."""
values = {
"version_id": active.version_id,
"ontology_id": active.ontology_id,
"revision": active.revision,
"validation_mode": active.validation_mode,
}
if (
not values["version_id"]
or not values["ontology_id"]
or type(values["revision"]) is not int
or values["revision"] < 1
or values["validation_mode"] not in {"permissive", "strict"}
):
raise RuntimeError("No restorable authoritative active binding; no ontology was changed")
values["document"] = active.document.model_dump(mode="json")
return values
async def read_active_binding(http):
"""Read a restorable binding directly from the public REST response.
SDK 0.7.0 get_active() reports the bound version but returns None metadata
when a response has no version record (releases before 0.7.0 inferred it
from the latest revision). This reader requires the actual binding returned
by /ontologies/active, validated against the released public document
model, and fails closed rather than substituting a list lookup.
"""
from neo4j_agent_memory.nams import OntologyDocument
response = await http.get("ontologies/active")
response.raise_for_status()
payload = response.json()
version = payload.get("version") if isinstance(payload, dict) else None
if (
not isinstance(version, dict)
or any(
not isinstance(version.get(key), str) or not version[key].strip()
for key in ("id", "ontology_id")
)
or type(version.get("revision")) is not int
or version["revision"] < 1
or version.get("validation_mode") not in {"permissive", "strict"}
):
raise RuntimeError("No restorable authoritative active binding; no ontology was changed")
def document(value):
if isinstance(value, str):
value = json.loads(value)
return OntologyDocument.model_validate(value).model_dump(mode="json")
try:
active_document = document(payload.get("ontology"))
version_document = version.get("schema_json")
if version_document is not None and document(version_document) != active_document:
raise ValueError("Active and version schemas differ")
except (ValueError, TypeError) as exc:
raise RuntimeError("Invalid or conflicting active ontology schema; retained state") from exc
return {
"version_id": version["id"],
"ontology_id": version["ontology_id"],
"revision": version["revision"],
"validation_mode": version["validation_mode"],
"document": active_document,
}
def _version_binding(version):
from types import SimpleNamespace
if version.document is None:
raise RuntimeError("Returned version has no inspectable schema; retained clone")
return ontology_binding(
SimpleNamespace(
version_id=version.id,
ontology_id=version.ontology_id,
revision=version.revision,
validation_mode=version.validation_mode,
document=version.document,
)
)
async def _require_binding(read_active, expected):
actual = await read_active()
if actual != expected:
raise RuntimeError("Active binding changed or readback differs; retained state and clone")
return actual
async def _clone_absent(ontology, clone_id):
from neo4j_agent_memory.core.exceptions import NotFoundError
try:
await ontology.get(clone_id)
except NotFoundError:
return True
return False
async def verify_ontology_restoration(ontology, state, *, read_active):
"""Fresh, read-only verification of the saved binding and exact clone absence."""
run = state.data.get("ontology", {})
if not run.get("previous"):
raise RuntimeError("State has no captured prior binding")
await _require_binding(read_active, run["previous"])
clone_id = run.get("clone_id")
if not clone_id or not await _clone_absent(ontology, clone_id):
raise RuntimeError("Clone absence is not verified; inspect or recover this run")
if state.data.get("pending"):
raise RuntimeError(
"An operation remains unresolved; run recover before claiming completion"
)
print(f"Verified: original version {run['previous']['version_id']} and mode restored")
print(f"Verified: exercise clone {clone_id} is absent")
async def recover_ontology(ontology, state, *, read_active):
"""Restore only a recognized binding, then delete the proven-owned clone.
An interrupted clone with no returned ID cannot be reconciled automatically.
A concurrent editor's binding is never replaced. There is no server-side
compare-and-swap here, so the exercise requires exclusive ontology editing.
"""
from neo4j_agent_memory.core.exceptions import NotFoundError
run = state.data.get("ontology", {})
previous = run.get("previous")
if not previous:
raise RuntimeError("State has no captured prior binding; no recovery writes were made")
current = await read_active()
if current != previous and current != run.get("strict"):
raise RuntimeError("Active binding changed unexpectedly; retained state and clone")
clone_id = run.get("clone_id")
pending = state.data.get("pending")
if not clone_id:
if pending:
raise RuntimeError(
"Clone outcome is uncertain with no returned ID; inspect retained state"
)
# A prerequisite failed before the first write.
print("Verified: original binding unchanged; no clone was created")
return
owned = state.data["resources"].get("ontology", {}).get(clone_id, {})
if owned.get("created_by_run") is not True:
raise RuntimeError("Clone ownership is not established; no recovery writes were made")
if pending:
if not isinstance(pending, dict) or pending.get("operation") not in {
"clone",
"update",
"activate",
"restore",
"delete",
}:
raise RuntimeError("Unrecognized pending operation; inspect retained state")
# Keep uncertainty recorded until deleting the known clone resolves all
# its revisions. Never replay a clone or an update with an unknown result.
run.setdefault("interrupted_operations", []).append(pending)
state.finish()
if current != previous:
state.begin({"operation": "restore", "version_id": previous["version_id"]})
await ontology.activate(previous["version_id"])
await _require_binding(read_active, previous)
state.finish()
else:
await _require_binding(read_active, previous)
run["restored"] = True
state.save()
print(f"Restored active version: {previous['version_id']} ({previous['validation_mode']})")
# Recheck immediately before deletion. This detects observed concurrent
# changes; it cannot make an unversioned service operation atomic.
await _require_binding(read_active, previous)
if not await _clone_absent(ontology, clone_id):
state.begin({"operation": "delete", "ontology_id": clone_id})
try:
await ontology.delete(clone_id)
except NotFoundError:
pass # A retry still requires the independent GET below.
if not await _clone_absent(ontology, clone_id):
raise RuntimeError(f"Deletion not verified; retained clone {clone_id}")
state.finish()
state.record("ontology", clone_id, status="deleted")
for version_id, entry in state.data["resources"].get("ontology_version", {}).items():
if entry.get("ontology_id") == clone_id:
entry["status"] = "deleted_with_ontology"
run["deleted"] = True
run["status"] = "complete"
state.save()
print(f"Deleted exercise clone: {clone_id} (absence verified)")
@asynccontextmanager
async def temporary_strict_ontology(ontology: Any, template: str, state, *, read_active):
"""Persist the restoration target before mutation and retain failed recovery."""
if state.data.get("ontology") or state.data.get("pending"):
raise RuntimeError("Ontology exercise already started; inspect or recover its state")
previous = await read_active()
catalog = await ontology.list()
if not any(item.name == template and item.is_system for item in catalog):
raise RuntimeError("Required system template is absent; no ontology was changed")
existing_ids = {item.id for item in catalog}
state.data["ontology"] = {"previous": previous, "template": template, "status": "started"}
state.save()
print(f"Previous active version: {previous['version_id']} ({previous['validation_mode']})")
try:
state.begin({"operation": "clone", "template": template})
clone = await ontology.clone(template)
run = state.data["ontology"]
run["clone_id"] = clone.ontology_id
state.record(
"ontology", clone.ontology_id, created_by_run=clone.ontology_id not in existing_ids
)
state.record("ontology_version", clone.id, ontology_id=clone.ontology_id)
state.finish()
if clone.ontology_id in existing_ids:
raise RuntimeError(
"Clone response identifies a preexisting ontology; ownership is unknown"
)
if clone.document is None:
raise RuntimeError("Clone has no inspectable schema")
print(f"Exercise clone: {clone.ontology_id}")
state.begin({"operation": "update", "ontology_id": clone.ontology_id})
strict = await ontology.update(clone.ontology_id, clone.document, validation_mode="strict")
state.record("ontology_version", strict.id, ontology_id=strict.ontology_id)
run["strict"] = _version_binding(strict)
state.finish()
if strict.ontology_id != clone.ontology_id or strict.validation_mode != "strict":
raise RuntimeError("Strict revision response differs from the requested clone and mode")
await _require_binding(read_active, previous)
state.begin({"operation": "activate", "version_id": strict.id})
await ontology.activate(strict.id)
await _require_binding(read_active, run["strict"])
state.finish()
print(f"Activated strict version: {strict.id}")
yield strict
except BaseException as exercise_error:
state.data["ontology"]["exercise_error"] = type(exercise_error).__name__
state.save()
try:
await recover_ontology(ontology, state, read_active=read_active)
except BaseException as recovery_error:
raise recovery_error from exercise_error
raise
else:
await recover_ontology(ontology, state, read_active=read_active)