Store and read back memory on NAMS

We will create a hosted conversation, read its two messages back, wait for extraction, and verify the same messages from a fresh process. We will retain the IDs returned by NAMS and use them for scoped cleanup. This path needs no model API key and creates no reasoning records.

This lesson uses neo4j-agent-memory 0.7.0 from PyPI and the complete programs on this page in a POSIX shell. Create the local files shown below; the checks below verify the results in your environment.

Before you begin

  • Python 3.10 or newer and a POSIX shell.

  • A dedicated NAMS test workspace and its workspace key, using https://memory.neo4jlabs.com/v1.

  • Permission to create tutorial records in that workspace. Keep the same workspace key throughout the lesson.

See Authentication for obtaining the key before beginning. This path uses a workspace key; admin keys and custom deployments belong in the NAMS guide.

Step 1: Install the published SDK

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'
python -c "from importlib.metadata import version; \
import neo4j_agent_memory; \
assert version('neo4j-agent-memory') == '0.7.0'; \
print('SDK 0.7.0 import verified')"

Expected: SDK 0.7.0 import verified. Keep agent-memory-tutorials as your working directory for every command below.

Step 2: Configure the hosted workspace

export MEMORY_API_KEY="replace-with-your-workspace-key"
export MEMORY_ENDPOINT="https://memory.neo4jlabs.com/v1"
unset MEMORY_WORKSPACE_ID
python -c "from neo4j_agent_memory import NamsSettings; \
s=NamsSettings(); \
print(s.backend, s.nams.endpoint)"

Expected: nams https://memory.neo4jlabs.com/v1. This validates configuration only; it does not authenticate or write data.

Before seeding, identify the workspace name or ID in the dashboard and its responsible operator. Supply those non-secret labels in the seed command below. They are saved as operator notes; they do not change request routing or prove workspace identity. Because this lesson unsets MEMORY_WORKSPACE_ID, the ledger’s selector can be null: the workspace key selects the workspace. Never put a key in either note.

Step 3: Write and read back one message

Before building the tutorial’s stateful ledger, confirm the SDK itself works end to end: build settings from the API key exported in Step 2, open a client, create a conversation, write a message, and read it back by the IDs NAMS returns. The program then deletes that conversation, so this check leaves nothing for the later cleanup to find.

Save as nams_first_write.py
"""Smallest hosted round trip: create, write, read back, then delete by returned ID."""

import asyncio

from neo4j_agent_memory import MemoryClient, NamsSettings, NotFoundError

# Entities NAMS extracted from this message. Deleting the conversation does not
# remove them, and an extracted entity can be shared with other workspace data.
EXTRACTED_ENTITIES = "MATCH (m {id: $id})-[:EXTRACTED_FROM]-(e:Entity) RETURN e.id AS id"


async def main() -> None:
    settings = NamsSettings()  # reads MEMORY_API_KEY from the environment
    async with MemoryClient(settings) as client:
        # NAMS assigns the conversation ID; this argument is not sent to the service.
        conversation = await client.short_term.create_conversation("docs-nams-first-write")
        conversation_id = str(conversation.id)
        print(f"Created conversation {conversation_id}")
        content = "Hello from my first hosted write."
        message = await client.short_term.add_message(conversation_id, "user", content)
        history = await client.short_term.get_conversation(conversation_id)
        stored = [item for item in history.messages if item.id == message.id]
        if len(stored) != 1 or stored[0].content != content:
            raise RuntimeError(f"Message {message.id} was not read back; conversation kept")
        print(f"Stored and read back: {stored[0].role.value} said {stored[0].content!r}")

        # Let background extraction finish before deleting its source conversation.
        settled = await client.long_term.wait_for_extraction(
            session_id=conversation_id, timeout=60.0
        )
        if not settled:
            raise TimeoutError(f"Extraction still pending; conversation {conversation_id} kept")
        rows = await client.query.cypher(EXTRACTED_ENTITIES, {"id": str(message.id)})

        await client.short_term.clear_session(conversation_id)
        try:
            await client.short_term.get_conversation(conversation_id)
        except NotFoundError:
            print(f"Deleted conversation {conversation_id}; readback returned 404")
        else:
            raise RuntimeError(f"Conversation {conversation_id} still exists after DELETE")
        if rows:
            ids = sorted(str(row["id"]) for row in rows)
            print(f"Kept entities extracted from the message (not deleted here): {ids}")


if __name__ == "__main__":
    asyncio.run(main())
python nams_first_write.py

Expected output (placeholder IDs):

Created conversation <conversation-id>
Stored and read back: user said 'Hello from my first hosted write.'
Deleted conversation <conversation-id>; readback returned 404

NAMS assigns the conversation ID. The docs-nams-first-write argument is not sent to the service, so it cannot be used to find or reopen the conversation; use the printed ID. The program waits for background extraction before deleting, because deleting a conversation does not delete entities extracted from it. The greeting names no person, place or organization, so extraction is unlikely to produce any. If it does, a Kept entities extracted from the message line lists entity IDs that remain in the workspace. An extracted entity can be shared with other records, so give those IDs to the workspace owner instead of deleting them. An error after the Created conversation line leaves that conversation in the workspace; record its ID for the workspace owner before continuing.

Step 4: Read the complete program

The rest of this lesson wraps that same round trip in scaffolding the later verify and cleanup steps use: a private ledger and scoped-deletion helpers. Create them below, then save nams_quickstart.py. Each block contains the complete file, including imports and client cleanup.

Copy these files into the same lesson folder. The final row is the program we will run; the other files are tutorial support, not SDK APIs. Reuse a helper only when its contents match. Finish pending cleanup/recovery before replacing code that owns an existing ledger.

File Purpose Continuing reader

hosted_tutorial_state.py

Private ledger and local inspection

Reuse unchanged

hosted_tutorial_cleanup.py

Scoped deletion and independent absence checks

Reuse unchanged

nams_quickstart.py

The learner-facing seed, verify and cleanup commands

New lesson file

Open the helper below and save its complete source as hosted_tutorial_state.py in ~/agent-memory-tutorials.

Show hosted_tutorial_state.py
Save as 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_cleanup.py in ~/agent-memory-tutorials.

Show hosted_tutorial_cleanup.py
Save as hosted_tutorial_cleanup.py
"""Scoped cleanup using the hosted API's supported conversation/entity routes."""

from __future__ import annotations

import asyncio
import time
from datetime import datetime, timezone
from urllib.parse import quote

# Start from retained message IDs, then inspect all sources and neighboring IDs.
PROVENANCE_QUERY = """
MATCH (m) WHERE m.id IN $message_ids
MATCH (m)-[:EXTRACTED_FROM]-(e:Entity)
WITH DISTINCT e
OPTIONAL MATCH (e)-[:EXTRACTED_FROM]-(source)
WITH e, collect(DISTINCT source.id) AS sources, count(DISTINCT source) AS source_count
OPTIONAL MATCH (e)--(neighbor)
RETURN e.id AS id, e.createdAt AS created_at, sources, source_count,
       collect(DISTINCT neighbor.id) AS neighbors, count(DISTINCT neighbor) AS neighbor_count
"""
RESIDUAL_QUERY = """
MATCH (n) WHERE n.id IN $ids OR n.conversationId = $conversation
RETURN n.id AS id, labels(n) AS labels
"""


async def request_json(http, method, route, **kwargs):
    response = await http.request(method, route, **kwargs)
    response.raise_for_status()
    return response.json()


async def query(http, cypher, params):
    result = await request_json(http, "POST", "query", json={"cypher": cypher, "params": params})
    rows = result.get("rows")
    if not isinstance(rows, list):
        raise RuntimeError("Query did not return rows; retained tutorial state")
    return rows


async def wait_for_extraction(http, conversation_id, count, *, timeout=60.0, interval=1.0):
    if timeout <= 0 or interval <= 0:
        raise ValueError("Polling timeout and interval must be positive")
    deadline = time.monotonic() + timeout
    route = f"conversations/{quote(conversation_id, safe='')}/extraction-status"
    while True:
        remaining = deadline - time.monotonic()
        if remaining <= 0:
            raise TimeoutError("Extraction did not settle; no cleanup was attempted")
        result = await asyncio.wait_for(request_json(http, "GET", route), remaining)
        summary = result.get("summary")
        if not isinstance(summary, dict) or any(
            type(value) is not int or value < 0 for value in summary.values()
        ):
            raise RuntimeError("Invalid extraction status; no cleanup was attempted")
        if any(summary.get(key, 0) for key in ("failed", "error", "cancelled")):
            raise RuntimeError("Extraction failed; retain state for inspection")
        terminal = {"done", "completed", "skipped"}
        pending = {"pending", "processing", "running", "in_progress", "queued"}
        if set(summary) - terminal - pending:
            raise RuntimeError("Unknown extraction status; retain state for inspection")
        if sum(summary.values()) == count and not any(summary.get(key, 0) for key in pending):
            return
        await asyncio.sleep(min(interval, max(0, deadline - time.monotonic())))


async def delete_and_verify(http, route):
    response = await http.delete(route)
    if response.status_code != 404:
        response.raise_for_status()
    readback = await http.get(route)
    if readback.status_code != 404:
        readback.raise_for_status()
        raise RuntimeError(f"Resource still exists after DELETE: {route}")


def exclusively_owned(row, message_ids, started_at):
    try:
        created = datetime.fromisoformat(row["created_at"].replace("Z", "+00:00"))
        started = datetime.fromisoformat(started_at)
        # collect(node.id) omits anonymous nodes. Compare node counts as well
        # so missing IDs cannot hide another source or adjacent resource.
        for field, count_field in (("sources", "source_count"), ("neighbors", "neighbor_count")):
            ids = row.get(field)
            count = row.get(count_field)
            if (
                not isinstance(ids, list)
                or any(not isinstance(value, str) or not value for value in ids)
                or type(count) is not int
                or count != len(set(ids))
            ):
                return False
        return (
            bool(row.get("sources"))
            and set(row["sources"]) <= message_ids
            and started <= created <= datetime.now(timezone.utc)
        )
    except (KeyError, ValueError, TypeError):
        return False


async def cleanup(http, state, *, timeout=60.0, interval=1.0) -> bool:
    """Retain shared/uncertain records; never use unsupported step/skill deletion."""
    if state.data.get("pending"):
        raise RuntimeError("Unresolved write; reconcile its outcome before cleanup")
    conversation_id = state.data.get("conversation_id")
    resources = state.data["resources"]
    if not conversation_id or conversation_id not in resources.get("conversation", {}):
        raise RuntimeError("No recorded conversation owned by this run")
    # A previous preflight cannot authorize a later deletion. Ownership and
    # neighboring sources may have changed while a failed run was interrupted.
    for entry in resources.get("entity", {}).values():
        if entry.get("status") != "deleted":
            entry["owned"] = False
    state.data["cleanup_complete"] = False
    state.save()
    route = f"conversations/{quote(conversation_id, safe='')}"
    response = await http.get(route)
    owned_messages = set(resources.get("message", {}))
    conversation = resources["conversation"][conversation_id]
    if response.status_code != 404:
        response.raise_for_status()
        metadata = response.json()
        expected_metadata = conversation.get("metadata", {})
        if (
            expected_metadata.get("tutorialRun") != state.data["run_token"]
            or metadata.get("id") != conversation_id
            or metadata.get("metadata") != expected_metadata
        ):
            raise RuntimeError("Conversation ownership mismatch; no cleanup attempted")
        history = await request_json(http, "GET", route + "/messages")
        messages = history.get("messages") if isinstance(history, dict) else history
        if not isinstance(messages, list) or {m.get("id") for m in messages} != owned_messages:
            raise RuntimeError("Unrecorded or missing messages; no cleanup attempted")
        await wait_for_extraction(
            http, conversation_id, len(owned_messages), timeout=timeout, interval=interval
        )
        rows = await query(http, PROVENANCE_QUERY, {"message_ids": sorted(owned_messages)})
        candidates = {
            r["id"] for r in rows if exclusively_owned(r, owned_messages, state.data["started_at"])
        }
        while True:
            exclusive = {
                r["id"]
                for r in rows
                if r["id"] in candidates
                and isinstance(r.get("neighbors"), list)
                and set(r["neighbors"]) <= owned_messages | candidates
            }
            if exclusive == candidates:
                break
            candidates = exclusive
        for row in rows:
            # A graph edge to data outside this run prevents entity deletion too.
            owned = row["id"] in candidates
            state.record("entity", row["id"], status="retained", owned=owned, provenance=row)
        state.data["cleanup_preflight_complete"] = True
        state.save()
    elif not state.data.get("cleanup_preflight_complete"):
        raise RuntimeError(
            "Conversation disappeared before provenance capture; inspect retained IDs"
        )

    for entity_id, entry in resources.get("entity", {}).items():
        entity_route = "entities/" + quote(entity_id, safe="")
        if entry.get("status") == "deleted" or not entry.get("owned"):
            readback = await http.get(entity_route)
            if readback.status_code == 404:
                entry["status"] = "deleted"
                entry["disposition"] = "absence independently verified"
                state.save()
                continue
            readback.raise_for_status()
            entry["status"] = "retained"
            entry["owned"] = False
        if not entry.get("owned"):
            entry["disposition"] = "retained: shared or insufficient ownership evidence"
            state.save()
            continue
        await delete_and_verify(http, entity_route)
        entry["status"] = "deleted"
        state.save()
    await delete_and_verify(http, route)
    conversation["status"] = "deleted"
    state.save()
    ids = [
        conversation_id,
        *owned_messages,
        *(
            key
            for key, value in resources.get("entity", {}).items()
            if value["status"] == "deleted"
        ),
        *(
            key
            for kind, entries in resources.items()
            if kind not in {"conversation", "message", "entity"}
            for key in entries
        ),
    ]
    rows = await query(http, RESIDUAL_QUERY, {"ids": ids, "conversation": conversation_id})
    state.data["remaining_rows"] = rows
    remaining_ids = {row.get("id") for row in rows}
    for entries in resources.values():
        for resource_id, entry in entries.items():
            if resource_id in remaining_ids:
                entry["status"] = "retained"
                entry["disposition"] = "present in final exact-ID readback"
    for message_id, entry in resources.get("message", {}).items():
        if message_id not in remaining_ids:
            entry["status"] = "deleted"
    retained = {
        kind: [key for key, value in entries.items() if value["status"] != "deleted"]
        for kind, entries in resources.items()
    }
    state.data["retained_resources"] = {kind: ids for kind, ids in retained.items() if ids}
    state.save()
    # Steps/skills can remain attached to the deleted conversation. These are
    # recorded dispositions, not a successful assertion that every row vanished.
    unsupported = {
        key
        for kind, entries in resources.items()
        if kind not in {"conversation", "message", "entity"}
        for key in entries
    }
    if any(row.get("id") not in unsupported for row in rows):
        raise RuntimeError("Unexpected run-owned residuals remain; inspect state")
    state.data["cleanup_complete"] = not state.data["retained_resources"]
    state.save()
    print("Verified: owned conversation and proven-owned entities are absent")
    print(f"Retained resources: {state.data['retained_resources']}")
    return not state.data["retained_resources"]
Save as nams_quickstart.py
"""Persist, restart/read back, and clean up a run-owned hosted conversation."""

from __future__ import annotations

import argparse
import asyncio
import json
from pathlib import Path

from hosted_tutorial_cleanup import cleanup
from hosted_tutorial_state import TutorialState, http_client, remember_message, verify_messages

STATE = Path(".tutorial-state/nams.json")


async def exercise(client, state):
    name = "docs-nams-" + state.data["run_token"]
    metadata = {"tutorialRun": state.data["run_token"], "tutorialLesson": "nams"}
    state.begin("create conversation " + name)
    conversation = await client.short_term.create_conversation(name, metadata=metadata)
    conversation_id = str(conversation.id)
    state.data["conversation_id"] = conversation_id
    state.record("conversation", conversation_id, metadata=metadata)
    state.finish()
    print(f"Conversation ID: {conversation_id}")
    token = state.data["run_token"][:8]
    transcript = [
        {"role": "user", "content": f"Maya-{token} is preparing the Robotics-{token} workshop."},
        {
            "role": "assistant",
            "content": f"The Robotics-{token} checklist includes Sensor-{token}.",
        },
    ]
    state.begin("store fixture messages")
    stored = await client.short_term.bulk_add_messages(conversation_id, transcript)
    for message in stored:
        remember_message(state, message)
    state.finish()
    if len(stored) != len(transcript):
        raise RuntimeError("Message write did not return both records; inspect retained state")
    await verify_messages(client, state)
    settled = await client.long_term.wait_for_extraction(
        session_id=conversation_id, timeout=60.0, interval=1.0
    )
    if not settled:
        raise TimeoutError(f"Extraction still pending for conversation {conversation_id}")
    entities = await client.long_term.search_entities(f"Robotics-{token} workshop", limit=5)
    print(f"Extraction settled; workspace search returned {len(entities)} candidate(s)")
    # Search candidates do not establish this run's entity provenance or quality.
    state.data["seed_verified"] = True
    state.save()
    print(f"Retained state: {state.path}; run verify in a new process, then cleanup")
    return conversation_id


async def main(argv=None):
    from neo4j_agent_memory import NamsSettings, connect

    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument(
        "command", choices=["seed", "inspect", "verify", "cleanup"], nargs="?", default="seed"
    )
    parser.add_argument("--state", type=Path, default=STATE)
    parser.add_argument(
        "--workspace-label", help="Operator-recorded workspace name/ID; not routing"
    )
    parser.add_argument("--workspace-owner", help="Operator responsible for resource disposition")
    args = parser.parse_args(argv)
    settings = NamsSettings()
    state = (
        TutorialState.create(
            args.state,
            settings,
            "nams",
            workspace_label=args.workspace_label,
            workspace_owner=args.workspace_owner,
        )
        if args.command == "seed"
        else TutorialState.load(args.state, settings, "nams")
    )
    if args.command == "inspect":
        print(json.dumps(state.inspect(), indent=2))
        return
    if args.command == "cleanup":
        async with http_client(settings) as http:
            complete = await cleanup(http, state)
        if not complete:
            raise SystemExit(2)
        return
    client = await connect(settings)
    try:
        if args.command == "seed":
            await exercise(client, state)
        else:
            await verify_messages(client, state)
    finally:
        await client.close()


if __name__ == "__main__":
    asyncio.run(main())

Before the first service operation, check the assembled files offline:

python -m py_compile \
    hosted_tutorial_state.py \
    hosted_tutorial_cleanup.py \
    nams_quickstart.py
python nams_quickstart.py --help

Expected: compilation prints nothing and exits successfully; --help prints the command choices. Compilation catches syntax errors; the CLI check also loads the local imports. Neither command authenticates, contacts NAMS or creates tutorial state. Fix a missing or truncated file before continuing.

The shared state and cleanup modules sit beside the script. Its private .tutorial-state/nams.json file records the selected endpoint, workspace and credential identity, returned IDs, content digests and operation status. Keep that directory private and out of version control. Reusing the file with a different key or workspace stops before saved IDs are used; reseeding an existing file is refused.

Step 5: Store and read back the fixture

python nams_quickstart.py seed --state .tutorial-state/nams.json \
  --workspace-label 'workspace name or ID from the dashboard' \
  --workspace-owner 'responsible workspace operator'

The program prints these milestones in order:

  1. Conversation ID: identifies the actual server-created conversation.

  2. Verified: exact IDs, roles and content for 2 stored messages confirms a separate read of both fixture messages.

  3. Extraction settled; workspace search returned …​ candidate(s) reports the search after a finite extraction wait. Candidate count and names depend on the service; a workspace search match does not prove that this conversation produced the entity.

  4. Retained state: identifies the file needed for verification and cleanup.

A timeout or exception means that stage is not verified. Keep the saved state rather than rerunning seed and creating another fixture.

Step 6: Verify from a fresh process

The seed command has exited. From the same agent-memory-tutorials directory, run:

python nams_quickstart.py inspect --state .tutorial-state/nams.json
python nams_quickstart.py verify --state .tutorial-state/nams.json

inspect displays the saved resource IDs and pending/disposition information. Expected from verify: Verified: exact IDs, roles and content for 2 stored messages. That command reads the original records without repeating seed, search or model generation. This checks a client-process restart; the hosted service is not restarted.

The introductory path creates no reasoning steps or tool calls. Those records currently have no verified public deletion route, and deleting a conversation must not be described as deleting them. See Reasoning API and Hosted cleanup limits before adding a reasoning exercise with an explicit retention arrangement.

Step 7: Clean up and inspect the disposition

Before switching workspaces or continuing to another lesson, run:

python nams_quickstart.py cleanup --state .tutorial-state/nams.json
printf 'Cleanup exit status: %s\n' "$?"
python nams_quickstart.py inspect --state .tutorial-state/nams.json

The helper waits for terminal extraction of the saved messages, examines only entities linked to them, and checks creation and all source/neighbor IDs before deletion. It uses supported conversation/entity DELETE routes, then verifies 404 responses and a query restricted to this run’s IDs. It does not use a workspace reset or write-Cypher workaround.

Expected: Verified: owned conversation and proven-owned entities are absent, followed by Retained resources:. An empty retained-resource mapping means no recorded resource remains. Shared, merged or ambiguous entities stay listed and must not be called deleted. Exit status 2 reports retained resources even when the supported deletions passed. Inspect that disposition before claiming this run is fully removed. An uncertain write, failed extraction, failed request or unexpected residual stops the check and preserves the state for inspection.

Two illustrative outcomes (placeholder IDs, not a recorded live run):

Verified: owned conversation and proven-owned entities are absent
Retained resources: {}
Cleanup exit status: 0
Verified: owned conversation and proven-owned entities are absent
Retained resources: {'entity': ['<shared-entity-id>']}
Cleanup exit status: 2

The second outcome means the supported deletions passed, but the shared entity was retained. The next action is owner review of inspect output, not another seed. If cleanup instead raises an error, preserve the state and inspect the unfinished operation; neither example output describes that failure.

Closing a client only releases its connections. Keep the private state until every returned ID has a verified deletion or retention disposition. If a write was interrupted before its ID was saved, establish that exact resource with the workspace owner instead of deleting by a similar name or reseeding.

Inspect locally after a credential change

If the original key expires or is rotated, normal lesson commands keep refusing to use its saved IDs with another credential. Use this separate local command:

python hosted_tutorial_state.py inspect-file .tutorial-state/nams.json

It needs no key and makes no service requests. The summary omits credential fingerprints, message digests, fixture metadata and review/provenance payloads. It still contains resource IDs and operator labels, so keep it private. Its local_only: true and service_checked: false mean it reports saved evidence, not current remote state or authorization.

Give the authorized owner the workspace label, endpoint, run token, exact resource IDs and pending operation through your agreed private channel. The owner must establish the actual workspace and reconcile those resources with a valid credential. Do not edit the identity hash, substitute a key in the ledger, or remove state to unlock a write. The helper cannot migrate a run to a new key or supply an administrative deletion route.

Keep the old run and start a separate exercise

See Recover from an interrupted tutorial run to check the recorded state before reseeding.

For a prerequisite failure before the first remote write, the new ledger’s remote_write_started: false, empty resources and absent pending operation establish that this program did not attempt a mutation. An empty ledger alone is not enough. Older empty ledgers without that marker, a lost returned ID, or any uncertain write fail the check and require owner inspection.

After correcting the prerequisite, choose this unused run directory. The first command only reads the old ledger and checks the new path. It preserves the original file and creates nothing. && prevents seed if that check fails:

python hosted_tutorial_state.py check-new-run .tutorial-state/nams.json \
  --next-state .tutorial-state/nams-run-2/state.json &&
python nams_quickstart.py seed --state .tutorial-state/nams-run-2/state.json \
  --workspace-label 'workspace name or ID from the dashboard' \
  --workspace-owner 'responsible workspace operator'

Expected before seed: Recorded new-run check: no_remote_write_started for an unstarted run, or recorded_disposition_complete for fully disposed records. This is a check of saved evidence, not a new service verification. The seed still checks its own prerequisites and refuses an existing target state file. Use the new --state path for every later command; reports and archives are written alongside that new state. Never copy the old IDs into the new ledger.

What we learned

The conversation name and returned ID have different roles. Exact message readback verifies persistence independently of asynchronous extraction and workspace-wide semantic search. A saved run ledger makes the later verification and cleanup refer to the same records.

The Python Bolt and NAMS APIs have different creation and return contracts. Do not switch the settings object and assume this script is portable. See Backend capabilities.

Explore the next hosted capabilities

Context assembly, graph expansion, ontology management, and read-only Cypher are separate tasks:

Continue with the ontology lesson after this conversation path succeeds.