Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion SECURITY.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
2 changes: 2 additions & 0 deletions api/gdpr.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
19 changes: 18 additions & 1 deletion api/memory.py
Original file line number Diff line number Diff line change
Expand Up @@ -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__)
Expand Down Expand Up @@ -258,13 +263,25 @@ 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,
"workspace_id": result.workspace_id,
"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:
Expand Down
2 changes: 2 additions & 0 deletions graph/gdpr.py
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,7 @@ class GdprErasureResult:
person_id: str
decisions_deleted: int
requested_by: str
decision_ids: tuple[str, ...] = ()


class GdprErasureService:
Expand Down Expand Up @@ -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),
)
44 changes: 44 additions & 0 deletions memory/episodic.py
Original file line number Diff line number Diff line change
Expand Up @@ -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))

30 changes: 30 additions & 0 deletions memory/semantic.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

29 changes: 29 additions & 0 deletions tests/memory/test_gdpr_store_purge.py
Original file line number Diff line number Diff line change
@@ -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
Loading