diff --git a/SECURITY.md b/SECURITY.md index 8a5f5a6..429a1a7 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -41,7 +41,7 @@ Cortex does not yet have tagged releases; security fixes are applied to the `POST /gdpr/erase` (`api/gdpr.py`) cascades deletes for a person across Decisions, Rationales, and Contradictions in the knowledge graph, and writes -a `GdprAuditLog` entry. It requires the `admin`, `gdpr_officer`, or `legal` +a `GdprAuditLog` entry, deletes matching Qdrant decision vectors, and purges TimescaleDB `cortex_raw_events` rows for that author. It requires the `admin`, `gdpr_officer`, or `legal` role — protect `CORTEX_API_KEYS` accordingly, since anyone who can obtain one of those roles can erase organizational memory. diff --git a/api/gdpr.py b/api/gdpr.py index adfa814..bc5b86e 100644 --- a/api/gdpr.py +++ b/api/gdpr.py @@ -40,6 +40,8 @@ class GdprEraseResponse(BaseModel): person_id: str decisions_deleted: int requested_by: str + vectors_deleted: int = 0 + raw_events_deleted: int = 0 @router.post( diff --git a/api/memory.py b/api/memory.py index d0cf5e8..f6159f5 100644 --- a/api/memory.py +++ b/api/memory.py @@ -12,7 +12,12 @@ from graph.gdpr import GdprErasureService from graph.query import GraphQueryService -from memory.semantic import search_decision_ids, semantic_enabled +from memory.episodic import purge_raw_events_for_person +from memory.semantic import ( + delete_decision_vectors, + search_decision_ids, + semantic_enabled, +) from scoring.trust_scorer import is_injectable log = structlog.get_logger(__name__) @@ -258,6 +263,16 @@ async def erase_gdpr_subject( caller_roles=caller_roles, reason=reason, ) + vectors_deleted = await asyncio.to_thread( + delete_decision_vectors, + list(result.decision_ids), + workspace_id, + ) + raw_events_deleted = await asyncio.to_thread( + purge_raw_events_for_person, + workspace_id, + person_id, + ) self.invalidate_workspace_cache(workspace_id) return { "audit_id": result.audit_id, @@ -265,6 +280,8 @@ async def erase_gdpr_subject( "person_id": result.person_id, "decisions_deleted": result.decisions_deleted, "requested_by": result.requested_by, + "vectors_deleted": vectors_deleted, + "raw_events_deleted": raw_events_deleted, } async def neo4j_health(self) -> str: diff --git a/graph/gdpr.py b/graph/gdpr.py index c39fa91..2d75ab5 100644 --- a/graph/gdpr.py +++ b/graph/gdpr.py @@ -74,6 +74,7 @@ class GdprErasureResult: person_id: str decisions_deleted: int requested_by: str + decision_ids: tuple[str, ...] = () class GdprErasureService: @@ -218,4 +219,5 @@ def _erase_transaction( person_id=person_id, decisions_deleted=len(decision_ids), requested_by=requested_by, + decision_ids=tuple(decision_ids), ) diff --git a/memory/episodic.py b/memory/episodic.py index 3be4453..2c009b4 100644 --- a/memory/episodic.py +++ b/memory/episodic.py @@ -76,3 +76,47 @@ def append_raw_event(raw: RawEvent) -> None: if _timescale_dsn() is None: return asyncio.run(_append_async(raw)) + + +async def _purge_raw_events_async(workspace_id: str, person_id: str) -> int: + import asyncpg + + dsn = _timescale_dsn() + if dsn is None: + return 0 + conn = await asyncpg.connect(dsn) + try: + await _ensure_table(conn) + # RawEvent.author is the canonical person id; also catch nested metadata. + result = await conn.execute( + """ + DELETE FROM cortex_raw_events + WHERE workspace_id = $1 + AND ( + payload->>'author' = $2 + OR payload->'metadata'->>'author' = $2 + OR payload->'metadata'->>'user_id' = $2 + ) + """, + workspace_id, + person_id, + ) + # asyncpg returns "DELETE N" + deleted = int(result.split()[-1]) if result else 0 + log.info( + "episodic.purge", + workspace_id=workspace_id, + person_id=person_id, + deleted=deleted, + ) + return deleted + finally: + await conn.close() + + +def purge_raw_events_for_person(workspace_id: str, person_id: str) -> int: + """Delete Timescale RawEvent rows authored by ``person_id`` in the workspace.""" + if _timescale_dsn() is None: + return 0 + return asyncio.run(_purge_raw_events_async(workspace_id, person_id)) + diff --git a/memory/semantic.py b/memory/semantic.py index 1562bbe..86f3bc8 100644 --- a/memory/semantic.py +++ b/memory/semantic.py @@ -116,3 +116,33 @@ def search_decision_ids(query: str, workspace_id: str, limit: int) -> list[str]: ), ) return [str(hit.payload.get("event_id", "")) for hit in hits if hit.payload] + + +def delete_decision_vectors(decision_ids: list[str], workspace_id: str) -> int: + """Remove Qdrant points for erased decision event ids. + + Point ids are deterministic ``uuid5(NAMESPACE_URL, event_id)`` (see upsert). + Returns the number of ids submitted for deletion (0 when semantic store is off). + """ + client = _client() + if client is None or not decision_ids: + return 0 + point_ids = [str(uuid.uuid5(uuid.NAMESPACE_URL, eid)) for eid in decision_ids] + try: + client.delete(collection_name=_COLLECTION, points_selector=point_ids) + except Exception as exc: + log.warning( + "semantic.delete_failed", + error=str(exc), + workspace_id=workspace_id, + count=len(point_ids), + ) + return 0 + log.info( + "semantic.delete", + workspace_id=workspace_id, + deleted=len(point_ids), + collection=_COLLECTION, + ) + return len(point_ids) + diff --git a/tests/memory/test_gdpr_store_purge.py b/tests/memory/test_gdpr_store_purge.py new file mode 100644 index 0000000..d289d37 --- /dev/null +++ b/tests/memory/test_gdpr_store_purge.py @@ -0,0 +1,29 @@ +"""Unit tests for Qdrant + Timescale GDPR purge helpers.""" + +from __future__ import annotations + +from unittest.mock import MagicMock, patch + +from memory.episodic import purge_raw_events_for_person +from memory.semantic import delete_decision_vectors + + +@patch("memory.semantic._client") +def test_delete_decision_vectors_uses_deterministic_point_ids(mock_client: MagicMock) -> None: + client = MagicMock() + mock_client.return_value = client + deleted = delete_decision_vectors(["dec-1", "dec-2"], "ws-1") + assert deleted == 2 + client.delete.assert_called_once() + kwargs = client.delete.call_args.kwargs + assert len(kwargs["points_selector"]) == 2 + + +@patch("memory.semantic._client", return_value=None) +def test_delete_decision_vectors_noop_when_disabled(_mock: MagicMock) -> None: + assert delete_decision_vectors(["dec-1"], "ws-1") == 0 + + +@patch("memory.episodic._timescale_dsn", return_value=None) +def test_purge_raw_events_noop_without_dsn(_mock: MagicMock) -> None: + assert purge_raw_events_for_person("ws-1", "alice@x.com") == 0