hermes(hux): close Wave A review findings in events, memory, privacy and organization
F3 memory edits go through the same privacy shaping as proposals; F5 seq is derived from the ledger tail so a crash between append and checkpoint never duplicates; F7 transitions re-read under the lock and always write with the loaded revision; F9 secret scrub on titles, passages, claims and notebooks and forget blanks the conversation document; F13 no ghost conversations from notices, idempotency under the lock, artifact titles searchable, normalised paths. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01RNPhwu2bsaRNg3DETSAZoM
This commit is contained in:
parent
6964a9d8c8
commit
681b040885
@ -12,6 +12,7 @@ from __future__ import annotations
|
||||
import json
|
||||
import threading
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from collections.abc import Iterator
|
||||
|
||||
@ -47,14 +48,45 @@ def is_private(store: TenantStore, conversation_id: str) -> bool:
|
||||
return False
|
||||
|
||||
|
||||
def _ledger_path(store: TenantStore, conversation_id: str) -> Path:
|
||||
return store.root / FAMILY / f"{conversation_id}.jsonl"
|
||||
|
||||
|
||||
def conversation_known(store: TenantStore, conversation_id: str) -> bool:
|
||||
"""Ownership check (SO-18): the conversation exists somewhere in this subject's tree."""
|
||||
"""Ownership check (SO-18): the conversation document, its seq checkpoint or its event ledger exists in this subject's tree."""
|
||||
return (
|
||||
store.exists("conversations", conversation_id)
|
||||
or store.exists(SEQ_FAMILY, _seq_id(conversation_id))
|
||||
or _ledger_path(store, conversation_id).is_file()
|
||||
)
|
||||
|
||||
|
||||
def last_seq(store: TenantStore, conversation_id: str) -> int:
|
||||
"""The seq of the newest complete ledger line, read from the file tail (0 when the ledger is empty or absent).
|
||||
|
||||
A crash between the ledger append and the checkpoint write leaves the
|
||||
checkpoint behind the ledger; ``emit`` takes the larger of the two so a seq
|
||||
is never handed out twice (F5). Only the last ``LINE_CAP_BYTES`` window is
|
||||
read, and a torn trailing line is skipped.
|
||||
"""
|
||||
path = _ledger_path(store, conversation_id)
|
||||
if not path.is_file():
|
||||
return 0
|
||||
with open(path, "rb") as handle:
|
||||
handle.seek(0, 2)
|
||||
size = handle.tell()
|
||||
handle.seek(max(0, size - redaction.LINE_CAP_BYTES - 2))
|
||||
tail = handle.read()
|
||||
for line in reversed(tail.split(b"\n")):
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
return int(json.loads(line).get("seq", 0))
|
||||
except (json.JSONDecodeError, AttributeError, ValueError):
|
||||
continue
|
||||
return 0
|
||||
|
||||
|
||||
def _actor(identity: Identity, kind: str) -> dict[str, str]:
|
||||
if kind in USER_KINDS and identity.trust != "worker":
|
||||
return {"type": "user", "id": identity.subject}
|
||||
@ -127,7 +159,7 @@ def emit_with_status(store: TenantStore, identity: Identity, conversation_id: st
|
||||
return existing, True
|
||||
seq_id = _seq_id(conversation_id)
|
||||
checkpoint = store.get(SEQ_FAMILY, seq_id) if store.exists(SEQ_FAMILY, seq_id) else {"id": seq_id, "next_seq": 1, "last_event_id": ""}
|
||||
record["seq"] = int(checkpoint["next_seq"])
|
||||
record["seq"] = max(int(checkpoint["next_seq"]), last_seq(store, conversation_id) + 1)
|
||||
problems = contracts.validate_record(record, SCHEMAS)
|
||||
if problems:
|
||||
raise Invalid("event failed contract validation", problems)
|
||||
|
||||
@ -76,20 +76,37 @@ def _persist(store: TenantStore, record: dict[str, Any], expected: int | None) -
|
||||
|
||||
def _transition(store: TenantStore, identity: Identity, record: dict[str, Any], target: str, action: str, actor: dict[str, str],
|
||||
note: str = "", expected: int | None = None, **changes: Any) -> dict[str, Any]:
|
||||
if not rules.transition_allowed(rules.MEMORY_TRANSITIONS, record["status"], target):
|
||||
raise Conflict(f"{record['status']} -> {target} is not a legal memory transition")
|
||||
stamp = now_iso()
|
||||
updated = {**record, **changes, "status": target, "updated_at": stamp, "audit": [*record["audit"], {"at": stamp, "action": action, "actor": actor, **({"note": note[:200]} if note else {})}]}
|
||||
if target in {"forgotten", "rejected", "expired"}:
|
||||
updated["content"], updated["retrievable"] = "", False
|
||||
if target == "active":
|
||||
updated["retrievable"] = changes.get("retrievable", True)
|
||||
stored = _persist(store, updated, expected)
|
||||
if target == "forgotten":
|
||||
_tombstone(store, stored["id"], action)
|
||||
"""Move an entry to ``target`` from its *current* stored state (F7).
|
||||
|
||||
The caller's ``record`` may be stale: the entry is re-read under the family
|
||||
lock, ``expected`` (If-Match) is checked against that copy, and the write is
|
||||
always conditional on the loaded revision so a stale unconditional write can
|
||||
never resurrect a forgotten entry (SO-22, SO-44).
|
||||
"""
|
||||
with store.lock(FAMILY):
|
||||
record = _current(store, record, expected)
|
||||
if not rules.transition_allowed(rules.MEMORY_TRANSITIONS, record["status"], target):
|
||||
raise Conflict(f"{record['status']} -> {target} is not a legal memory transition")
|
||||
stamp = now_iso()
|
||||
updated = {**record, **changes, "status": target, "updated_at": stamp, "audit": [*record["audit"], {"at": stamp, "action": action, "actor": actor, **({"note": note[:200]} if note else {})}]}
|
||||
if target in {"forgotten", "rejected", "expired"}:
|
||||
updated["content"], updated["retrievable"] = "", False
|
||||
if target == "active":
|
||||
updated["retrievable"] = changes.get("retrievable", True)
|
||||
stored = _persist(store, updated, record["revision"])
|
||||
if target == "forgotten":
|
||||
_tombstone(store, stored["id"], action)
|
||||
return stored
|
||||
|
||||
|
||||
def _current(store: TenantStore, record: dict[str, Any], expected: int | None) -> dict[str, Any]:
|
||||
"""Fresh copy of ``record`` from the store; Conflict when If-Match no longer matches it."""
|
||||
current = store.get(FAMILY, record["id"])
|
||||
if expected is not None and expected != current["revision"]:
|
||||
raise Conflict(f"revision {expected} does not match current revision {current['revision']}", [str(current["revision"])])
|
||||
return current
|
||||
|
||||
|
||||
def load(store: TenantStore, memory_id: str, now: datetime | None = None) -> dict[str, Any]:
|
||||
"""Newest snapshot with lazy expiry (SO-27): an entry past its expiry is written as expired before it is served."""
|
||||
record = store.get(FAMILY, memory_id)
|
||||
@ -175,17 +192,12 @@ def _actor(request: Request) -> dict[str, str]:
|
||||
return {"type": "assistant", "id": "hermes"}
|
||||
|
||||
|
||||
def _shape(request: Request, body: dict[str, Any]) -> tuple[dict[str, Any], str]:
|
||||
"""Build the entry from a proposal body and decide its approval mode from policy; returns (entry, gate_reason)."""
|
||||
def _classify(store: TenantStore, body: dict[str, Any], content: str, conversation_id: str | None) -> dict[str, Any]:
|
||||
"""Scrub ``content`` and decide topic, sensitivity floor and the privacy gate; shared by create and edit (F3)."""
|
||||
from hux import privacy
|
||||
|
||||
if "id" in body or "revision" in body or "status" in body:
|
||||
raise Invalid("id, revision and status are server-assigned")
|
||||
content = body.get("content")
|
||||
if not isinstance(content, str) or not content.strip() or len(content) > 2000:
|
||||
raise Invalid("content must be a non-empty string of at most 2000 chars")
|
||||
hits: list[str] = []
|
||||
content = redaction.scrub_value(content.strip(), hits)
|
||||
content = redaction.scrub_value(content.strip(), hits)[:2000]
|
||||
detected = privacy.detect_topics(content) + (["credentials"] if hits else [])
|
||||
topic = body.get("topic") if body.get("topic") in rules.PRIVACY_TOPICS or body.get("topic") == "general" else None
|
||||
topic = topic if topic and topic != "general" else (detected[0] if detected else "general")
|
||||
@ -195,6 +207,35 @@ def _shape(request: Request, body: dict[str, Any]) -> tuple[dict[str, Any], str]
|
||||
floor = privacy.topic_sensitivity(detected + ([topic] if topic != "general" else []))
|
||||
if privacy.SENSITIVITY_RANK[floor] > privacy.SENSITIVITY_RANK[sensitivity]:
|
||||
sensitivity = floor
|
||||
allowed, why = privacy.memory_write_allowed(store, conversation_id, {"sensitivity": sensitivity, "topic": topic})
|
||||
if why == "private_mode":
|
||||
raise Forbidden("memory writes are refused in private mode")
|
||||
return {"content": content, "topic": topic, "sensitivity": sensitivity, "allowed": allowed, "why": why}
|
||||
|
||||
|
||||
def _decline(entry: dict[str, Any], why: str, actor: dict[str, str], stamp: str) -> dict[str, Any]:
|
||||
"""Content-free ``no_store`` decision record.
|
||||
|
||||
The topic moves into the note because rules.memory_policy_violations reads a
|
||||
deny topic on any non-rejected status as a write, and a sensitive entry must
|
||||
say approval_mode=ask.
|
||||
"""
|
||||
topic = entry.pop("topic")
|
||||
entry.update(status="no_store", approval_mode="ask" if entry["sensitivity"] == "sensitive" else "no_store", content="", retrievable=False,
|
||||
audit=[{"at": stamp, "action": "no_store", "actor": actor, "note": f"{why}; topic={topic}"[:200]}])
|
||||
return entry
|
||||
|
||||
|
||||
def _shape(request: Request, body: dict[str, Any]) -> tuple[dict[str, Any], str]:
|
||||
"""Build the entry from a proposal body and decide its approval mode from policy; returns (entry, gate_reason)."""
|
||||
if "id" in body or "revision" in body or "status" in body:
|
||||
raise Invalid("id, revision and status are server-assigned")
|
||||
content = body.get("content")
|
||||
if not isinstance(content, str) or not content.strip() or len(content) > 2000:
|
||||
raise Invalid("content must be a non-empty string of at most 2000 chars")
|
||||
conversation_id = body.get("conversation_id") if isinstance(body.get("conversation_id"), str) else None
|
||||
verdict = _classify(request.store, body, content, conversation_id)
|
||||
topic, sensitivity = verdict["topic"], verdict["sensitivity"]
|
||||
scope = body.get("scope") if isinstance(body.get("scope"), dict) else {"level": "global"}
|
||||
source = body.get("source") if isinstance(body.get("source"), dict) else {"kind": "message", "id": "unspecified"}
|
||||
ttl = body.get("ttl") if isinstance(body.get("ttl"), dict) else None
|
||||
@ -203,7 +244,6 @@ def _shape(request: Request, body: dict[str, Any]) -> tuple[dict[str, Any], str]
|
||||
ttl = {"policy": "decay", "decay_days": days}
|
||||
actor = _actor(request)
|
||||
stamp = now_iso()
|
||||
conversation_id = body.get("conversation_id") if isinstance(body.get("conversation_id"), str) else None
|
||||
provenance = {"surface": request.identity.surface, "actor": actor, "recorded_at": stamp}
|
||||
if conversation_id:
|
||||
provenance["conversation_id"] = conversation_id
|
||||
@ -211,7 +251,7 @@ def _shape(request: Request, body: dict[str, Any]) -> tuple[dict[str, Any], str]
|
||||
provenance["run_id"] = body["run_id"][:120]
|
||||
entry: dict[str, Any] = {
|
||||
"schema": "hux.memory.v1", "id": new_id("mem"), "owner": request.identity.subject, "scope": scope,
|
||||
"kind": body.get("kind") if body.get("kind") in KINDS else "fact", "content": content, "status": "proposed",
|
||||
"kind": body.get("kind") if body.get("kind") in KINDS else "fact", "content": verdict["content"], "status": "proposed",
|
||||
"approval_mode": "ask", "sensitivity": sensitivity, "topic": topic, "ttl": ttl, "source": source,
|
||||
"provenance": provenance, "created_at": stamp, "updated_at": stamp,
|
||||
"audit": [{"at": stamp, "action": "proposed", "actor": actor}], "reason": str(body.get("reason") or "Proposed to be remembered.")[:280],
|
||||
@ -219,19 +259,12 @@ def _shape(request: Request, body: dict[str, Any]) -> tuple[dict[str, Any], str]
|
||||
}
|
||||
if isinstance(body.get("supersedes"), str):
|
||||
entry["supersedes"] = body["supersedes"]
|
||||
allowed, why = privacy.memory_write_allowed(request.store, conversation_id, entry)
|
||||
if why == "private_mode":
|
||||
raise Forbidden("memory writes are refused in private mode")
|
||||
references_forgotten = entry.get("supersedes") in tombstoned(request.store) or source.get("id") in tombstoned(request.store)
|
||||
wants = body.get("approval_mode", "ask" if actor["type"] == "assistant" else "automatic")
|
||||
if not allowed or wants == "no_store":
|
||||
# Content-free decision record. The topic moves into the note because
|
||||
# rules.memory_policy_violations reads a deny topic on any non-rejected
|
||||
# status as a write, and a sensitive entry must say approval_mode=ask.
|
||||
why = why if not allowed else "declined"
|
||||
entry.pop("topic")
|
||||
entry.update(status="no_store", approval_mode="ask" if sensitivity == "sensitive" else "no_store", content="", retrievable=False,
|
||||
audit=[{"at": stamp, "action": "no_store", "actor": actor, "note": f"{why}; topic={topic}"[:200]}])
|
||||
why = verdict["why"]
|
||||
if not verdict["allowed"] or wants == "no_store":
|
||||
why = why if not verdict["allowed"] else "declined"
|
||||
_decline(entry, why, actor, stamp)
|
||||
elif sensitivity == "sensitive" or wants == "ask" or references_forgotten or actor["type"] != "user":
|
||||
entry["approval_mode"] = "ask"
|
||||
else:
|
||||
@ -240,26 +273,37 @@ def _shape(request: Request, body: dict[str, Any]) -> tuple[dict[str, Any], str]
|
||||
return entry, why
|
||||
|
||||
|
||||
def _replay(request: Request, key: str) -> Response | None:
|
||||
for row in request.store.read(IDEM_FAMILY, "keys"):
|
||||
if row["idempotency_key"] == key:
|
||||
request.audit("memory.propose", row["memory_id"], "allow", "replayed")
|
||||
return Response(200, load(request.store, row["memory_id"]), {"HUX-Replayed": "true"})
|
||||
return None
|
||||
|
||||
|
||||
def post_memory(request: Request) -> Response:
|
||||
"""``POST /hux/v1/memory``: propose an entry; policy decides active, proposed (ask) or no_store (202)."""
|
||||
"""``POST /hux/v1/memory``: propose an entry; policy decides active, proposed (ask) or no_store (202).
|
||||
|
||||
The Idempotency-Key lookup, the write and the key mapping all happen under
|
||||
the family lock so concurrent retries with one key yield one record (F13b).
|
||||
"""
|
||||
body = request.body if isinstance(request.body, dict) else None
|
||||
if body is None:
|
||||
raise Invalid("body must be an object")
|
||||
key = request.idempotency_key()
|
||||
if key:
|
||||
for row in request.store.read(IDEM_FAMILY, "keys"):
|
||||
if row["idempotency_key"] == key:
|
||||
request.audit("memory.propose", row["memory_id"], "allow", "replayed")
|
||||
return Response(200, load(request.store, row["memory_id"]), {"HUX-Replayed": "true"})
|
||||
entry, why = _shape(request, body)
|
||||
try:
|
||||
stored = _persist(request.store, entry, None)
|
||||
except Invalid:
|
||||
request.store.append(FAMILY, TOMBSTONES, {"memory_id": entry["id"], "at": now_iso(), "reason": "policy_violation", "purged": True})
|
||||
_emit(request.store, request.identity, entry, "memory.suppressed", "Memory proposal suppressed by policy")
|
||||
raise
|
||||
if key:
|
||||
request.store.append(IDEM_FAMILY, "keys", {"idempotency_key": key, "memory_id": stored["id"], "at": stored["created_at"]})
|
||||
with request.store.lock(FAMILY):
|
||||
replayed = _replay(request, key) if key else None
|
||||
if replayed is not None:
|
||||
return replayed
|
||||
entry, why = _shape(request, body)
|
||||
try:
|
||||
stored = _persist(request.store, entry, None)
|
||||
except Invalid:
|
||||
request.store.append(FAMILY, TOMBSTONES, {"memory_id": entry["id"], "at": now_iso(), "reason": "policy_violation", "purged": True})
|
||||
_emit(request.store, request.identity, entry, "memory.suppressed", "Memory proposal suppressed by policy")
|
||||
raise
|
||||
if key:
|
||||
request.store.append(IDEM_FAMILY, "keys", {"idempotency_key": key, "memory_id": stored["id"], "at": stored["created_at"]})
|
||||
if stored["status"] == "no_store":
|
||||
_tombstone(request.store, stored["id"], why)
|
||||
_emit(request.store, request.identity, stored, "memory.suppressed", f"Not remembered ({why})")
|
||||
@ -299,26 +343,53 @@ def get_memory(request: Request) -> Response:
|
||||
|
||||
|
||||
def _edit(request: Request, record: dict[str, Any], actor: dict[str, str], expected: int | None) -> dict[str, Any]:
|
||||
"""Supersede ``record`` with edited content under the same privacy shaping as a proposal (F3).
|
||||
|
||||
A deny verdict (deny topic, restricted, memory disabled, conversation
|
||||
forgotten) yields a content-free ``no_store`` record and leaves the old
|
||||
entry untouched; a sensitive floor makes the replacement ``proposed`` with
|
||||
``approval_mode: ask``. Private mode is refused outright.
|
||||
"""
|
||||
body = request.body if isinstance(request.body, dict) else {}
|
||||
content = body.get("content")
|
||||
if not isinstance(content, str) or not content.strip():
|
||||
raise Invalid("edit needs new content")
|
||||
if not isinstance(content, str) or not content.strip() or len(content) > 2000:
|
||||
raise Invalid("edit needs new content of at most 2000 chars")
|
||||
conversation_id = record["provenance"].get("conversation_id")
|
||||
verdict = _classify(request.store, {**body, "sensitivity": record["sensitivity"]}, content, conversation_id)
|
||||
stamp = now_iso()
|
||||
fresh = {
|
||||
**{k: v for k, v in record.items() if k not in {"revision", "supersedes"}}, "id": new_id("mem"), "content": verdict["content"],
|
||||
"topic": verdict["topic"], "sensitivity": verdict["sensitivity"], "status": "active", "approval_mode": "automatic", "retrievable": True,
|
||||
"created_at": stamp, "updated_at": stamp, "audit": [{"at": stamp, "action": "edited", "actor": actor, "note": f"edit of {record['id']}"}],
|
||||
}
|
||||
if not verdict["allowed"]:
|
||||
stored = _persist(request.store, _decline(fresh, verdict["why"], actor, stamp), None)
|
||||
_tombstone(request.store, stored["id"], verdict["why"])
|
||||
return stored
|
||||
old = _transition(request.store, request.identity, record, "forgotten", "superseded", actor, "edited", expected)
|
||||
from hux import events
|
||||
|
||||
events.redact_memory_references(request.store, old["id"])
|
||||
stamp = now_iso()
|
||||
fresh = {
|
||||
**{k: v for k, v in record.items() if k not in {"revision", "supersedes"}}, "id": new_id("mem"), "content": redaction.scrub_value(content.strip(), [])[:2000],
|
||||
"status": "active", "approval_mode": "automatic", "retrievable": True, "supersedes": old["id"], "created_at": stamp, "updated_at": stamp,
|
||||
"audit": [{"at": stamp, "action": "edited", "actor": actor, "note": f"supersedes {old['id']}"}, {"at": stamp, "action": "approved", "actor": actor}],
|
||||
}
|
||||
fresh["supersedes"] = old["id"]
|
||||
fresh["audit"][0]["note"] = f"supersedes {old['id']}"
|
||||
if fresh["sensitivity"] == "sensitive":
|
||||
fresh.update(status="proposed", approval_mode="ask", retrievable=False)
|
||||
fresh["audit"] = fresh["audit"][:1]
|
||||
else:
|
||||
fresh["audit"].append({"at": stamp, "action": "approved", "actor": actor})
|
||||
return _persist(request.store, fresh, None)
|
||||
|
||||
|
||||
def _set_retrievable(store: TenantStore, record: dict[str, Any], retrievable: bool, actor: dict[str, str], expected: int | None) -> dict[str, Any]:
|
||||
"""Flip ``retrievable`` on an active entry, re-reading it under the lock so a stale copy cannot overwrite a forget (F7)."""
|
||||
with store.lock(FAMILY):
|
||||
record = _current(store, record, expected)
|
||||
if record["status"] != "active":
|
||||
raise Conflict("retrieval can only change on active entries")
|
||||
stamp = now_iso()
|
||||
entry = {"at": stamp, "action": "approved" if retrievable else "retrieval_removed", "actor": actor, "note": "retrieval restored" if retrievable else ""}
|
||||
return _persist(store, {**record, "retrievable": retrievable, "updated_at": stamp, "audit": [*record["audit"], entry]}, record["revision"])
|
||||
|
||||
|
||||
def act_memory(request: Request) -> Response:
|
||||
"""``POST /hux/v1/memory/{id}/{action}``: approve, reject, forget, edit, remove_retrieval, restore_retrieval; every move emits a memory.* event."""
|
||||
memory_id, action = request.params["id"], request.params["action"]
|
||||
@ -349,13 +420,14 @@ def act_memory(request: Request) -> Response:
|
||||
_emit(store, identity, stored, "memory.forgotten", "Memory forgotten")
|
||||
elif action == "edit":
|
||||
stored = _edit(request, record, actor, expected)
|
||||
if stored["status"] == "no_store":
|
||||
_emit(store, identity, stored, "memory.suppressed", f"Edit not remembered ({stored['audit'][0]['note'].split(';')[0]})")
|
||||
request.audit("memory.edit", memory_id, "allow", "no_store")
|
||||
return Response(202, stored)
|
||||
_emit(store, identity, stored, "memory.committed" if stored["status"] == "active" else "memory.proposed", f"Memory edited, supersedes {record['id']}")
|
||||
else:
|
||||
if record["status"] != "active":
|
||||
raise Conflict("retrieval can only change on active entries")
|
||||
retrievable = action == "restore_retrieval"
|
||||
stamp = now_iso()
|
||||
stored = _persist(store, {**record, "retrievable": retrievable, "updated_at": stamp, "audit": [*record["audit"], {"at": stamp, "action": "retrieval_removed" if not retrievable else "approved", "actor": actor, "note": "retrieval restored" if retrievable else ""}]}, expected)
|
||||
stored = _set_retrievable(store, record, retrievable, actor, expected)
|
||||
_emit(store, identity, stored, "memory.retrieval_removed" if not retrievable else "memory.committed", "Retrieval removed" if not retrievable else "Retrieval restored")
|
||||
request.audit(f"memory.{action}", memory_id, "allow", "" if expected is not None else "unconditional_write")
|
||||
return Response(200, stored, {"ETag": str(stored["revision"])})
|
||||
|
||||
@ -12,7 +12,7 @@ from __future__ import annotations
|
||||
import re
|
||||
from typing import Any
|
||||
|
||||
from hux import contracts
|
||||
from hux import contracts, redaction
|
||||
from hux.errors import Conflict, Invalid, NotFound
|
||||
from hux.http import Request, Response, Router, page
|
||||
from hux.store import TenantStore, check_id, new_id, now_iso
|
||||
@ -63,8 +63,8 @@ def body_dict(request: Request) -> dict[str, Any]:
|
||||
|
||||
|
||||
def pick(body: dict[str, Any], fields: tuple[str, ...]) -> dict[str, Any]:
|
||||
"""Only the client-settable fields; ids, owner, timestamps and revision are server-set."""
|
||||
return {k: body[k] for k in fields if k in body}
|
||||
"""Only the client-settable fields, secret-scrubbed (F9); ids, owner, timestamps and revision are server-set."""
|
||||
return redaction.scrub_value({k: body[k] for k in fields if k in body}, [])
|
||||
|
||||
|
||||
def replay(store: TenantStore, family: str, key: str) -> dict[str, Any] | None:
|
||||
@ -233,14 +233,39 @@ def lineage(request: Request) -> Response:
|
||||
|
||||
# -- search --------------------------------------------------------------------
|
||||
|
||||
def artifact_titles(store: TenantStore, conversation: dict[str, Any]) -> list[str]:
|
||||
"""Artifact titles for a conversation via the artifacts lane, or nothing when it is absent."""
|
||||
def artifact_index(store: TenantStore) -> dict[str, list[str]]:
|
||||
"""Conversation id -> titles of the artifacts filed under it (F13c): both the conversation's ``artifact_ids`` and artifacts that name the conversation."""
|
||||
try:
|
||||
from hux import artifacts
|
||||
except ModuleNotFoundError:
|
||||
return []
|
||||
helper = getattr(artifacts, "artifact_titles", None)
|
||||
return list(helper(store, conversation["id"])) if helper else []
|
||||
return {}
|
||||
index: dict[str, list[str]] = {}
|
||||
for artifact in store.scan(artifacts.FAMILY):
|
||||
if isinstance(artifact.get("conversation_id"), str) and isinstance(artifact.get("title"), str):
|
||||
index.setdefault(artifact["conversation_id"], []).append(artifact["title"])
|
||||
for conversation in store.scan(CONVERSATIONS):
|
||||
ids = [i for i in conversation.get("artifact_ids", []) if isinstance(i, str)]
|
||||
titles = artifacts.artifact_titles(store, ids).values() if ids else []
|
||||
for title in titles:
|
||||
index.setdefault(conversation["id"], []).append(title)
|
||||
return index
|
||||
|
||||
|
||||
def artifact_titles(store: TenantStore, conversation: dict[str, Any], index: dict[str, list[str]] | None = None) -> list[str]:
|
||||
"""Artifact titles for a conversation via the artifacts lane, or nothing when it is absent."""
|
||||
index = artifact_index(store) if index is None else index
|
||||
return list(dict.fromkeys(index.get(conversation["id"], [])))
|
||||
|
||||
|
||||
def mark_forgotten(store: TenantStore, conversation_id: str) -> bool:
|
||||
"""Blank a forgotten conversation's title and tags and archive it (F9); False when there is no document."""
|
||||
with store.lock(CONVERSATIONS):
|
||||
if not conversation_exists(store, conversation_id):
|
||||
return False
|
||||
current = store.get(CONVERSATIONS, conversation_id)
|
||||
record = {**current, "title": "[forgotten]", "tags": [], "archived": True, "updated_at": now_iso()}
|
||||
store.put(CONVERSATIONS, checked(record), expected_revision=current["revision"])
|
||||
return True
|
||||
|
||||
|
||||
def tokens(text: str) -> list[str]:
|
||||
@ -269,11 +294,12 @@ def search(request: Request) -> Response:
|
||||
raise Invalid("q is required")
|
||||
project_filter = request.query.get("project_id")
|
||||
names = {p["id"]: p["name"] for p in request.store.scan(PROJECTS)}
|
||||
index = artifact_index(request.store)
|
||||
ranked = []
|
||||
for record in request.store.scan(CONVERSATIONS):
|
||||
if project_filter and record.get("project_id") != project_filter:
|
||||
continue
|
||||
points = score(record, terms, names.get(record.get("project_id", ""), ""), artifact_titles(request.store, record))
|
||||
points = score(record, terms, names.get(record.get("project_id", ""), ""), artifact_titles(request.store, record, index))
|
||||
if points:
|
||||
ranked.append((points, record))
|
||||
ranked.sort(key=lambda pair: (-pair[0], pair[1]["updated_at"]))
|
||||
|
||||
@ -186,11 +186,20 @@ def forget_conversation(store: TenantStore, identity: Identity, conversation_id:
|
||||
set_flag(store, conversation_id, "forgotten", True)
|
||||
forgotten_ids = memory.forget_from_conversation(store, identity, conversation_id)
|
||||
redacted = events.rewrite_full(store, conversation_id, "[forgotten conversation]", "conversation forgotten")
|
||||
counts = {"memory_forgotten": len(forgotten_ids), "events_redacted": redacted}
|
||||
counts = {"memory_forgotten": len(forgotten_ids), "events_redacted": redacted, "document_blanked": _blank_document(store, conversation_id)}
|
||||
store.append(FAMILY, "forgotten", {"conv_id": conversation_id, "requested_at": now_iso(), "purged_at": "", "counts": counts})
|
||||
return {"conversation_id": conversation_id, "forgotten": True, **counts}
|
||||
|
||||
|
||||
def _blank_document(store: TenantStore, conversation_id: str) -> bool:
|
||||
"""Blank the HUX-03 conversation document (title, tags, archived) when the organisation lane holds one (F9)."""
|
||||
try:
|
||||
from hux import organization
|
||||
except ModuleNotFoundError:
|
||||
return False
|
||||
return organization.mark_forgotten(store, conversation_id)
|
||||
|
||||
|
||||
# -- routes ---------------------------------------------------------------------
|
||||
|
||||
def _strip(record: dict[str, Any]) -> dict[str, Any]:
|
||||
@ -225,8 +234,12 @@ def post_notice(request: Request) -> Response:
|
||||
set_flag(request.store, conversation_id, "memory_disabled", True)
|
||||
elif chosen == "forget_this_conversation":
|
||||
forget_conversation(request.store, request.identity, conversation_id)
|
||||
events.emit(request.store, request.identity, conversation_id, "privacy.notice", text, notice, sensitivity=rules.PRIVACY_TOPICS[topic]["sensitivity"])
|
||||
request.audit("privacy.notice", conversation_id)
|
||||
# F13a: a notice for a conversation this subject never had must not conjure
|
||||
# an event ledger for it; the notice itself is still on record.
|
||||
known = events.conversation_known(request.store, conversation_id)
|
||||
if known:
|
||||
events.emit(request.store, request.identity, conversation_id, "privacy.notice", text, notice, sensitivity=rules.PRIVACY_TOPICS[topic]["sensitivity"])
|
||||
request.audit("privacy.notice", conversation_id, "allow", "" if known else "unknown_conversation")
|
||||
return Response(201, notice)
|
||||
|
||||
|
||||
|
||||
@ -136,8 +136,12 @@ def filter_detail(kind: str, detail: dict[str, Any] | None) -> dict[str, Any]:
|
||||
kept["steps"] = [str(step)[:200] for step in steps[:20]]
|
||||
if kind == "tool.call":
|
||||
# Raw arguments never persist: only their names, a hash and a byte count.
|
||||
path = kept.get("target_path")
|
||||
if not (isinstance(path, str) and path.startswith(WORKSPACE_PREFIX + "/")):
|
||||
# F13d: normalise before the prefix check so ``workspace/../.env`` cannot
|
||||
# pass as a workspace path; the stored value is the normalised one.
|
||||
path = os.path.normpath(kept["target_path"]) if isinstance(kept.get("target_path"), str) else ""
|
||||
if path.startswith(WORKSPACE_PREFIX + "/"):
|
||||
kept["target_path"] = path
|
||||
else:
|
||||
kept.pop("target_path", None)
|
||||
if "argument_names" in kept:
|
||||
names = kept["argument_names"] if isinstance(kept["argument_names"], list) else []
|
||||
|
||||
@ -16,6 +16,7 @@ from typing import Any
|
||||
from urllib.parse import urlsplit, urlunsplit
|
||||
|
||||
from hux.artifacts import remember, replay
|
||||
from hux import redaction
|
||||
from hux.contracts import load_all, validate_record
|
||||
from hux.errors import Conflict, Invalid, NotFound, TooLarge
|
||||
from hux.http import Request, Response, Router, page
|
||||
@ -62,6 +63,11 @@ def _string(body: dict[str, Any], key: str, limit: int, required: bool = True) -
|
||||
return value
|
||||
|
||||
|
||||
def _clean(value: Any) -> Any:
|
||||
"""Secret-scrub free text before it is stored (F9); ids and enums never carry secrets so they skip this."""
|
||||
return redaction.scrub_value(value, [])
|
||||
|
||||
|
||||
def _sha(*parts: str) -> str:
|
||||
return "sha256:" + hashlib.sha256("\n".join(parts).encode("utf-8")).hexdigest()
|
||||
|
||||
@ -127,7 +133,7 @@ def create_source(request: Request) -> Response:
|
||||
classification = body.get("classification", "unknown")
|
||||
if classification not in CLASSIFICATIONS:
|
||||
raise Invalid("unknown classification")
|
||||
title = _string(body, "title", 300)
|
||||
title = _clean(_string(body, "title", 300))
|
||||
uri = _string(body, "uri", 2000, required=False)
|
||||
dedupe = _sha("uri", normalise_uri(uri)) if uri is not None else _sha("title", kind, title.strip().lower())
|
||||
stamp = now_iso()
|
||||
@ -180,7 +186,7 @@ def create_passage(request: Request) -> Response:
|
||||
body = _body(request)
|
||||
key = request.idempotency_key()
|
||||
source = _resolve(request.store, SOURCES, body.get("source_id"))
|
||||
text = _string(body, "text", 4000)
|
||||
text = _clean(_string(body, "text", 4000))
|
||||
text_hash = "sha256:" + hashlib.sha256(text.encode("utf-8")).hexdigest()
|
||||
dedupe = _sha("passage", source["id"], text_hash)
|
||||
with request.store.lock(FAMILY):
|
||||
@ -215,7 +221,7 @@ def attach_citation(request: Request) -> Response:
|
||||
message_id = request.params["id"]
|
||||
if len(message_id) > 120:
|
||||
raise Invalid("message id is too long")
|
||||
claim = _string(body, "claim", 1000)
|
||||
claim = _clean(_string(body, "claim", 1000))
|
||||
support = body.get("support", "unverified")
|
||||
if support not in SUPPORT:
|
||||
raise Invalid("unknown support verdict")
|
||||
@ -231,7 +237,7 @@ def attach_citation(request: Request) -> Response:
|
||||
return Response(200, _resolve(request.store, CITATIONS, existing), {"HUX-Replayed": "true"})
|
||||
record: dict[str, Any] = {"schema": "hux.citation.v1", "id": new_id("cit"), "message_id": message_id, "claim": claim, "passage_ids": passage_ids, "support": support, "dedupe_key": dedupe}
|
||||
if _string(body, "note", 500, required=False) is not None:
|
||||
record["note"] = body["note"]
|
||||
record["note"] = _clean(body["note"])
|
||||
_check(record)
|
||||
stored = record | {}
|
||||
request.store.put(CITATIONS, record)
|
||||
@ -287,7 +293,7 @@ def _text_list(body: dict[str, Any], key: str, current: list[str]) -> list[str]:
|
||||
return current
|
||||
if not isinstance(values, list) or not all(isinstance(v, str) and v.strip() for v in values):
|
||||
raise Invalid(f"{key} must be a list of non-empty strings")
|
||||
return values
|
||||
return _clean(values)
|
||||
|
||||
|
||||
def create_notebook(request: Request) -> Response:
|
||||
@ -301,7 +307,7 @@ def create_notebook(request: Request) -> Response:
|
||||
return Response(200, request.store.get(NOTEBOOKS, existing), {"HUX-Replayed": "true"})
|
||||
record: dict[str, Any] = {
|
||||
"schema": "hux.research_notebook.v1", "id": new_id("nb"), "conversation_id": check_id(body.get("conversation_id")),
|
||||
"question": _string(body, "question", 1000), "status": "open", "source_ids": [], "passage_ids": [], "citation_ids": [],
|
||||
"question": _clean(_string(body, "question", 1000)), "status": "open", "source_ids": [], "passage_ids": [], "citation_ids": [],
|
||||
"assumptions": _text_list(body, "assumptions", []), "unresolved_questions": _text_list(body, "unresolved_questions", []),
|
||||
"updated_at": now_iso(), "notes": [], "revision": 1,
|
||||
}
|
||||
@ -350,7 +356,7 @@ def patch_notebook(request: Request) -> Response:
|
||||
if not isinstance(notes, list) or not all(isinstance(n, dict) for n in notes):
|
||||
raise Invalid("add_notes must be a list of objects")
|
||||
for note in notes:
|
||||
entry = {"at": stamp, "text": _string(note, "text", 2000)}
|
||||
entry = {"at": stamp, "text": _clean(_string(note, "text", 2000))}
|
||||
if note.get("source_id") is not None:
|
||||
entry["source_id"] = _resolve(request.store, SOURCES, note["source_id"])["id"]
|
||||
updated["notes"] = [*updated["notes"], entry]
|
||||
|
||||
@ -361,6 +361,68 @@ def test_flag_off_hides_event_routes(tmp_path):
|
||||
assert call(router, "GET", f"/hux/v1/conversations/{CONV}/events")[1]["code"] == "flag_off"
|
||||
|
||||
|
||||
def test_f5_crash_between_append_and_checkpoint_never_duplicates_a_seq(tmp_path, monkeypatch):
|
||||
"""F5 / SO-15: with the checkpoint left behind the ledger, the next emit takes max(checkpoint, tail + 1) so paging stays complete."""
|
||||
router = router_for(tmp_path)
|
||||
s = tenant(tmp_path)
|
||||
events.emit(s, ident(), CONV, "message.user", "one")
|
||||
original = store.TenantStore.put
|
||||
|
||||
def crash(self, family, record, expected_revision=None):
|
||||
if family == events.SEQ_FAMILY:
|
||||
raise RuntimeError("simulated crash after the ledger append")
|
||||
return original(self, family, record, expected_revision)
|
||||
|
||||
monkeypatch.setattr(store.TenantStore, "put", crash)
|
||||
with pytest.raises(RuntimeError):
|
||||
events.emit(s, ident(), CONV, "message.user", "two")
|
||||
monkeypatch.undo()
|
||||
assert s.get(events.SEQ_FAMILY, events._seq_id(CONV))["next_seq"] == 2 and events.last_seq(s, CONV) == 2
|
||||
third = events.emit(s, ident(), CONV, "message.user", "three")
|
||||
assert third["seq"] == 3 and s.get(events.SEQ_FAMILY, events._seq_id(CONV))["next_seq"] == 4
|
||||
items = call(router, "GET", f"/hux/v1/conversations/{CONV}/events")[1]["items"]
|
||||
assert [(e["seq"], e["summary"]) for e in items] == [(1, "one"), (2, "two"), (3, "three")]
|
||||
assert [e["seq"] for e in call(router, "GET", f"/hux/v1/conversations/{CONV}/events?after_seq=2")[1]["items"]] == [3]
|
||||
|
||||
|
||||
def test_f5_last_seq_reads_only_the_tail_and_skips_torn_lines(tmp_path):
|
||||
"""F5: a torn trailing line (crash mid-append) and a missing ledger are both handled."""
|
||||
s = tenant(tmp_path)
|
||||
assert events.last_seq(s, "conv_none0001") == 0
|
||||
for index in range(3):
|
||||
events.emit(s, ident(), CONV, "message.user", f"m{index}")
|
||||
path = s.root / events.FAMILY / f"{CONV}.jsonl"
|
||||
with open(path, "ab") as handle:
|
||||
handle.write(b'{"schema":"hux.event.v1","seq":9')
|
||||
assert events.last_seq(s, CONV) == 3
|
||||
with open(path, "ab") as handle:
|
||||
handle.write(b"\n[]\n\n")
|
||||
assert events.last_seq(s, CONV) == 3
|
||||
path.write_bytes(b"\n\n")
|
||||
assert events.last_seq(s, CONV) == 0
|
||||
|
||||
|
||||
def test_conversation_known_accepts_a_ledger_without_a_checkpoint(tmp_path):
|
||||
"""F5 / SO-18: a conversation whose checkpoint write was lost is still this subject's conversation."""
|
||||
s = tenant(tmp_path)
|
||||
events.emit(s, ident(), CONV, "message.user", "one")
|
||||
s.delete(events.SEQ_FAMILY, events._seq_id(CONV))
|
||||
assert events.conversation_known(s, CONV) is True
|
||||
assert events.conversation_known(s, "conv_other0001") is False
|
||||
assert events.emit(s, ident(), CONV, "message.user", "two")["seq"] == 2
|
||||
|
||||
|
||||
def test_f13d_target_path_is_normalised_before_the_workspace_check(tmp_path):
|
||||
"""F13d / SO-11: ``workspace/../.env`` is not a workspace path; a dotted path inside the workspace is stored normalised."""
|
||||
s = tenant(tmp_path)
|
||||
escaped = events.emit(s, ident(), CONV, "tool.call", "ran", {"tool": "bash", "target_path": "/opt/data/workspace/../.env"})
|
||||
assert "target_path" not in escaped["detail"]
|
||||
inside = events.emit(s, ident(), CONV, "tool.call", "ran", {"tool": "bash", "target_path": "/opt/data/workspace/a/../notes.md"})
|
||||
assert inside["detail"]["target_path"] == "/opt/data/workspace/notes.md"
|
||||
assert "target_path" not in events.emit(s, ident(), CONV, "tool.call", "ran", {"tool": "bash", "target_path": 7})["detail"]
|
||||
assert "target_path" not in events.emit(s, ident(), CONV, "tool.call", "ran", {"tool": "bash", "target_path": "/opt/data/workspace"})["detail"]
|
||||
|
||||
|
||||
def test_events_sources_stay_under_500_lines():
|
||||
for name in ("events.py", "redaction.py"):
|
||||
assert len((FOUNDATION / "hux" / name).read_text().splitlines()) <= 500, name
|
||||
|
||||
@ -10,7 +10,6 @@ from __future__ import annotations
|
||||
|
||||
import json
|
||||
import sys
|
||||
import types
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
@ -213,23 +212,42 @@ def test_search_ranks_title_over_tag_over_project_and_says_what_is_indexed(route
|
||||
assert call(router, "GET", "/hux/v1/search?q=cabinet", headers=OTHER)[1]["items"] == []
|
||||
|
||||
|
||||
def test_search_uses_artifact_titles_when_the_artifacts_lane_offers_them(router, monkeypatch):
|
||||
def test_f13c_search_finds_conversations_by_artifact_title(router, monkeypatch):
|
||||
"""F13c: an artifact filed under a conversation makes that conversation searchable by the artifact title; the lane may be absent."""
|
||||
conversation = make_conversation(router, title="Plain")
|
||||
fake = types.ModuleType("hux.artifacts")
|
||||
fake.artifact_titles = lambda store, conversation_id: ["Supplier comparison sheet"] if conversation_id == conversation["id"] else []
|
||||
monkeypatch.setitem(sys.modules, "hux.artifacts", fake)
|
||||
monkeypatch.setattr(hux, "artifacts", fake, raising=False)
|
||||
assert [c["id"] for c in call(router, "GET", "/hux/v1/search?q=supplier")[1]["items"]] == [conversation["id"]]
|
||||
bare = types.ModuleType("hux.artifacts")
|
||||
monkeypatch.setitem(sys.modules, "hux.artifacts", bare)
|
||||
monkeypatch.setattr(hux, "artifacts", bare, raising=False)
|
||||
assert call(router, "GET", "/hux/v1/search?q=supplier")[1]["items"] == []
|
||||
other = make_conversation(router, title="Other")
|
||||
status, artifact, _ = call(router, "POST", "/hux/v1/artifacts", {"type": "markdown", "title": "Supplier comparison sheet", "content": "x", "conversation_id": conversation["id"]})
|
||||
assert status == 201, artifact
|
||||
status, body, _ = call(router, "GET", "/hux/v1/search?q=supplier")
|
||||
assert [c["id"] for c in body["items"]] == [conversation["id"]] and body["scores"][conversation["id"]] == 1
|
||||
s = store.TenantStore(router.data_root, identity.resolve(HEADERS, {}))
|
||||
assert organization.artifact_titles(s, conversation) == ["Supplier comparison sheet"] and organization.artifact_titles(s, other) == []
|
||||
# The conversation's own artifact_ids list is indexed as well, without duplicates.
|
||||
s.put(organization.CONVERSATIONS, {**s.get(organization.CONVERSATIONS, other["id"]), "artifact_ids": [artifact["id"], "art_missing0001"]}, 1)
|
||||
assert organization.artifact_titles(s, other) == ["Supplier comparison sheet"]
|
||||
assert sorted(c["id"] for c in call(router, "GET", "/hux/v1/search?q=supplier")[1]["items"]) == sorted([conversation["id"], other["id"]])
|
||||
monkeypatch.setitem(sys.modules, "hux.artifacts", None)
|
||||
monkeypatch.delattr(hux, "artifacts", raising=False)
|
||||
assert call(router, "GET", "/hux/v1/search?q=supplier")[1]["items"] == []
|
||||
assert call(router, "GET", "/hux/v1/search?q=plain")[1]["items"] == [conversation]
|
||||
|
||||
|
||||
def test_f9_titles_tags_names_and_descriptions_are_secret_scrubbed(router):
|
||||
"""F9 / SO-07: a token pasted into a title, tag, project name or description is scrubbed before it is stored or indexed."""
|
||||
token = "ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"
|
||||
project = make_project(router, name=f"proj {token}", description="password: hunter2")
|
||||
assert token not in project["name"] and project["description"] == "[redacted:password]"
|
||||
# A tag can only be [a-z0-9-], so a scrubbed tag fails the contract instead of being stored.
|
||||
assert call(router, "POST", "/hux/v1/conversations", {"title": "t", "tags": [token]})[0] == 400
|
||||
conversation = make_conversation(router, title=f"token {token}", tags=["ok"])
|
||||
assert conversation["title"] == "token [redacted:github_token]" and conversation["tags"] == ["ok"]
|
||||
status, patched, _ = call(router, "PATCH", f"/hux/v1/conversations/{conversation['id']}", {"title": f"again {token}"}, {**HEADERS, "If-Match": "1"})
|
||||
assert status == 200 and token not in patched["title"]
|
||||
assert call(router, "GET", f"/hux/v1/search?q={token}")[1]["items"] == []
|
||||
leaked = [path for path in Path(router.data_root).rglob("*.json*") if token in path.read_text()]
|
||||
assert leaked == []
|
||||
|
||||
|
||||
def test_flag_off_hides_the_card(tmp_path):
|
||||
off = build_router(tmp_path, {"HUX_FLAGS": ""})
|
||||
status, body, _ = call(off, "GET", "/hux/v1/projects")
|
||||
|
||||
209
testing/tests/test_hermes_hux_memory_repairs.py
Normal file
209
testing/tests/test_hermes_hux_memory_repairs.py
Normal file
@ -0,0 +1,209 @@
|
||||
"""HUX-02 memory review repairs: edit shaping (F3), stale-write resurrection (F7), idempotency race (F13b).
|
||||
|
||||
Security obligations exercised: SO-21 (deny topics never reach memory, even
|
||||
through an edit), SO-22 (no_store and forget are content-free and final),
|
||||
SO-23 (a forgotten entry never retrieves again), SO-28 (private mode refuses
|
||||
memory writes), SO-44 (revision checks hold under concurrency).
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import sys
|
||||
import threading
|
||||
from pathlib import Path
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[2]
|
||||
FOUNDATION = ROOT / "dockerfiles" / "hermes-hux-foundation"
|
||||
if str(FOUNDATION) not in sys.path:
|
||||
sys.path.insert(0, str(FOUNDATION))
|
||||
|
||||
from hux import contracts, identity, memory, privacy, store # noqa: E402
|
||||
from hux.errors import Conflict # noqa: E402
|
||||
from hux.server import build_router # noqa: E402
|
||||
|
||||
SCHEMAS = contracts.load_all()
|
||||
HEADERS = {"X-Hermes-Tenant-Identity": "slot-3", "X-Hux-Subject": "usr_0123456789abcdef", "X-Hux-Surface": "chat"}
|
||||
ALL_ON = ",".join(card["flag"] for card in contracts.load_flags()["cards"])
|
||||
USER = {"proposed_by": "user", "approval_mode": "automatic"}
|
||||
|
||||
|
||||
def ident() -> identity.Identity:
|
||||
return identity.Identity(tenant_slot="slot-3", subject="usr_0123456789abcdef", surface="chat", trust="router")
|
||||
|
||||
|
||||
def router_for(tmp_path):
|
||||
return build_router(tmp_path, {"HUX_FLAGS": ALL_ON})
|
||||
|
||||
|
||||
def call(router, method, path, body=None, headers=None):
|
||||
raw = json.dumps(body).encode() if body is not None else b""
|
||||
response = router.dispatch(method, path, {**HEADERS, **(headers or {})}, raw)
|
||||
return response.status, response.body
|
||||
|
||||
|
||||
def tenant(tmp_path) -> store.TenantStore:
|
||||
return store.TenantStore(tmp_path, ident())
|
||||
|
||||
|
||||
def conversation(router) -> str:
|
||||
status, body = call(router, "POST", "/hux/v1/conversations", {"title": "tea"})
|
||||
assert status == 201
|
||||
return body["id"]
|
||||
|
||||
|
||||
def active(router, conv, content="I like tea") -> dict:
|
||||
status, body = call(router, "POST", "/hux/v1/memory", {"content": content, "conversation_id": conv, **USER})
|
||||
assert (status, body["status"]) == (201, "active"), body
|
||||
return body
|
||||
|
||||
|
||||
def valid(record) -> dict:
|
||||
assert contracts.validate_record(record, SCHEMAS) == [], record
|
||||
return record
|
||||
|
||||
|
||||
# --- F3: edit goes through the same shaping as a proposal ---------------------------
|
||||
|
||||
def test_f3_edit_with_deny_topic_becomes_no_store_and_leaves_the_original(tmp_path):
|
||||
"""F3 / SO-21 / SO-22: a location + minors edit is refused content-free; the old entry is not superseded."""
|
||||
router = router_for(tmp_path)
|
||||
conv = conversation(router)
|
||||
original = active(router, conv)
|
||||
status, body = call(router, "POST", f"/hux/v1/memory/{original['id']}/edit", {"content": "my home address is 12 Main Street; my daughter is 6"}, {"If-Match": "1"})
|
||||
assert status == 202 and valid(body)["status"] == "no_store" and body["content"] == "" and "topic" not in body
|
||||
assert body["sensitivity"] == "restricted" and body["retrievable"] is False and "supersedes" not in body
|
||||
assert body["audit"][0]["note"] == "restricted; topic=location" and "Main Street" not in json.dumps(body)
|
||||
assert body["id"] in memory.tombstoned(tenant(tmp_path))
|
||||
status, again = call(router, "GET", f"/hux/v1/memory/{original['id']}")
|
||||
assert status == 200 and again["status"] == "active" and again["revision"] == 1 and again["content"] == "I like tea"
|
||||
assert [hit["content"] for hit in memory.retrieve(tenant(tmp_path), ["address"])] == []
|
||||
assert [hit["id"] for hit in memory.retrieve(tenant(tmp_path), ["tea"])] == [original["id"]]
|
||||
kinds = [row["kind"] for row in tenant(tmp_path).read("events", conv)]
|
||||
assert kinds[-1] == "memory.suppressed"
|
||||
|
||||
|
||||
def test_f3_edit_with_sensitive_content_is_proposed_with_ask(tmp_path):
|
||||
"""F3 / SO-21: the sensitivity floor from detected topics applies to edits, so a health edit waits for approval."""
|
||||
router = router_for(tmp_path)
|
||||
conv = conversation(router)
|
||||
original = active(router, conv)
|
||||
status, body = call(router, "POST", f"/hux/v1/memory/{original['id']}/edit", {"content": "I take lithium for my diagnosis"})
|
||||
assert status == 200 and valid(body)["status"] == "proposed" and body["approval_mode"] == "ask"
|
||||
assert body["sensitivity"] == "sensitive" and body["topic"] == "health" and body["retrievable"] is False
|
||||
assert body["supersedes"] == original["id"] and [row["action"] for row in body["audit"]] == ["edited"]
|
||||
assert call(router, "GET", f"/hux/v1/memory/{original['id']}")[1]["status"] == "forgotten"
|
||||
|
||||
|
||||
def test_f3_edit_scrubs_secrets_and_floors_credentials(tmp_path):
|
||||
"""F3 / SO-07: a secret in edited content is scrubbed and the credentials topic makes the edit no_store."""
|
||||
router = router_for(tmp_path)
|
||||
conv = conversation(router)
|
||||
original = active(router, conv)
|
||||
status, body = call(router, "POST", f"/hux/v1/memory/{original['id']}/edit", {"content": "token ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"})
|
||||
assert status == 202 and body["status"] == "no_store" and body["sensitivity"] == "restricted"
|
||||
leaked = [path for path in tmp_path.rglob("*.json*") if "ghp_ABCDEFGHIJ" in path.read_text()]
|
||||
assert leaked == []
|
||||
|
||||
|
||||
def test_f3_edit_refuses_private_mode_and_declines_when_memory_disabled(tmp_path):
|
||||
"""F3 / SO-28: private mode is a 403 for edits too; disable_memory_here turns the edit into no_store."""
|
||||
router = router_for(tmp_path)
|
||||
conv = conversation(router)
|
||||
original = active(router, conv)
|
||||
privacy.set_flag(tenant(tmp_path), conv, "memory_disabled", True)
|
||||
status, body = call(router, "POST", f"/hux/v1/memory/{original['id']}/edit", {"content": "I like coffee"})
|
||||
assert status == 202 and body["status"] == "no_store" and body["audit"][0]["note"].startswith("memory_disabled")
|
||||
assert call(router, "GET", f"/hux/v1/memory/{original['id']}")[1]["status"] == "active"
|
||||
privacy.set_flag(tenant(tmp_path), conv, "memory_disabled", False)
|
||||
assert call(router, "PATCH", f"/hux/v1/conversations/{conv}", {"mode": "private"})[0] == 200
|
||||
status, body = call(router, "POST", f"/hux/v1/memory/{original['id']}/edit", {"content": "I like coffee"})
|
||||
assert status == 403 and body["code"] == "forbidden"
|
||||
assert call(router, "GET", f"/hux/v1/memory/{original['id']}")[1]["status"] == "active"
|
||||
|
||||
|
||||
def test_f3_edit_still_needs_content(tmp_path):
|
||||
router = router_for(tmp_path)
|
||||
original = active(router, conversation(router))
|
||||
assert call(router, "POST", f"/hux/v1/memory/{original['id']}/edit", {"content": " "})[0] == 400
|
||||
assert call(router, "POST", f"/hux/v1/memory/{original['id']}/edit", {"content": "x" * 2001})[0] == 400
|
||||
|
||||
|
||||
# --- F7: a stale record cannot resurrect a forgotten entry ------------------------
|
||||
|
||||
def test_f7_stale_approve_after_forget_is_a_conflict(tmp_path, monkeypatch):
|
||||
"""F7 / SO-22 / SO-44: approve holding a pre-reject snapshot (the a2-B interleaving) is refused and the entry stays rejected."""
|
||||
router = router_for(tmp_path)
|
||||
conv = conversation(router)
|
||||
status, proposed = call(router, "POST", "/hux/v1/memory", {"content": "I take lithium daily", "conversation_id": conv, "proposed_by": "user", "approval_mode": "ask"})
|
||||
assert (status, proposed["status"]) == (201, "proposed")
|
||||
stale = dict(proposed)
|
||||
assert call(router, "POST", f"/hux/v1/memory/{proposed['id']}/reject")[1]["status"] == "rejected"
|
||||
monkeypatch.setattr(memory, "load", lambda store, memory_id, now=None: stale)
|
||||
status, body = call(router, "POST", f"/hux/v1/memory/{proposed['id']}/approve")
|
||||
assert status == 409 and body["code"] == "conflict"
|
||||
monkeypatch.undo()
|
||||
status, after = call(router, "GET", f"/hux/v1/memory/{proposed['id']}")
|
||||
assert after["status"] == "rejected" and after["content"] == "" and after["revision"] == 2
|
||||
assert memory.retrieve(tenant(tmp_path), ["lithium"]) == []
|
||||
assert call(router, "GET", "/hux/v1/memory/export")[1]["items"] == []
|
||||
|
||||
|
||||
def test_f7_stale_remove_retrieval_cannot_overwrite_a_forget(tmp_path, monkeypatch):
|
||||
"""F7: the unconditional retrieval branch re-reads the entry under the lock, so the a5 race ends in 409 and no restore works."""
|
||||
router = router_for(tmp_path)
|
||||
conv = conversation(router)
|
||||
entry = active(router, conv, "I take lithium daily")
|
||||
stale = dict(entry)
|
||||
assert call(router, "POST", f"/hux/v1/memory/{entry['id']}/forget")[1]["status"] == "forgotten"
|
||||
monkeypatch.setattr(memory, "load", lambda store, memory_id, now=None: stale)
|
||||
assert call(router, "POST", f"/hux/v1/memory/{entry['id']}/remove_retrieval")[0] == 409
|
||||
monkeypatch.undo()
|
||||
assert call(router, "POST", f"/hux/v1/memory/{entry['id']}/restore_retrieval")[0] == 409
|
||||
assert call(router, "GET", f"/hux/v1/memory/{entry['id']}")[1]["status"] == "forgotten"
|
||||
assert memory.retrieve(tenant(tmp_path), ["lithium"]) == []
|
||||
|
||||
|
||||
def test_f7_transition_checks_if_match_against_the_fresh_copy(tmp_path):
|
||||
"""F7 / SO-44: If-Match is compared with the stored revision, not the caller's stale snapshot."""
|
||||
router = router_for(tmp_path)
|
||||
entry = active(router, conversation(router))
|
||||
s = tenant(tmp_path)
|
||||
memory._set_retrievable(s, entry, False, {"type": "user", "id": "usr_0123456789abcdef"}, 1)
|
||||
try:
|
||||
memory._transition(s, ident(), entry, "forgotten", "forgotten", {"type": "user", "id": "usr_0123456789abcdef"}, expected=1)
|
||||
except Conflict as error:
|
||||
assert "revision 1 does not match current revision 2" in str(error)
|
||||
else:
|
||||
raise AssertionError("stale If-Match must conflict")
|
||||
assert memory.load(s, entry["id"])["status"] == "active"
|
||||
stored = memory._transition(s, ident(), entry, "forgotten", "forgotten", {"type": "user", "id": "usr_0123456789abcdef"})
|
||||
assert stored["status"] == "forgotten" and stored["revision"] == 3
|
||||
|
||||
|
||||
# --- F13b: Idempotency-Key is checked under the lock ------------------------------
|
||||
|
||||
def test_f13b_concurrent_proposals_with_one_key_create_exactly_one_record(tmp_path):
|
||||
"""F13b: eight threads racing one Idempotency-Key produce one memory entry; seven are 200 replays."""
|
||||
router = router_for(tmp_path)
|
||||
results: list[tuple[int, str]] = []
|
||||
barrier = threading.Barrier(8)
|
||||
|
||||
def propose() -> None:
|
||||
barrier.wait(5)
|
||||
status, body = call(router, "POST", "/hux/v1/memory", {"content": "I like tea", **USER}, {"Idempotency-Key": "same-key-0001"})
|
||||
results.append((status, body["id"]))
|
||||
|
||||
threads = [threading.Thread(target=propose) for _ in range(8)]
|
||||
for thread in threads:
|
||||
thread.start()
|
||||
for thread in threads:
|
||||
thread.join(10)
|
||||
assert sorted(status for status, _ in results) == [200] * 7 + [201]
|
||||
assert len({memory_id for _, memory_id in results}) == 1
|
||||
assert len(call(router, "GET", "/hux/v1/memory")[1]["items"]) == 1
|
||||
assert len(tenant(tmp_path).read(memory.IDEM_FAMILY, "keys")) == 1
|
||||
|
||||
|
||||
def test_memory_module_stays_under_500_lines():
|
||||
assert len((FOUNDATION / "hux" / "memory.py").read_text().splitlines()) <= 500
|
||||
@ -90,6 +90,7 @@ def test_policy_route_serves_rules_and_reports_stale_audit(tmp_path):
|
||||
|
||||
def test_notice_is_recorded_scoped_and_emitted(tmp_path):
|
||||
router = router_for(tmp_path)
|
||||
events.emit(tenant(tmp_path), ident(), CONV, "message.user", "hello")
|
||||
status, body, _ = call(router, "POST", "/hux/v1/privacy/notices", HEADERS, {"topic": "health", "conversation_id": CONV, "controls": ["dismiss", "forget_this_conversation", "bogus"]})
|
||||
assert status == 201 and body["schema"] == "hux.privacy_notice.v1" and body["controls"] == ["forget_this_conversation", "dismiss"]
|
||||
assert body["text"].startswith("This looks like a health topic")
|
||||
@ -98,7 +99,7 @@ def test_notice_is_recorded_scoped_and_emitted(tmp_path):
|
||||
assert s.read(privacy.FAMILY, "notices") == [body]
|
||||
state = privacy.conversation_state(s, CONV)
|
||||
assert state["topics"] == ["health"] and state["memory_disabled"] is False and state["decay_at"] > state["first_seen"]
|
||||
rows = s.read(events.FAMILY, CONV)
|
||||
rows = s.read(events.FAMILY, CONV)[1:]
|
||||
assert [r["kind"] for r in rows] == ["privacy.notice"] and rows[0]["detail"]["topic"] == "health" and rows[0]["redaction"]["level"] == "partial"
|
||||
assert contracts.validate_record(rows[0], SCHEMAS) == []
|
||||
status, again, _ = call(router, "POST", "/hux/v1/privacy/notices", HEADERS, {"topic": "relationships", "conversation_id": CONV})
|
||||
@ -142,7 +143,7 @@ def test_forget_route_redacts_events_and_marks_state(tmp_path):
|
||||
events.emit(s, ident(), CONV, "message.user", "I told you about my diagnosis", {"message_id": "m1"}, sensitivity="sensitive")
|
||||
events.emit(s, ident(), CONV, "message.assistant", "noted")
|
||||
status, body, _ = call(router, "POST", f"/hux/v1/conversations/{CONV}/forget")
|
||||
assert status == 200 and body == {"conversation_id": CONV, "forgotten": True, "memory_forgotten": 0, "events_redacted": 2}
|
||||
assert status == 200 and body == {"conversation_id": CONV, "forgotten": True, "memory_forgotten": 0, "events_redacted": 2, "document_blanked": False}
|
||||
rows = s.read(events.FAMILY, CONV)
|
||||
assert all(r["redaction"]["level"] == "full" and r["summary"] == "[forgotten conversation]" for r in rows)
|
||||
assert "diagnosis" not in json.dumps(rows)
|
||||
@ -170,3 +171,38 @@ def test_flag_off_hides_privacy_routes(tmp_path):
|
||||
|
||||
def test_privacy_module_stays_under_500_lines():
|
||||
assert len((FOUNDATION / "hux" / "privacy.py").read_text().splitlines()) <= 500
|
||||
|
||||
|
||||
def test_f13a_notice_for_an_unknown_conversation_does_not_create_a_ghost_ledger(tmp_path):
|
||||
"""F13a / SO-18: the notice is stored and scoped, but no event ledger or checkpoint is conjured for an unknown conversation."""
|
||||
router = router_for(tmp_path)
|
||||
s = tenant(tmp_path)
|
||||
status, body, _ = call(router, "POST", "/hux/v1/privacy/notices", HEADERS, {"topic": "health", "conversation_id": "conv_ghost0001"})
|
||||
assert status == 201 and s.read(privacy.FAMILY, "notices") == [body]
|
||||
assert privacy.conversation_state(s, "conv_ghost0001")["topics"] == ["health"]
|
||||
assert events.conversation_known(s, "conv_ghost0001") is False and s.read(events.FAMILY, "conv_ghost0001") == []
|
||||
assert [(r["action"], r["reason"]) for r in audit.recent(s) if r["action"] == "privacy.notice"] == [("privacy.notice", "unknown_conversation")]
|
||||
assert call(router, "GET", "/hux/v1/conversations/conv_ghost0001/events", HEADERS)[0] == 404
|
||||
|
||||
|
||||
def test_f9_forget_blanks_the_conversation_document(tmp_path):
|
||||
"""F9 / SO-22: forgetting a HUX-03 conversation blanks its title and tags and archives it, so nothing of it is searchable."""
|
||||
router = router_for(tmp_path)
|
||||
status, doc, _ = call(router, "POST", "/hux/v1/conversations", HEADERS, {"title": "Doctor visit notes", "tags": ["health", "private"]})
|
||||
assert status == 201
|
||||
status, body, _ = call(router, "POST", f"/hux/v1/conversations/{doc['id']}/forget", HEADERS)
|
||||
assert status == 200 and body["document_blanked"] is True
|
||||
status, after, _ = call(router, "GET", f"/hux/v1/conversations/{doc['id']}", HEADERS)
|
||||
assert after["title"] == "[forgotten]" and after["tags"] == [] and after["archived"] is True and after["revision"] == 2
|
||||
assert contracts.validate_record(after, SCHEMAS) == []
|
||||
assert call(router, "GET", "/hux/v1/search?q=doctor", HEADERS)[1]["items"] == []
|
||||
assert privacy._blank_document(tenant(tmp_path), "conv_none000001") is False
|
||||
|
||||
|
||||
def test_f9_forget_tolerates_a_missing_organization_lane(tmp_path, monkeypatch):
|
||||
"""F9: without the HUX-03 module there is no document to blank and forget still succeeds."""
|
||||
import hux
|
||||
|
||||
monkeypatch.setitem(sys.modules, "hux.organization", None)
|
||||
monkeypatch.delattr(hux, "organization", raising=False)
|
||||
assert privacy._blank_document(tenant(tmp_path), CONV) is False
|
||||
|
||||
@ -227,6 +227,20 @@ def test_validate_citations_reports_integrity_problems(router, monkeypatch):
|
||||
assert cite(router, [passage(router, source(router, "https://b.example/")["id"])["id"]], claim="silent", conversation_id="conv_0001abcd")["support"] == "supports"
|
||||
|
||||
|
||||
def test_f9_source_title_passage_text_and_citation_claim_are_scrubbed(router):
|
||||
"""F9 / SO-07: secrets in free-text research fields never persist; the passage hash covers the scrubbed text."""
|
||||
token = "ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"
|
||||
src = source(router, title=f"notes {token}")
|
||||
assert src["title"] == "notes [redacted:github_token]"
|
||||
psg = passage(router, src["id"], f"password: hunter2 and {token}")
|
||||
assert psg["text"] == "[redacted:password] and [redacted:github_token]"
|
||||
assert psg["hash"] == "sha256:" + hashlib.sha256(psg["text"].encode()).hexdigest()
|
||||
cit = cite(router, [psg["id"]], claim=f"claim {token}", note=f"note {token}")
|
||||
assert token not in cit["claim"] and token not in cit["note"]
|
||||
leaked = [path for path in Path(router.data_root).rglob("*.json*") if token in path.read_text()]
|
||||
assert leaked == []
|
||||
|
||||
|
||||
def test_research_modules_stay_under_500_lines():
|
||||
assert len((FOUNDATION / "hux" / "research.py").read_text().splitlines()) <= 500
|
||||
|
||||
|
||||
@ -160,3 +160,15 @@ def test_notebook_paths_are_tenant_scoped(router):
|
||||
assert (tenant.root / "notebooks" / f"{record['id']}.json").exists()
|
||||
assert call(router, "GET", "/hux/v1/notebooks/nb_missing000001")[0] == 404
|
||||
assert call(router, "GET", "/hux/v1/notebooks/NB")[0] == 404
|
||||
|
||||
|
||||
def test_f9_notebook_free_text_is_scrubbed(router, evidence):
|
||||
"""F9 / SO-07: question, assumptions, unresolved questions and notes are secret-scrubbed on create and patch."""
|
||||
token = "ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"
|
||||
record = notebook(router, question=f"q {token}", assumptions=[f"a {token}"], unresolved_questions=[f"u {token}"])
|
||||
assert record["question"] == "q [redacted:github_token]" and record["assumptions"] == ["a [redacted:github_token]"]
|
||||
assert record["unresolved_questions"] == ["u [redacted:github_token]"]
|
||||
status, patched, _ = call(router, "PATCH", f"/hux/v1/notebooks/{record['id']}", {"add_notes": [{"text": f"n {token}"}], "assumptions": [f"b {token}"]}, {"If-Match": "1"})
|
||||
assert status == 200 and patched["notes"][0]["text"] == "n [redacted:github_token]" and patched["assumptions"] == ["b [redacted:github_token]"]
|
||||
leaked = [path for path in Path(router.data_root).rglob("*.json*") if token in path.read_text()]
|
||||
assert leaked == []
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user