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.
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 |
|---|---|---|
|
Private ledger and local inspection |
Reuse unchanged |
|
Scoped deletion and independent absence checks |
Reuse unchanged |
|
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
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
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"]
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:
-
Conversation ID:identifies the actual server-created conversation. -
Verified: exact IDs, roles and content for 2 stored messagesconfirms a separate read of both fixture messages. -
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. -
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:
-
Hosted context, observations and reflections: NAMS conversation operations.
-
Graph expansion:
expand_graph. -
Read-only Cypher:
client.query.cypherin the MemoryClient properties. -
Endpoint paths for all of these: REST API.
-
Ontology management: Activate and verify a workspace ontology.
client.ontologyalso exists on a Bolt client, where the ontology is stored in your own Neo4j database and drives local extraction; see ontology-driven extraction.
Continue with the ontology lesson after this conversation path succeeds.